Roadmap To Be A Data Engineer / Lesson 07

Lesson 07 Orchestration Fundamentals §5 Orchestration & Scheduling About 5 min read

The Green Board

Every task succeeded and the data was still wrong. DAGs and check gates.

Sequel to Lesson 06. That one was nothing fired — a dead pipe, no alarm. This is the opposite and more dangerous failure: everything fired, every task went green, and the data was wrong anyway. The orchestrator did exactly what it promises. Its promise is just narrower than everyone assumes.

1The scene

Tuesday 06:00 CEST. The daily_repricing DAG wakes: four tasks in a line — extract_listings → enrich_brand → compute_prices → publish_prices — pull active listings, attach brand demand signal, compute markdowns, publish prices to the table the app reads.

At 05:55 an upstream job truncated raw_listings to rebuild it (a full refresh that normally finishes by 05:58; that morning it ran long). At 06:00 extract read the table mid-rebuild, found it empty, wrote an empty file, exited 0. enrich joined nothing → 0. compute marked down nothing → 0. publish ran its usual CREATE OR REPLACE TABLE dim_price AS SELECT … on an empty batch and replaced the entire live price table with zero rows.

Board all green by 06:01. Meanwhile 1,082 stale items that should have been marked down weren’t, and the price table was empty until a merchandiser noticed the storefront showing full prices on everything.

Dependency success = the upstream task’s process exited 0. Data success = the upstream task produced correct, complete data. An orchestrator guarantees the first. It has no opinion about the second.

2What an orchestrator actually does

Airflow / Dagster / Prefect / a cron of shell scripts — same category. It is the conductor, not a musician: plays no note itself, only decides what runs, in what order, when, and what happens on failure. Five jobs, only five:

# Job Meaning
1 Schedule start on a clock or a signal (“every day 06:00”, “when a file lands”)
2 Sequence run in dependency order — B waits for A. This is the DAG.
3 Retry task threw? try again N times, backing off
4 Backfill re-run a range of past dates (Lesson 04’s logical run date)
5 Observe surface state — which ran, how long, pass/fail — and alert

Not on the list: “check that the data is right.” Never was the conductor’s job. It sequences and watches exit codes. Whether SELECT returned 8,000 rows or 0 is invisible — both are “the task finished.”

3The DAG, precisely

DAG = Directed Acyclic Graph. Each word load-bearing:

  • Graph — tasks are nodes, dependencies are edges.
  • Directed — edges point one way; enrich depends on extract, never the reverse.
  • Acyclic — no loops. A ⇄ B never starts. The scheduler rejects it.

An edge means exactly one thing: do not start me until my upstream reports success. It does not mean “until my upstream produced good data.” Every edge in the scene did its job perfectly. That is the trap.

4The tells nobody was computing (both a single query away)

Same lesson as the Flatline Weekend: the failure was obvious, just unwatched.

Signal Normal day Incident day Orchestrator’s verdict
Exit code 0 0 SUCCESS
Rows published ~8,011 0 didn’t look
compute duration ~40.1 s ~0.2 s didn’t look
Markdowns applied 1,082 0 didn’t look
Markdown value moved €30,864.22 €0.00 didn’t look

Row count: 0 sits far outside the 14-day band [7,800 – 8,143] (median 8,011). Duration: a task that suddenly runs 200× faster almost always did nothing. Both trivially computable, neither computed. All figures produced deterministically by seed.py — they match to the cent.

5Making the wire hold: trigger rules

Every task’s default rule is all_success — run me only if all upstreams succeeded. That is the rule that walked the empty batch to publish. But the rule is a knob:

Trigger rule Task runs when… Use for
all_success every upstream succeeded (default) happy path
all_done every upstream finished, pass or fail cleanup, “always report”
one_success any one upstream succeeded “whichever source came back first”
all_failed every upstream failed fallback / alert-on-collapse
none_failed_min_one_success none failed and ≥1 ran branches that may skip, not fail

Add one node the scene lacked: a check task between compute and publish asserting “row count ≥ 60% of the 14-day median”. On the incident day it sees 0 vs a floor of 0.6 × 8,011 = 4,807, raises, turns red. Because publish depends on it under all_success, publish is skipped and yesterday’s good price table stays live. This is Lesson 06’s Write–Audit–Publish, wired into the DAG as a real edge.

6Schedules vs sensors — don’t run on a clock that lies

Root cause was timing: extract fired at 06:00 while the source was still rebuilding. A schedule answers “is it time?”; a sensor answers “is the input actually ready?” — it pokes until a file/partition/row-count threshold is met, then releases downstream.

Schedule: “It’s 06:00, go.” → ran into an empty table. Sensor: “Wait until raw_listings has > 4,000 rows for today, then go.”

Belt (sensor in front of extract) and braces (check task before publish): don’t start on stale input, don’t publish an output that fails its own audit.

7Ask your team

  1. When a task in our most important DAG goes green, what does that prove — a process exited 0, or the data is right? Where’s the gap?
  2. Do any DAGs run CREATE OR REPLACE / unconditional overwrite with no row-count check first? An empty batch there wipes the table.
  3. Which tasks start on a clock that could fire before their input exists? Which should be sensors?
  4. If a task’s runtime dropped 100× overnight, would anything notice — or would it look like a fast day?
  5. When a mid-pipeline check fails, does the bad batch get skipped (last good data stays) or published then fixed?

8Hands on (~45 min)

No Airflow, no Docker. A DAG is a dict of {task: [upstreams]} plus a runner that respects trigger rules — ~40 lines and you understand the engine.

  1. Run the deterministic seed.py (in the artifact). Confirm 8,000 listings, 1,082 changed, €30,864.22 of markdown.
  2. Write the 4-task pipeline as functions. Give extract a switch returning 0 rows. Run it, watch all four “succeed” and publish an empty table — reproduce the green board.
  3. Add a check task between compute and publish: fail if rows < 0.6 × 8011. Make publish depend on it under all_success. Re-run broken → publish gets skipped, table not wiped.
  4. Add a sensor in front of extract that waits for ≥ 4,000 source rows. Show the empty-source run now never produces a bad batch.
  5. Log each task’s row count and duration; print the run as a status board. That board plus the check is the whole idea.

Push as de-practice/07-orchestration with a README stating, for one real DAG you rely on: its schedule vs its true readiness signal, and the one check that would stop a bad batch publishing.

9Takeaway

An orchestrator is a conductor, not an inspector. It guarantees order and reports exit codes — it never promised the data was right. A green DAG means “the tasks ran in the order I told them to,” nothing more. Correctness is an edge you have to add: a check task that turns red on bad data, so the wire actually holds.

Smallest useful action: put one row-count check before every task that overwrites a table. Three lines of SQL, and the difference between a skipped run and a wiped table.

10Vocabulary

  • DAG — directed acyclic graph: tasks (nodes) + dependency edges (directed), no cycles.
  • Orchestrator — schedules, sequences, retries, backfills, observes tasks. Runs no business logic itself.
  • Dependency success vs data success — a task exiting 0 vs a task producing correct data. The orchestrator only knows the first.
  • Trigger rule — the condition under which a task runs given upstreams’ states (all_success, all_done, …).
  • Sensor — a task that waits for an external condition (file, partition, row count) before releasing downstream work.
  • Check / audit task — a node that asserts something about the data and fails the run when it’s wrong; the DAG’s Write–Audit–Publish.
  • Skipped — a task state that is neither success nor failure; downstreams under all_success don’t run, so the last good output stays live.
  • Silent-green failure — every task succeeds while the data is empty, partial or wrong. The failure mode this lesson exists to make loud.
Back to top