Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing

Orchestration & Reliable Pipelines

Progress0/14
x

What is a Pipeline

  • What is a data pipeline?12m
  • Scheduling and cron12m
  • Monitoring and alerting12m

DAG Mental Model

  • Why a scheduler12m
  • Tasks, edges, topological order14m
  • Operators, sensors, and XCom12m

Retries, Idempotency, Dedup

  • Retries in the DAG runner14m
  • Idempotent loads14m
  • Deduplicate before load12m

Backfills, Catchup, Intervals

  • Backfills and catchup12m
  • Data interval vs execution clock12m

CDC Concepts

  • CDC: inserts, updates, deletes12m
  • Apply CDC idempotently14m

Capstone

  • Capstone: a reliable night job16m
Back to track
  1. Learn
  2. Orchestration & Reliable Pipelines
  3. DAG Mental Model
  4. Tasks, edges, topological order

Lesson 5 of 14 · Theory first, then run it

Tasks, edges, topological order

orchestrationpythonbeginner14 min

Overview

A DAG is tasks plus directed edges with no cycles. The runner walks a topological order.

On this page7 sections›
  1. 1The picture
  2. 2Why this shape
  3. 3Walk the boxes
  4. 4A matching example
  5. 5Common beginner questions
  6. 6What comes next
  7. 7Practice

The picture

A simple three-step pipeline
extractloadtransform

Each arrow means 'must finish before.' Extract must finish before load. Load must finish before transform.

DAG stands for Directed Acyclic Graph. That sounds complicated, but it is a simple idea. A DAG is a collection of tasks connected by arrows that show which task must finish before another can start. 'Directed' means the arrows have a direction (A points to B, not the other way). 'Acyclic' means there are no loops (you cannot follow the arrows and end up back where you started).

Every orchestrator uses DAGs to represent pipelines. When you write an Airflow pipeline, you are literally defining a DAG. When you use Prefect or Dagster, you are defining the same structure with different syntax.

This lesson is a simulation: you will build a DAG as a Python dictionary and implement topological sort (the algorithm that finds a valid execution order). No Airflow parsing happens here, but the constraint 'load must wait for extract' is the same in every orchestrator.

Why this shape

Imagine a warehouse pipeline with three steps: extract data from an API, load it into a staging table, and transform it into a gold report. If you run transform before load finishes, you get stale or empty data. If you run load before extract finishes, you load nothing. The DAG encodes these rules so the orchestrator can enforce them automatically.

Without a DAG, you either run everything sequentially (slow and wasteful when tasks are independent) or you guess at parallelism and hope nothing breaks. A DAG gives you maximum safe parallelism: any tasks without a dependency relationship can run at the same time.

Walk the boxes

Think of a DAG like a recipe with steps that have prerequisites. You cannot frost a cake before you bake it, and you cannot bake it before you mix the batter. But you can preheat the oven while mixing the batter, because those two steps do not depend on each other.

In code, a DAG is typically represented as an adjacency list: a dictionary where each key is a task, and the value is a list of tasks that come after it (downstream tasks).

Trace the overnight failure

Arrows first. Then a timeout at 02:14. Predict what you click before the DAG moves.

A matching example

Run the example below in this tab. Read the input, follow the code, then check the output matches what you expect.

Python
# DAG as an adjacency list (downstream edges)
EDGES = {
    "extract": ["load"],
    "load": ["transform"],
    "transform": [],
}

# Compute indegree (number of incoming edges) for each task
def compute_indegree(edges):
    indegree = {task: 0 for task in edges}
    for task, downstreams in edges.items():
        for d in downstreams:
            indegree[d] += 1
    return indegree

print(compute_indegree(EDGES))
# {"extract": 0, "load": 1, "transform": 1}

Topological sort is the algorithm that finds a valid execution order. It works by repeatedly finding tasks with zero incoming edges (no unfinished dependencies), executing them, and then removing their outgoing edges.

Python
from collections import deque

def topo_order(edges):
    """Return tasks in a valid execution order using Kahn's algorithm."""
    indegree = {task: 0 for task in edges}
    for task, downstreams in edges.items():
        for d in downstreams:
            indegree[d] += 1

    queue = deque(t for t, deg in indegree.items() if deg == 0)
    order = []

    while queue:
        task = queue.popleft()
        order.append(task)
        for downstream in edges[task]:
            indegree[downstream] -= 1
            if indegree[downstream] == 0:
                queue.append(downstream)

    if len(order) != len(edges):
        raise ValueError("Cycle detected: not a valid DAG")
    return order

print(topo_order(EDGES))  # ["extract", "load", "transform"]

If there is a cycle (A depends on B, B depends on A), the algorithm detects it because some tasks will never reach zero indegree. This is why the 'acyclic' part matters: a cycle means the pipeline can never start.

TermMeaningExample
NodeA task in the pipelineextract, load, transform
EdgeA dependency relationshipextract -> load
IndegreeNumber of upstream dependenciesload has indegree 1
Topological orderValid execution sequenceextract, load, transform
CycleCircular dependency (invalid)A -> B -> A

Common beginner questions

Why is it called 'acyclic'?

Because cycles (loops) are not allowed. If task A depends on task B, and task B depends on task A, neither can ever start. The orchestrator would be stuck forever. DAGs forbid this by definition.

Can two tasks run at the same time?

Yes, if neither depends on the other. In a DAG where extract_orders and extract_customers both feed into a join task, the two extracts can run in parallel because they have no edge between them.

Cycles crash the scheduler

If you accidentally create a circular dependency, the orchestrator refuses to run the DAG at all. Airflow will show a 'cycle detected' error in the UI.

SQL joins are downstream of both tables

You already know from the SQL track that a JOIN needs both tables ready. In orchestration terms, the join task has two upstream dependencies.

What comes next

Now that you understand the graph structure, the next lesson introduces Airflow's vocabulary: operators, sensors, and XCom. These are the building blocks you use to fill each node in the DAG.

Say the main idea in one sentence without jargon. If that is hard, re-read Why this shape before moving on.

Say the main idea in one sentence without jargon. If that is hard, re-read Why this shape before moving on.

Say the main idea in one sentence without jargon. If that is hard, re-read Why this shape before moving on.

Say the main idea in one sentence without jargon. If that is hard, re-read Why this shape before moving on.

Say the main idea in one sentence without jargon. If that is hard, re-read Why this shape before moving on.

Practice

Implement topo_order on the EDGES dictionary. Expect the output: extract, load, transform.

Practicals · load into the editor

After you read the theory, run these in the pane on the right. They execute in this tab, no cluster.

Practice this

Same ideas as interview drills. These challenges open in the studio with a dataset and tests already set up.

  • Order DAG tasks so every dependency runs firstProduction ticket: A scheduler can't just run tasks in the order they were declared - it has to respect the dependency graph.Studiointermediatepython15 minPro
Rate:
Was this useful?
Why a schedulerOperators, sensors, and XCom