airflow
Learning Path Diagrams Book Mode
15 units • 5 levels
Change Guide
Welcome to airflow
Your guide to understanding the codebase
Start Reading
5 levels
15 learning units
Orchestration & APIs
- Pipelines, workflows, public interfaces • 3 units
DAGs and Task Fundamentals
- Data Model·12
Data Persistence Layer
- Data Model
Configuration and CLI Basics
- Infrastructure·11
Core Logic & Data
- Business rules, schemas, models • 3 units
REST API and Service Layer
- API
Scheduling and Execution
- Workflow·5
Task Communication and Dependencies
- Workflow·7
Interaction & Integration
- UI components, external connectors • 3 units
Core Architecture and Executors
- Architecture·8
UI Architecture and Visualization
- Architecture·6
Security and Integration Systems
- Security·29
Cross-Cutting Concerns
- Auth, logging, config, testing • 2 units
Concurrency and Resource Management
- Concurrency·4
Observability and Plugin Systems
- Infrastructure·3
Edge Cases & Resilience
- Error handling, fault tolerance • 1 units
Additional Data Model Patterns
- Data Model·5
Test Your Knowledge
Test your deep understanding of the codebase.
Start Quiz
Progress 0/26 answered
Question Tiers - ordered for deep understanding
- Why 8
- Purpose & Problem Architecture 10
- Design & Patterns Code 8
- Implementation
Hands-On Assignment
2-3 hours
Start Assignment
Your Challenge
Build a custom executor that implements a "Priority Queue" execution strategy, where tasks can be assigned priority levels and higher-priority tasks are executed before lower-priority ones when resources are constrained. Your executor should integrate with Airflow's existing executor framework, respect concurrency limits, and allow priority to be set via task configuration. This will require understanding how executors interact with the scheduler, task queuing mechanisms, and the core execution loop.
Starting Points
airflow-core/docs/img/diagram_basic_airflow_architecture.py:1-103
Study the architecture diagram to understand how executors fit into the overall system - note the interaction between Scheduler, Executor, and Workersairflow/executors/base_executor.py
Examine the BaseExecutor class that all executors inherit from - understand the core methods like execute_async, sync, heartbeat, and the queued_tasks structureairflow/executors/local_executor.py
Study this as a reference implementation - see how it manages task queuing, worker processes, and result handlingairflow/executors/sequential_executor.py
Look at the simplest executor implementation to understand the minimal contract you need to fulfillairflow/executors/
Create your new PriorityQueueExecutor class here, following the naming and structure conventions of existing executorsairflow/models/taskinstance.py
Explore how task instances store metadata - you may need to understand task attributes to extract priority informationairflow/config_templates/default_airflow.cfg
See how executor configuration is handled - you'll need to add configuration options for your executor
Success Criteria
- Executor successfully inherits from BaseExecutor and implements all required methods
- Tasks with higher priority values execute before lower priority tasks when parallelism limit is reached
- Executor respects the configured parallelism limit (doesn't exceed max concurrent tasks)
- Tasks complete successfully and state changes are properly reported back to the scheduler
- A test DAG with 10 tasks (5 high priority, 5 low priority) and parallelism=2 shows high-priority tasks completing first
- Executor can be configured via airflow.cfg and selected as the active executor
- Existing Airflow functionality remains unaffected when using standard executors
Hints
(click to reveal)
- Hint 1: Understanding the executor contract conceptual
- Hint 2: Priority queue data structure implementation
- Hint 3: Respecting concurrency limits code location
- Hint 4: Task execution mechanism implementation
- Hint 5: Testing your executor implementation
Prerequisites:
- Python multiprocessing or subprocess
- Priority queue data structures
- Airflow executor architecture
- Task lifecycle and state management