Data
Lesson 7 of 8About 4 min readSuggest an edit

Orchestrating pipelines

A data platform runs hundreds of jobs: extract, load, transform, test, publish. Each depends on others finishing first. An orchestrator runs them in the right order, on time, and tells you when something breaks.

DAGs and dependencies

Pipelines are modelled as a DAG, a directed acyclic graph. Each node is a task, and each edge says “this must finish before that starts”. Acyclic means no loops, so there is always a valid order to run things in.

extract_orders ──┐
                 ├── build_revenue ── test_revenue ── publish
extract_fx ──────┘

Explicit dependencies let the orchestrator run independent tasks in parallel, skip tasks whose inputs failed, and rerun only part of the graph.

Schedulers

Common open-source orchestrators take slightly different views:

  • Airflow is task-centric: you define DAGs of tasks in Python, and it schedules and runs them.
  • Dagster is asset-centric: you define the tables and files you want to exist, and it works out what to run to produce them.
  • Prefect turns ordinary Python functions into flows and tasks, with scheduling and retries added.

The concepts below apply to all of them.

Parameterise by run date

A scheduled run should process a specific slice of time, passed in as a parameter, not “whatever is new since now”. Airflow, for example, gives each run a logical date and a data interval, and exposes the date to templates as {{ ds }}.

def load_daily_orders(conn, run_date: str) -> None:
    with conn.transaction():
        conn.execute(
            "DELETE FROM orders_daily WHERE day = %s",
            (run_date,),
        )
        conn.execute(
            """INSERT INTO orders_daily (day, orders, revenue)
               SELECT %s, COUNT(*), SUM(total)
               FROM orders
               WHERE created_at >= %s::date
                 AND created_at < %s::date + 1""",
            (run_date, run_date, run_date),
        )

This task is idempotent: running it twice for the same date gives the same result, because it replaces that day’s data rather than appending. It never reads the wall clock, so running it today for last March produces what last March’s run should have produced.

Backfills

A backfill runs a pipeline for past dates: after fixing a bug, adding a table, or recovering from an outage. With date-parameterised, idempotent tasks, a backfill is the normal task run for a range of dates. Limit how many backfill runs execute at once, so you don’t overload the source or the warehouse.

Retries

Networks drop, APIs rate-limit, and warehouses queue. Configure retries with a delay, ideally increasing between attempts, for tasks that call external systems. Retries are only safe when tasks are idempotent; a retried append creates duplicates. Set a timeout on every task, so a hung job fails visibly instead of blocking everything downstream.

Schedules, sensors and events

  • A schedule runs at fixed times, such as 02:00 daily. Simple, but it assumes inputs are ready by then.
  • A sensor waits for a condition, such as a file landing in a bucket or a partition appearing, before starting downstream work.
  • Event or data-aware triggers start a pipeline when an upstream dataset is updated, so consumers run as soon as producers finish.

Prefer triggering on the data you depend on over guessing a time with a safety margin.

Hidden dependencies

The most painful failures come from dependencies the orchestrator does not know about: a job that reads a table another DAG builds, relying on the fact that it “usually finishes by 3am”. One day it doesn’t, and you publish partial data without any error.

Declare cross-pipeline dependencies explicitly with sensors, dataset triggers or a single graph, and make each task read only the inputs it declares.

Observability

Track for every pipeline:

  • run status and duration, with alerts on failure and on runs far slower than usual;
  • freshness of the outputs, which catches pipelines that silently stopped being scheduled;
  • rows read and written per task, so an empty load stands out;
  • logs linked from the alert, with the run date and parameters.

Checklist

  • Every task takes its run date as a parameter and never reads the current time.
  • Every task is idempotent and safe to retry.
  • Dependencies are declared, including those across pipelines.
  • Retries, timeouts and backfill concurrency are configured deliberately.
  • Failures and stale outputs alert the owning team.

Next: Data governance and privacy

Classifying personal data, access control, masking and pseudonymisation, retention, lineage and practical privacy principles.