Chapter 16 · Orchestration, Dependencies, Scheduling, Backfills, SLAs, and Failure Recovery

DAGs, Dependencies, Data Availability Sensors/Checks, and Avoiding Clock-Only Scheduling

Model orchestration as dependency-aware execution over data readiness, not as a collection of clock-triggered scripts; make task state, source readiness, run identity, and certification evidence observable.

Intermediate → Advanced125–150 minutesDependency/readiness DAG labPython 3 stdlib · local/syntheticLast reviewed: September 2026

Learning outcomes

At 08:00 UTC AtlasMart expects its daily warehouse workflow to begin. On 2026-09-21 the ERP extract is not actually complete until 08:12. A cron-like scheduler can still fire at 08:00, but that tells the warehouse nothing about whether the source is ready. The first orchestration problem is therefore not “how do I schedule six tasks?” It is “what evidence allows each task to become runnable?”

01

Define a DAG, task dependency, data dependency, readiness check, schedule trigger, run ID, and partition before relying on orchestration jargon.

02

Explain why time-based scheduling can trigger work without proving that required source data is available or complete.

03

Model task states and dependency transitions so blocked, running, failed, retried, published, and certified states are distinguishable.

04

Use source manifests/readiness checks to prevent incomplete data from reaching transforms and downstream certification.

05

Separate vendor-neutral orchestration semantics from an optional product implementation such as Apache Airflow.

Chapter 16 continuity contract

Chapter 16 does not change the accepted Chapter 15 analytical state. The current certified AtlasMart sales target remains nine paid order-line facts, seven paid orders, eleven units, 740 USD paid GMV, 450 USD cost-at-sale, and 290 USD gross profit at committed source sequence 206. Orchestration adds run/task state, dependency evidence, partition manifests, SLO measurements, resource-pool labels, and operator notes around that data state. The historical backfill fixture deliberately starts from one corrupted 2026-09-20 publication (cost 435 USD instead of 425 USD) and repairs only that historical partition from immutable raw evidence; the current 2026-09-21 checksum must remain unchanged.

Execution and guarantee boundary

The mandatory labs are synthetic, local, and free. They use Python 3 standard-library modules and local JSON/filesystem state; generation-time validation ran with Python 3.13.5. Logical timestamps are simulated—no real waiting, cluster scheduler, queue, cloud warehouse, distributed lock, or resource manager is involved. A labeled current versus backfill resource pool proves orchestration intent in the fixture, not operating-system or cloud compute isolation. Apache Airflow is referenced only as an optional later-course implementation example and is not a prerequisite.

1. A DAG expresses allowed dependency order, not business correctness by itself

A directed acyclic graph (DAG) represents work as nodes connected by directed dependency edges, with no cycle that makes a task depend on itself through a chain. A task dependency says one task may run only after another reaches an accepted state. A data dependency says required input data has reached a defined readiness state. Those are related but not interchangeable.

A schedule trigger proposes when a workflow should be evaluated. A run ID identifies one orchestration attempt. A partition in this chapter is a bounded processing unit—here a logical daily AtlasMart publication such as 2026-09-21. It need not map one-to-one to a database engine's physical partition feature.

Task Depends on Success evidence Why the dependency exists
check_source — Source manifest is ready and complete enough to evaluate. A clock time does not prove the upstream data exists.
ingest_partition check_source Immutable input hash recorded. Do not transform an incomplete/unidentified receipt.
transform_partition ingest_partition Partition result produced from that input identity. Transformation must bind to a specific input version.
quality_gate transform_partition Counts/sums/checksums and business invariants pass. Task completion alone does not prove data correctness.
publish_partition quality_gate Atomic partition artifact becomes visible. Consumers must not see failed/uncertified output.
certify_partition publish_partition SLO/lineage/certification metadata recorded. Publication and consumer trust are distinct state transitions.

2. Clock time is a trigger; readiness is evidence

The synthetic source manifest says the 2026-09-21 input becomes ready at 08:12Z. The first run, RUN-20260921-0800-A, observes the source at 08:05Z. A clock-only design would start transformation because “08:00 has passed.” The dependency-aware design records check_source = WAITING and allows no downstream task to run.

