In Airflow, the cornerstone of building data pipelines is the DAG, which stands for Directed Acyclic Graph. Think of a DAG as the blueprint for your entire workflow; it's a collection of all the tasks you want to run, organized in a specific sequence. "Directed" means the flow of work moves in one direction, from start to finish, while "Acyclic" means there are no loops, preventing tasks from running indefinitely. Each individual piece of work within this blueprint is called a task. A task represents an atomic unit of work, like fetching data from an API, cleaning a dataset, or loading results into a database. These tasks are typically defined using Airflow operators (e.g., BashOperator, PythonOperator).
What makes a DAG powerful is not just the list of tasks, but how they relate to each other through dependencies. Dependencies define the precise order in which tasks must execute. For example, you might define that a "transform data" task can only run after a "fetch data" task has successfully completed. In Airflow, you express these relationships using simple syntax like task_a >> task_b (meaning task_b runs after task_a) or task_c << task_d (meaning task_c runs before task_d). This structure establishes the execution order and ensures your data pipeline progresses logically, preventing errors and maintaining data integrity. Tasks without dependencies on each other can even run in parallel, maximizing efficiency.
Ultimately, your DAG definition is a Python script that Airflow reads, parses, and then uses to orchestrate your pipeline runs. When you schedule a DAG, Airflow’s scheduler will execute the tasks according to their defined dependencies. It manages retries, monitors progress, and ensures that the execution order you’ve specified is respected, even across failures. Understanding how to define DAGs, break down work into tasks, and establish their dependencies is fundamental to designing robust and efficient data pipelines with Airflow.
Key Takeaways
- A DAG (Directed Acyclic Graph) is the blueprint for your entire data pipeline workflow in Airflow.
- Tasks are individual, atomic units of work within a DAG, often defined using Airflow operators.
- Dependencies define the sequential relationships and execution order between tasks, ensuring logical progression.
- Airflow uses the DAG's definition of tasks and dependencies to orchestrate and manage pipeline runs, including parallel execution and error handling.
Code Example
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime
with DAG(
dag_id='simple_pipeline_example',
start_date=datetime(2023, 1, 1),
schedule_interval=None,
catchup=False,
tags=['pipeline_basics'],
) as dag:
extract_data = BashOperator(
task_id='extract_data_from_source',
bash_command='echo "Fetching raw data..."',
)
transform_data = BashOperator(
task_id='clean_and_transform_data',
bash_command='echo "Processing and cleaning data..."',
)
load_data = BashOperator(
task_id='load_to_data_warehouse',
bash_command='echo "Loading transformed data..."',
)
# Define dependencies: Extract -> Transform -> Load
extract_data >> transform_data >> load_dataHow this code works
This code defines a simple data processing pipeline using Airflow, demonstrating how to sequence different steps, like an Extract-Transform-Load (ETL) workflow. It begins by creating a DAG object, which represents the entire pipeline. The dag_id gives it a unique name, and start_date sets when it becomes active. Crucially, schedule_interval=None means this pipeline won't run automatically on a schedule, while catchup=False prevents Airflow from attempting to run past, unscheduled instances between the start_date and now—a common beginner gotcha that can lead to unexpected runs.
Within the DAG's context, individual steps are defined as tasks using BashOperators. Each BashOperator has a task_id for identification and a bash_command that specifies the action to perform, in this case, simple print statements simulating data operations like fetching, processing, and loading. The lines creating extract_data, transform_data, and load_data set up these distinct steps. Finally, the extract_data >> transform_data >> load_data line explicitly defines the execution order, ensuring that data is extracted before it's transformed, and transformed before it's loaded into its destination.