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.
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?”
Define a DAG, task dependency, data dependency, readiness check, schedule trigger, run ID, and partition before relying on orchestration jargon.
Explain why time-based scheduling can trigger work without proving that required source data is available or complete.
Model task states and dependency transitions so blocked, running, failed, retried, published, and certified states are distinguishable.
Use source manifests/readiness checks to prevent incomplete data from reaching transforms and downstream certification.
Separate vendor-neutral orchestration semantics from an optional product implementation such as Apache Airflow.
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.
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. |
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.
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.
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
- What does a DAG edge prove?
- Why is 08:00 insufficient evidence for AtlasMart?
- What state does RUN-20260921-0800-A reach?
- Why keep quality_gate separate from transform_partition?
- 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
- Python documentation — graphlibStandard-library topological ordering concepts used to explain dependency graphs without requiring an orchestrator product.
- Python documentation — os.replaceLocal atomic file-replacement primitive used by the fixture to demonstrate partition publication without partial output files.
- Python documentation — hashlibDeterministic SHA-256 evidence for replay/backfill comparisons in the local lab.
- Google SRE Book — Service Level ObjectivesFoundational distinction among service indicators/objectives and externally meaningful reliability goals.
- Apache Airflow documentation — DAGsOptional later-course example of a production orchestrator's DAG concept; no Airflow command or installation is required here.
- Apache Airflow documentation — BackfillOptional implementation reference for historical run concepts; Chapter 16 teaches the vendor-neutral semantics first.
- Kimball Group — Dimensional Modeling TechniquesBackground for the grains and dimensional facts whose correctness orchestration must preserve.