How to Build Data Pipelines with Apache Airflow Orchestration: A Complete Guide
Apache Airflow orchestrates complex data pipelines as Directed Acyclic Graphs (DAGs) using Python-based configuration, explicit task dependencies, and built-in monitoring to ensure reliable, scheduled data workflows.
The DataExpert-io/data-engineer-handbook repository identifies Airflow as a core orchestration tool for modern data engineering, listing it among essential workflow managers in the README.md orchestration section. While the repository focuses on conceptual foundations and pipeline maintenance best practices rather than hands-on Airflow tutorials, its intermediate bootcamp materials provide the architectural principles necessary to build production-grade pipelines.
Understanding Airflow's Core Architecture
Apache Airflow structures data workflows as DAGs (Directed Acyclic Graphs)—Python files that define tasks and their execution order. Each DAG resides in the dags/ directory of your Airflow deployment and consists of discrete tasks connected by dependencies. Tasks use operators such as PostgresOperator for SQL execution, PythonOperator for custom logic, or hooks like S3Hook for cloud storage interactions. According to the handbook's pipeline maintenance philosophy, this explicit dependency management and task isolation form the foundation of reliable data engineering.
Step-by-Step Guide to Building Airflow Pipelines
1. Define the Pipeline Scope
Before writing code, identify your source systems (databases, APIs, file stores), transformation logic (cleaning, enrichment, aggregation), and destination targets (data warehouses, data lakes). The handbook emphasizes that clear scoping prevents downstream maintenance issues documented in intermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md.
2. Create the DAG File
Place your pipeline definition in a Python file within the dags/ folder. Import Airflow's core classes and configure default arguments for ownership, retries, and alerting:
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.amazon.aws.operators.s3 import S3CreateObjectOperator
from airflow.operators.python import PythonOperator
default_args = {
"owner": "data-eng-team",
"depends_on_past": False,
"email_on_failure": True,
"email": ["alerts@example.com"],
"retries": 2,
"retry_delay": timedelta(minutes=5),
}
3. Implement Idempotent Tasks
Ensure each task can safely re-run without duplicating data. Use UPSERT logic, checkpoint files, or overwrite patterns. This aligns with the handbook's pipeline reliability standards, which stress that idempotency is critical for failure recovery and backfilling operations.
4. Configure Scheduling and Dependencies
Define your DAG's schedule_interval using cron expressions or presets like @daily. Establish execution order using the >> operator or set_upstream()/set_downstream() methods:
with DAG(
dag_id="example_etl",
default_args=default_args,
description="Simple ETL pipeline with Airflow",
schedule_interval="@daily",
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["etl", "airflow"],
) as dag:
# Task definitions here
extract >> transform >> load
5. Add Monitoring and Alerting
Leverage Airflow's web UI for task status, log inspection, and SLA monitoring. For production environments, integrate external observability tools like Prometheus and Grafana. The handbook references data quality tools such as Great Expectations and Soda in its README data quality section, suggesting these complement Airflow's native monitoring for comprehensive pipeline health checks.
6. Version Control and CI/CD
Store DAG files in Git repositories following the handbook's example structure. Implement CI pipelines to lint Python code using flake8 and black, and run unit tests before deployment to prevent broken DAGs from reaching production.
7. Deploy and Scale
Deploy Airflow on Kubernetes via Helm charts or use managed services like Astronomer or Google Cloud Composer. Select your executor type based on workload:
- LocalExecutor: Single-machine development
- CeleryExecutor: Distributed worker queues
- KubernetesExecutor: Dynamic pod creation for elastic scaling
8. Maintain and Document
Create run-books documenting ownership, failure scenarios, and recovery procedures. The intermediate-bootcamp/materials/6-data-pipeline-maintenance/homework/homework.md file provides templates for pipeline documentation standards, emphasizing that operational knowledge must outlive original developer assignments.
Production-Ready Airflow DAG Example
Below is a complete, self-contained ETL pipeline that extracts from PostgreSQL, transforms data using Python, and loads to Amazon S3. Place this file in your dags/ directory:
# dags/example_etl.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.amazon.aws.operators.s3 import S3CreateObjectOperator
from airflow.operators.python import PythonOperator
default_args = {
"owner": "data-eng-team",
"depends_on_past": False,
"email_on_failure": True,
"email": ["alerts@example.com"],
"retries": 2,
"retry_delay": timedelta(minutes=5),
}
def transform_data(**context):
# Example idempotent transformation – replace with real logic
raw = context["ti"].xcom_pull(task_ids="extract")
transformed = raw.upper() # placeholder
return transformed
with DAG(
dag_id="example_etl",
default_args=default_args,
description="Simple ETL pipeline with Airflow",
schedule_interval="@daily",
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["etl", "airflow"],
) as dag:
extract = PostgresOperator(
task_id="extract",
postgres_conn_id="source_pg",
sql="SELECT * FROM public.raw_events;",
)
transform = PythonOperator(
task_id="transform",
python_callable=transform_data,
)
load = S3CreateObjectOperator(
task_id="load",
s3_bucket="data-warehouse",
s3_key="events/{{ ds }}/transformed.csv",
data="{{ ti.xcom_pull(task_ids='transform') }}",
aws_conn_id="aws_default",
)
# Define dependencies
extract >> transform >> load
Key Resources from the Data Engineer Handbook
The DataExpert-io/data-engineer-handbook repository provides conceptual backing for these implementation steps:
README.md(Orchestration section): Lists Airflow among essential orchestration tools and provides the high-level tooling overview for data engineering stacksintermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md: Documents best practices for pipeline ownership, failure handling, and maintenance proceduresintermediate-bootcamp/materials/6-data-pipeline-maintenance/homework/homework.md: Contains exercises for designing documented pipelines with proper run-books and operational guidelinesbooks.md: References additional reading on pipeline theory and data engineering architecture
Summary
- Apache Airflow uses Python-defined DAGs to orchestrate complex data workflows with explicit task dependencies and scheduling
- Store DAG files in the
dags/directory and use built-in operators likePostgresOperatorandPythonOperatorfor task implementation - Implement idempotent tasks to ensure safe reruns and align with the handbook's reliability standards
- Configure monitoring through Airflow's UI and external tools like Prometheus, following the repository's data quality principles
- Use version control and CI/CD pipelines to lint and test code before deployment
- Choose appropriate executors (Local, Celery, or Kubernetes) based on production scaling requirements
- Document run-books and ownership according to
intermediate-bootcamp/materials/6-data-pipeline-maintenance/guidelines
Frequently Asked Questions
What is the difference between a DAG and a Task in Airflow?
A DAG (Directed Acyclic Graph) is the complete workflow definition—a Python file containing the entire pipeline structure, schedule, and default arguments. A Task is a single node within that DAG representing one unit of work, such as executing a SQL query or running a Python function. Tasks are instantiated using operators and connected via dependencies to form the DAG's execution graph.
How do I make Airflow tasks idempotent?
Design tasks to produce the same result whether run once or multiple times by implementing UPSERT logic for database writes, using deterministic primary keys, or writing to partition paths that include execution dates (like s3://bucket/data/{{ ds }}/). Avoid append-only operations that duplicate data on reruns, and use Airflow's {{ ds }} macros to create date-specific outputs that prevent collisions.
Where should I store my Airflow DAG files?
Store DAG files in a Git repository under a dags/ directory, then deploy them to Airflow's configured DAGs folder (typically /opt/airflow/dags in Docker deployments or synced via GitSync in Kubernetes). The DataExpert-io/data-engineer-handbook repository demonstrates this pattern by maintaining structured materials in version control, emphasizing that DAG code requires the same testing and review standards as application code.
Which executor should I use for production Airflow deployments?
Choose CeleryExecutor for traditional distributed processing with fixed worker pools, or KubernetesExecutor for dynamic, elastic scaling where each task runs in its own pod. Avoid LocalExecutor for production as it runs on a single machine and cannot distribute workload. Managed services like Google Cloud Composer or Astronomer abstract executor configuration while providing production-ready infrastructure.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →