Phase 3: Data Pipelines & ETL

DAGs, tasks, dependencies & execution order

Intermediate ~3 min read
Think of it this way A friendly analogy. Read this if the technical version feels dense. Show Hide

Imagine you want to bake a delicious chocolate cake. You can't just throw all the ingredients into the oven at once, right? You need a plan, a step-by-step guide to make sure it turns out perfectly. In the world of computers and big data, when we want to do a bunch of jobs in a specific order, we create something like a super-smart recipe. We call this a "DAG" – which is short for "Directed Acyclic Graph". Don't worry about the big name! Just think of it as your ultimate cake recipe.

Each instruction in your cake recipe is like a "task". A task is one small, clear job you need to do. For example, "melt the butter," "mix the dry ingredients," or "bake the cake in the oven." These tasks are individual steps, and your recipe tells you exactly what to do for each one. But what's super important is the order! You can't "bake the cake" before you "mix the batter." This ordering is what we call a "dependency." One task depends on another finishing first. Your recipe makes sure the steps are "directed" – meaning you always move forward, never backward, and there are "no loops" – you won't accidentally get stuck in an endless cycle of mixing the same ingredients forever!

So, your cake recipe (the DAG) maps out every task and all the dependencies. It says: "First, prepare the pans. THEN, melt the butter. Meanwhile, you can also mix the dry ingredients. ONCE the butter is melted, and the dry ingredients are mixed, THEN combine them with the wet ingredients." This careful planning creates an "execution order" – the exact sequence the computer will follow to complete all the jobs, just like you'd follow your recipe to bake the cake. This way, everything gets done in the right order, and no step is missed or done too early.

This is really powerful because it means you can build complicated systems where many things need to happen automatically. For example, a big company might have a recipe (a DAG) for how to get all its daily sales numbers, clean them up, and then put them into a report for the bosses. Each step is a task, and the dependencies make sure the report is accurate and only uses data that's already been cleaned. So, when you build these "recipes" for computers, you're making sure complex work gets done reliably, step by step, just like baking a perfect cake every time!

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

python
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_data

How 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.