Apache Airflow Orchestration
A comprehensive skill for mastering Apache Airflow workflow orchestration. This skill covers DAG development, operators, sensors, task dependencies, dynamic workflows, XCom communication, scheduling patterns, and production deployment strategies.
When to Use This Skill
Use this skill when:
- Building and managing complex data pipelines with task dependencies
- Orchestrating ETL/ELT workflows across multiple systems
- Scheduling and monitoring batch processing jobs
- Coordinating multi-step data transformations
- Managing workflows with conditional execution and branching
- Implementing event-driven or asset-based workflows
- Deploying production-grade workflow automation
- Creating dynamic workflows that generate tasks programmatically
- Coordinating distributed task execution across clusters
- Building data engineering platforms with workflow orchestration
Core Concepts
What is Apache Airflow?
Apache Airflow is an open-source platform for programmatically authoring, scheduling, and monitoring workflows. It allows you to define workflows as Directed Acyclic Graphs (DAGs) using Python code, making complex workflow orchestration maintainable and version-controlled.
Key Principles:
- Dynamic: Workflows are defined in Python, enabling dynamic generation
- Extensible: Rich ecosystem of operators, sensors, and hooks
- Scalable: Can scale from single machine to large clusters
- Observable: Comprehensive UI for monitoring and troubleshooting
DAGs (Directed Acyclic Graphs)
A DAG is a collection of tasks organized to reflect their relationships and dependencies.
DAG Properties:
- dag_id: Unique identifier for the DAG
- start_date: When the DAG should start being scheduled
- schedule: How often to run (cron, timedelta, or asset-based)
- catchup: Whether to run missed intervals on DAG activation
- tags: Labels for organization and filtering
- default_args: Default parameters for all tasks in the DAG
DAG Definition Example:
Tasks and Operators
Tasks are the basic units of execution in Airflow. Operators are templates for creating tasks.
Common Operator Types:
- BashOperator: Execute bash commands
- PythonOperator: Execute Python functions
- EmailOperator: Send emails
- EmptyOperator: Placeholder/dummy tasks
- Custom Operators: User-defined operators for specific needs
Operator vs. Task:
- Operator: Template/class definition
- Task: Instantiation of an operator with specific parameters
Task Dependencies
Task dependencies define the execution order and workflow structure.
Dependency Operators:
>>: Sets downstream dependency (task1 >> task2)<<: Sets upstream dependency (task2 << task1)chain(): Sequential dependencies for multiple taskscross_downstream(): Many-to-many relationships
Dependency Examples:
Executors
Executors determine how and where tasks run.
Executor Types:
- SequentialExecutor: Single-threaded, local (default, not for production)
- LocalExecutor: Multi-threaded, single machine
- CeleryExecutor: Distributed execution using Celery
- KubernetesExecutor: Each task runs in a separate Kubernetes pod
- DaskExecutor: Distributed execution using Dask
Scheduler
The Airflow scheduler:
- Monitors all DAGs and their tasks
- Triggers task instances based on dependencies and schedules
- Submits tasks to executors for execution
- Handles retries and task state management
Starting the Scheduler:
DAG Development Patterns
Basic DAG Structure
Every DAG follows this structure:
Task Dependencies and Chains
Linear Chain:
Dynamic Chain:
Pairwise Chain:
Cross Downstream:
Branching and Conditional Execution
BranchPythonOperator:
Custom Branch Operator:
TaskGroups for Organization
TaskGroups help organize related tasks hierarchically:
Edge Labeling
Add labels to dependency edges for clarity:
LatestOnlyOperator
Skip tasks if not the latest DAG run:
Operators Deep Dive
BashOperator
Execute bash commands:
Complex Bash Command:
PythonOperator
Execute Python functions:
Traditional ETL with PythonOperator:
EmailOperator
Send email notifications:
EmptyOperator
Placeholder for workflow structure:
Custom Operators
Create reusable custom operators:
Sensors Deep Dive
Sensors are a special type of operator that wait for a certain condition to be met before proceeding.
ExternalTaskSensor
Wait for tasks in other DAGs:
Deferrable ExternalTaskSensor:
FileSensor
Wait for files to appear:
TimeDeltaSensor
Wait for a specific time period:
BigQuery Table Sensor
Wait for BigQuery table to exist:
Custom Sensors
Create custom sensors for specific conditions:
Deferrable Sensors
Deferrable sensors release worker slots while waiting:
XComs (Cross-Communication)
XComs enable task-to-task communication by storing and retrieving data.
Basic XCom Usage
Pushing to XCom:
Pulling from XCom:
XCom with TaskFlow API
TaskFlow API automatically manages XComs:
XCom Best Practices
Size Limitations:
- XComs are stored in the metadata database
- Keep XCom data small (< 1MB recommended)
- For large data, store in external systems and pass references
Example with External Storage:
XCom with Operators
Reading XCom in Templates:
XCom with EmailOperator:
Dynamic Workflows
Create tasks dynamically based on runtime conditions or external data.
Dynamic Task Generation with Loops
Dynamic Task Mapping
Map over task outputs to create dynamic parallel tasks:
Mapping with Classic Operators:
Task Group Mapping
Map over entire task groups:
Partial Parameters with Mapping
Mix static and dynamic parameters:
TaskFlow API
The modern way to write Airflow DAGs with automatic XCom handling and cleaner syntax.
Basic TaskFlow Example
Multiple Outputs
Return multiple values from tasks:
Mixing TaskFlow with Traditional Operators
Virtual Environment for Tasks
Isolate task dependencies:
Asset-Based Scheduling
Schedule DAGs based on data assets (formerly datasets) rather than time.
Producer-Consumer Pattern
Multiple Asset Dependencies
AND Logic (all assets must update):
OR Logic (any asset update triggers):
Complex Logic:
Asset Aliases
Use aliases for flexible asset references:
Accessing Asset Event Information
Scheduling Patterns
Cron Expressions
Timedelta Scheduling
Preset Schedules
Catchup and Backfilling
Catchup:
Manual Backfilling:
Production Patterns
Error Handling and Retries
Task-Level Retries:
DAG-Level Default Args:
Task Concurrency Control
Per-Task Concurrency:
DAG-Level Concurrency:
Idempotency
Make tasks idempotent for safe retries:
SLAs and Alerts
Task Callbacks
Docker Deployment
Docker Compose for Local Development:
Kubernetes Executor
KubernetesExecutor Configuration:
Pod Override for Specific Task:
Monitoring and Logging
Structured Logging:
StatsD Metrics:
Best Practices
DAG Design
- Keep DAGs Simple: Break complex workflows into multiple DAGs
- Use Descriptive Names: dag_id and task_id should be self-explanatory
- Idempotent Tasks: Tasks should produce same result when re-run
- Small XComs: Keep XCom data under 1MB
- External Storage: Use S3/GCS for large data, pass references
- Proper Dependencies: Model true dependencies, avoid unnecessary ones
- Error Handling: Use retries, callbacks, and proper error logging
- Resource Management: Set appropriate task concurrency limits
Code Organization
Testing DAGs
Unit Testing:
Integration Testing:
Performance Optimization
- Use Deferrable Operators: For sensors and long-running waits
- Dynamic Task Mapping: For parallel processing
- Appropriate Executor: Choose based on scale (Local, Celery, Kubernetes)
- Connection Pooling: Reuse database connections
- Task Parallelism: Set max_active_runs and concurrency appropriately
- Lazy Loading: Don't execute heavy logic at DAG parse time
- External Storage: Keep metadata database light
Security
- Secrets Management: Use Airflow Secrets Backend (not hardcoded)
- Connection Encryption: Use encrypted connections for databases
- RBAC: Enable role-based access control
- Audit Logging: Enable audit logs for compliance
- Network Isolation: Restrict worker network access
- Credential Rotation: Regularly rotate credentials
Configuration Management
Common Patterns and Examples
See EXAMPLES.md for 18+ detailed real-world examples including:
- ETL pipelines
- Machine learning workflows
- Data quality checks
- Multi-cloud orchestration
- Event-driven architectures
- Complex branching logic
- Dynamic task generation
- Asset-based scheduling
- Sensor patterns
- Error handling strategies
Troubleshooting
DAG Not Appearing in UI
- Check for Python syntax errors in DAG file
- Verify DAG file is in correct directory
- Check dag_id is unique
- Ensure schedule is not None if you expect it to run
- Check scheduler logs for import errors
Tasks Not Running
- Check task dependencies are correct
- Verify upstream tasks succeeded
- Check task concurrency limits
- Ensure executor has available slots
- Review task logs for errors
Performance Issues
- Reduce DAG complexity (break into multiple DAGs)
- Optimize SQL queries in tasks
- Use appropriate executor for scale
- Enable task parallelism
- Check for slow sensors (use deferrable mode)
- Monitor metadata database performance
Common Errors
Import Errors:
Circular Dependencies:
Large XComs:
Resources
- Official Documentation: https://airflow.apache.org/docs/
- Airflow GitHub: https://github.com/apache/airflow
- Astronomer Guides: https://docs.astronomer.io/learn
- Community Slack: https://apache-airflow.slack.com
- Stack Overflow: Tag
apache-airflow - Awesome Airflow: https://github.com/jghoman/awesome-apache-airflow
Skill Version: 1.0.0 Last Updated: January 2025 Apache Airflow Version: 2.7+ Skill Category: Data Engineering, Workflow Orchestration, Pipeline Management