Observation Bad interpretation Correct interpretation
Wall clock ≥ 08:00 The dataset is ready. Only the schedule trigger is eligible.
Source file exists The extract is complete. Existence needs completeness/manifest semantics.
Upstream job says SUCCESS All expected rows arrived. Job state is evidence about work, not necessarily data completeness.
Manifest says ready_at=08:12 and expected_rows=9 Transform immediately. Now compare actual input identity/count/hash to that contract.
dependency_gate.py
from datetime import datetime, timezoneready_at = datetime.fromisoformat("2026-09-21T08:12:00+00:00")observed_at = datetime.fromisoformat("2026-09-21T08:05:00+00:00")if observed_at < ready_at:    state = "BLOCKED_UPSTREAM"    runnable_downstream = []else:    state = "READY"    runnable_downstream = ["ingest_partition"]print(state)print(runnable_downstream)

Expected output is BLOCKED_UPSTREAM and an empty downstream list. This proves the local readiness rule is evaluated before work; it does not prove a production upstream manifest is trustworthy, atomic, or correctly generated.

3. Readiness checks need semantics, not just polling

A sensor/check is any mechanism that evaluates whether a dependency is ready. It might poll a manifest, query a control table, verify an object-generation marker, inspect a source sequence, or consume an event. The mechanism is less important than the contract: what exactly means ready, who owns that signal, how stale signals are detected, and what happens when the signal never arrives.

For AtlasMart the fixture combines readiness time with an expected row count. A production contract may additionally include extraction batch ID, source high-water mark, byte count, checksum, schema version, and completeness declaration. If a source publishes a “done” flag before its data files are durable, the sensor merely automates the wrong assumption.

4. Controlled failure: publish because every task process exited 0

Another common error is to define dependency success entirely in terms of process exit codes. Suppose transform_partition reads eight of nine expected rows and exits successfully. The orchestration graph is green but completeness is 88.89%. Chapter 16 therefore places a quality gate before publication and a separate certification step after publication.

The causal chain is important: scheduler → readiness check → identified input → transform → quality evidence → atomic publication → certification. A task graph is a control plane; the data product remains a separate state surface that must be measured.

5. Hands-on: inspect the AtlasMart DAG and upstream delay

Save the following as dag_gate.py and run it with your local Python 3. It uses only the standard library.

dag_gate.py
from graphlib import TopologicalSorterfrom datetime import datetimeDAG = {  "check_source": set(),  "ingest_partition": {"check_source"},  "transform_partition": {"ingest_partition"},  "quality_gate": {"transform_partition"},  "publish_partition": {"quality_gate"},  "certify_partition": {"publish_partition"},}print("order:", list(TopologicalSorter(DAG).static_order()))ready = datetime.fromisoformat("2026-09-21T08:12:00+00:00")observed = datetime.fromisoformat("2026-09-21T08:05:00+00:00")print("source_ready:", observed >= ready)print("decision:", "RUN" if observed >= ready else "BLOCK")

Expected signal: the topological order begins with check_source; source_ready is False; decision is BLOCK. Change only the observed time to 08:12 and rerun. That makes the source check eligible, but it still does not bypass later quality/certification gates.

Verification checklist

Record your Python version; confirm all six tasks appear once; confirm the 08:05 run blocks; confirm no downstream task is logically runnable; then confirm 08:12 changes only readiness, not the later data-quality requirements.

6. Production judgment and bridge

Choose readiness evidence according to the source's guarantees. Event-driven notifications can reduce polling latency but still require deduplication and durable state. Polling can be perfectly adequate when the freshness target allows it. The production question is not which orchestration style sounds modern; it is whether dependency state is unambiguous, observable, recoverable, and protected against stale or premature signals.

Lesson 2 assumes the source is ready and moves to the next failure surface: a task starts, then fails. Whether an orchestrator may retry automatically depends on timeout semantics, idempotency, partial side effects, and operator policy.

Knowledge check

Check your understanding

  1. What does a DAG edge prove?
  2. Why is 08:00 insufficient evidence for AtlasMart?
  3. What state does RUN-20260921-0800-A reach?
  4. Why keep quality_gate separate from transform_partition?
  5. Must a production implementation use Airflow?
Review the answers

1. It proves an allowed execution dependency in the control plane; it does not by itself prove the upstream data is complete or correct.

2. Because the source manifest is not ready until 08:12; time eligibility is different from data readiness.

3. BLOCKED_UPSTREAM; no publication is created.

4. A process can finish successfully while producing incomplete or semantically invalid data.

5. No. The DAG/readiness semantics are vendor-neutral; Airflow is only an optional later-course implementation example.

Authoritative references

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.