Chapter 16 · Orchestration, Dependencies, Scheduling, Backfills, SLAs, and Failure Recovery
Design an Operational Runbook for a Failed Daily Load Including Root Cause, Repair, Replay, and Communication
Turn a failed AtlasMart daily load into a repeatable incident workflow: detect, classify, contain, repair, replay, reconcile, certify, communicate, and document evidence without widening the blast radius.
Learning outcomes
An operator is paged because AtlasMart’s daily data product is late. There are two distinct incidents in the local scenario: the first run is blocked by an upstream delay, and the later transform fails once before succeeding on retry. Separately, reconciliation discovers a historical cost defect requiring a backfill. A useful runbook must tell the operator what to freeze, what evidence to collect, how to identify the affected state surface, when retry is safe, how to prove repair, and what consumers need to know.
Design a failed-load runbook that starts from containment and evidence preservation rather than blind reruns.
Classify root cause across source readiness, transformation, data quality/reconciliation, storage/warehouse, and publication/consumer surfaces.
Repair and replay the smallest safe partition while retaining run IDs, hashes, task traces, SLO measurements, and operator notes.
Communicate impact using consumer-visible states instead of internal task jargon alone.
Run the complete local acceptance harness and prove blocked scheduling, safe retry, historical backfill isolation, SLO attainment, and deterministic replay.
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. Incident response begins by identifying the failed state surface
A warehouse incident can originate in the source, extraction/ingestion, transformation, quality rules, orchestration control plane, physical warehouse/storage, semantic layer, or downstream report. Root cause and blast radius are different: a transform bug may affect one partition or many; a source delay can leave all warehouse tasks healthy but data stale; a semantic error can produce wrong dashboards while every table and task is technically available.
| Surface | Evidence to capture | Do not assume |
|---|---|---|
| Source readiness | Manifest/batch ID, expected counts, source high-water mark, ready timestamp. | That clock time or source job SUCCESS means completeness. |
| Orchestration | Run ID, task states, attempts, dependencies, timeouts, resource pool. | That FAILED means no side effect occurred. |
| Data state | Partition/hash, counts/sums, rejects, watermark, lineage. | That task SUCCESS means data correctness. |
| Publication/certification | Published version, certification flag/time, consumer access path. | That a written file/table is safe for consumers. |
| Consumer impact | Affected reports/data products, last good version, freshness/completeness impact. | That internal task names are meaningful to consumers. |
2. Runbook: root cause → contain → repair → replay → reconcile → communicate
- Contain. Freeze certification/publication for the affected partition or withdraw certification if a bad version is already visible.
- Preserve evidence. Capture run ID, partition, input hash/version, task trace, failure text, current publication hash, SLO state, and operator notes.
- Classify root cause. Decide whether the problem is upstream readiness, transient infrastructure, deterministic transform logic, data-contract/quality failure, resource exhaustion, or publication/consumer logic.
- Choose the smallest repair. Change data/code/config only where needed. If semantics change, version/approve the correction rather than silently editing history.
- Replay safely. Retry only tasks with known idempotency/compensation semantics. For historical repair, scope a backfill partition/range and isolate resources.
- Reconcile. Re-run grain-level counts, sums, keys, checksums, quality rules, history/incremental invariants, and SLO measurements.
- Certify and communicate. Publish/recertify only after evidence passes; state affected consumer data, time range, last-good/current version, and whether metrics were restated.
- Close with prevention. Record root cause, contributing conditions, monitoring gap, follow-up owner, and rollback/backout evidence.
3. Complete local acceptance lab
Save the script below as atlasmart_ch16.py. It uses
only Python 3 standard-library modules. Run it from a disposable
directory. It creates atlasmart_ch16_lab/, writes
synthetic immutable raw partitions, records a task-state graph
and run traces, blocks an upstream-late run, retries a failed
transform, publishes/certifies current data, repairs one
historical partition in a separate logical resource pool, and
proves replay determinism.
from __future__ import annotationsimport hashlibimport jsonimport osimport shutilimport sysfrom datetime import datetime, timezonefrom pathlib import PathCURRENT_ROWS = [ {"order_id":"O0999","line_no":1,"event_ts":"2026-09-17T13:00:00Z","customer_id":"C003","product_id":"P300","quantity":1,"extended_amount":80,"extended_cost":50,"status":"paid"}, {"order_id":"O1000","line_no":1,"event_ts":"2026-09-18T15:00:00Z","customer_id":"C001","product_id":"P200","quantity":1,"extended_amount":75,"extended_cost":45,"status":"paid"}, {"order_id":"O1001","line_no":1,"event_ts":"2026-09-18T09:15:00Z","customer_id":"C001","product_id":"P100","quantity":2,"extended_amount":100,"extended_cost":60,"status":"paid"}, {"order_id":"O1001","line_no":2,"event_ts":"2026-09-18T09:15:00Z","customer_id":"C001","product_id":"P200","quantity":1,"extended_amount":25,"extended_cost":15,"status":"paid"}, {"order_id":"O1002","line_no":1,"event_ts":"2026-09-18T11:30:00Z","customer_id":"C002","product_id":"P300","quantity":1,"extended_amount":195,"extended_cost":125,"status":"paid"}, {"order_id":"O1003","line_no":1,"event_ts":"2026-09-19T15:00:00Z","customer_id":"C001","product_id":"P400","quantity":1,"extended_amount":100,"extended_cost":60,"status":"paid"}, {"order_id":"O1003","line_no":2,"event_ts":"2026-09-19T15:00:00Z","customer_id":"C001","product_id":"P200","quantity":2,"extended_amount":55,"extended_cost":30,"status":"paid"}, {"order_id":"O1005","line_no":1,"event_ts":"2026-09-20T07:10:00Z","customer_id":"C004","product_id":"P100","quantity":1,"extended_amount":50,"extended_cost":30,"status":"paid"}, {"order_id":"O1006","line_no":1,"event_ts":"2026-09-21T08:00:00Z","customer_id":"C002","product_id":"P100","quantity":1,"extended_amount":60,"extended_cost":35,"status":"paid"},]BASELINE_ROWS = [ {"order_id":"O1000","line_no":1,"event_ts":"2026-09-18T15:00:00Z","customer_id":"C001","product_id":"P200","quantity":1,"extended_amount":75,"extended_cost":45,"status":"paid"}, {"order_id":"O1001","line_no":1,"event_ts":"2026-09-18T09:15:00Z","customer_id":"C001","product_id":"P100","quantity":2,"extended_amount":100,"extended_cost":60,"status":"paid"}, {"order_id":"O1001","line_no":2,"event_ts":"2026-09-18T09:15:00Z","customer_id":"C001","product_id":"P200","quantity":1,"extended_amount":25,"extended_cost":15,"status":"paid"}, {"order_id":"O1002","line_no":1,"event_ts":"2026-09-18T11:30:00Z","customer_id":"C002","product_id":"P300","quantity":1,"extended_amount":190,"extended_cost":125,"status":"paid"}, {"order_id":"O1003","line_no":1,"event_ts":"2026-09-19T15:00:00Z","customer_id":"C001","product_id":"P400","quantity":1,"extended_amount":100,"extended_cost":60,"status":"paid"}, {"order_id":"O1003","line_no":2,"event_ts":"2026-09-19T15:00:00Z","customer_id":"C001","product_id":"P200","quantity":2,"extended_amount":50,"extended_cost":30,"status":"paid"}, {"order_id":"O1005","line_no":1,"event_ts":"2026-09-20T07:10:00Z","customer_id":"C004","product_id":"P100","quantity":1,"extended_amount":50,"extended_cost":30,"status":"paid"}, {"order_id":"O1005","line_no":2,"event_ts":"2026-09-20T07:10:00Z","customer_id":"C004","product_id":"P400","quantity":1,"extended_amount":100,"extended_cost":60,"status":"paid"},]DAG = { "check_source": [], "ingest_partition": ["check_source"], "transform_partition": ["ingest_partition"], "quality_gate": ["transform_partition"], "publish_partition": ["quality_gate"], "certify_partition": ["publish_partition"],}SOURCE = { "2026-09-20": {"ready_at":"2026-09-20T08:00:00Z", "expected_rows":8}, "2026-09-21": {"ready_at":"2026-09-21T08:12:00Z", "expected_rows":9},}SLOS = {"freshness_minutes":30, "completeness_pct":100, "duration_minutes":10, "availability_minutes":30}def parse_ts(s: str) -> datetime: return datetime.fromisoformat(s.replace("Z", "+00:00")).astimezone(timezone.utc)def canonical_hash(obj) -> str: payload = json.dumps(obj, sort_keys=True, separators=(",", ":"), ensure_ascii=False) return hashlib.sha256(payload.encode("utf-8")).hexdigest()def controls(rows): paid = [r for r in rows if r["status"] == "paid"] return ( len(paid), len({r["order_id"] for r in paid}), sum(r["quantity"] for r in paid), sum(r["extended_amount"] for r in paid), sum(r["extended_cost"] for r in paid), sum(r["extended_amount"] - r["extended_cost"] for r in paid), )def write_json(path: Path, obj): path.parent.mkdir(parents=True, exist_ok=True) tmp = path.with_suffix(path.suffix + ".tmp") tmp.write_text(json.dumps(obj, indent=2, sort_keys=True), encoding="utf-8") os.replace(tmp, path)def append_jsonl(path: Path, obj): path.parent.mkdir(parents=True, exist_ok=True) with path.open("a", encoding="utf-8") as f: f.write(json.dumps(obj, sort_keys=True) + "\n")def log_task(root: Path, run_id: str, partition: str, task: str, state: str, logical_time: str, attempt: int, pool: str, detail: str = ""): append_jsonl(root / "task_events.jsonl", { "run_id": run_id, "partition": partition, "task": task, "state": state, "logical_time": logical_time, "attempt": attempt, "resource_pool": pool, "detail": detail, })def publish_payload(root: Path, partition: str, rows, run_id: str, certified_at: str, pool: str): payload = { "partition": partition, "grain": "one current paid AtlasMart order line per (order_id, line_no)", "rows": sorted(rows, key=lambda r: (r["order_id"], r["line_no"])), "controls": controls(rows), "run_id": run_id, "certified_at": certified_at, "resource_pool": pool, } payload["sha256"] = canonical_hash({"rows": payload["rows"], "controls": payload["controls"]}) write_json(root / "published" / f"partition={partition}.json", payload) return payloaddef run_partition(root: Path, partition: str, run_id: str, start_at: str, *, pool: str, fail_transform_once: bool = False, scheduled_at: str | None = None): manifest = SOURCE[partition] ready_at = parse_ts(manifest["ready_at"]) start = parse_ts(start_at) scheduled = parse_ts(scheduled_at or start_at) log_task(root, run_id, partition, "check_source", "RUNNING", start_at, 1, pool, "check manifest readiness and expected row count") if start < ready_at: log_task(root, run_id, partition, "check_source", "WAITING", start_at, 1, pool, f"source ready_at={manifest['ready_at']}") append_jsonl(root / "runs.jsonl", {"run_id":run_id,"partition":partition,"status":"BLOCKED_UPSTREAM","resource_pool":pool,"scheduled_at":scheduled_at or start_at,"observed_at":start_at}) return {"status":"BLOCKED_UPSTREAM", "run_id":run_id, "partition":partition} log_task(root, run_id, partition, "check_source", "SUCCESS", start_at, 1, pool, "source manifest ready") raw_path = root / "raw" / f"partition={partition}.json" raw = json.loads(raw_path.read_text(encoding="utf-8")) log_task(root, run_id, partition, "ingest_partition", "SUCCESS", start_at, 1, pool, f"raw_sha256={canonical_hash(raw)}") transform_attempts = 0 if fail_transform_once: transform_attempts += 1 log_task(root, run_id, partition, "transform_partition", "FAILED", "2026-09-21T08:14:00Z", transform_attempts, pool, "synthetic parser exception before publish; no target mutation") log_task(root, run_id, partition, "transform_partition", "RETRYING", "2026-09-21T08:15:00Z", transform_attempts, pool, "retry allowed because task writes partition atomically from immutable input") transform_attempts += 1 transformed = [dict(r) for r in raw["rows"]] log_task(root, run_id, partition, "transform_partition", "SUCCESS", "2026-09-21T08:16:00Z" if partition == "2026-09-21" else start_at, transform_attempts, pool, f"rows={len(transformed)}") actual = controls(transformed) expected_rows = manifest["expected_rows"] completeness = round(100 * len(transformed) / expected_rows, 2) if expected_rows else 100.0 if len(transformed) != expected_rows: log_task(root, run_id, partition, "quality_gate", "FAILED", start_at, 1, pool, f"expected_rows={expected_rows} actual_rows={len(transformed)}") raise AssertionError("completeness gate failed") log_task(root, run_id, partition, "quality_gate", "SUCCESS", "2026-09-21T08:17:00Z" if partition == "2026-09-21" else start_at, 1, pool, f"controls={actual} completeness_pct={completeness}") certified_at = "2026-09-21T08:18:00Z" if partition == "2026-09-21" else start_at log_task(root, run_id, partition, "publish_partition", "RUNNING", certified_at, 1, pool, "atomic temp-file replace") payload = publish_payload(root, partition, transformed, run_id, certified_at, pool) log_task(root, run_id, partition, "publish_partition", "SUCCESS", certified_at, 1, pool, f"sha256={payload['sha256']}") cert = parse_ts(certified_at) freshness_min = int((cert - ready_at).total_seconds() // 60) duration_min = int((cert - start).total_seconds() // 60) availability_min = int((cert - scheduled).total_seconds() // 60) slo = { "freshness_min": freshness_min, "completeness_pct": completeness, "duration_min": duration_min, "availability_min": availability_min, "freshness_met": freshness_min <= SLOS["freshness_minutes"], "completeness_met": completeness >= SLOS["completeness_pct"], "duration_met": duration_min <= SLOS["duration_minutes"], "availability_met": availability_min <= SLOS["availability_minutes"], } slo["all_met"] = all(slo[k] for k in ("freshness_met","completeness_met","duration_met","availability_met")) log_task(root, run_id, partition, "certify_partition", "SUCCESS", certified_at, 1, pool, json.dumps(slo, sort_keys=True)) append_jsonl(root / "runs.jsonl", {"run_id":run_id,"partition":partition,"status":"SUCCESS","resource_pool":pool,"slo":slo,"sha256":payload["sha256"]}) return {"status":"SUCCESS","run_id":run_id,"partition":partition,"controls":actual,"sha256":payload["sha256"],"slo":slo,"transform_attempts":transform_attempts}def seed(root: Path): write_json(root / "dag.json", {"dag_id":"atlasmart_daily","tasks":DAG,"current_pool":{"max_concurrency":2},"backfill_pool":{"max_concurrency":1}}) write_json(root / "raw" / "partition=2026-09-20.json", {"partition":"2026-09-20","rows":BASELINE_ROWS,"source":"immutable synthetic raw","source_seq_high":200}) write_json(root / "raw" / "partition=2026-09-21.json", {"partition":"2026-09-21","rows":CURRENT_ROWS,"source":"immutable synthetic raw","source_seq_high":206}) # Intentionally wrong historical publication: one cost value is corrupted. Backfill must repair only this partition. broken = [dict(r) for r in BASELINE_ROWS] for row in broken: if row["order_id"] == "O1005" and row["line_no"] == 1: row["extended_cost"] = 40 publish_payload(root, "2026-09-20", broken, "LEGACY-BROKEN", "2026-09-20T08:10:00Z", "legacy")def main(): root = Path(os.environ.get("ATLASMART_LAB_HOME", "atlasmart_ch16_lab")) if root.exists(): shutil.rmtree(root) root.mkdir(parents=True) seed(root) assert controls(BASELINE_ROWS) == (8, 5, 10, 690, 425, 265) assert controls(CURRENT_ROWS) == (9, 7, 11, 740, 450, 290) broken_hist = json.loads((root / "published" / "partition=2026-09-20.json").read_text(encoding="utf-8")) broken_controls = tuple(broken_hist["controls"]) assert broken_controls == (8, 5, 10, 690, 435, 255) # 1) Clock fires, but source is not ready. No downstream task is allowed to run. blocked = run_partition(root, "2026-09-21", "RUN-20260921-0800-A", "2026-09-21T08:05:00Z", pool="current", scheduled_at="2026-09-21T08:00:00Z") assert blocked["status"] == "BLOCKED_UPSTREAM" assert not (root / "published" / "partition=2026-09-21.json").exists() # 2) Data becomes ready. Transform fails once, then a safe retry publishes atomically. current = run_partition(root, "2026-09-21", "RUN-20260921-0812-B", "2026-09-21T08:12:00Z", pool="current", fail_transform_once=True, scheduled_at="2026-09-21T08:00:00Z") assert tuple(current["controls"]) == (9, 7, 11, 740, 450, 290) assert current["transform_attempts"] == 2 and current["slo"]["all_met"] current_hash_before_backfill = current["sha256"] # 3) Backfill one historical partition using a separate logical resource pool. append_jsonl(root / "operator_notes.jsonl", {"incident":"INC-016-01","action":"backfill","partition":"2026-09-20","reason":"legacy cost-at-sale corruption detected by reconciliation","resource_pool":"backfill","approval":"synthetic lab operator"}) backfill = run_partition(root, "2026-09-20", "BF-20260920-0900-R1", "2026-09-21T09:00:00Z", pool="backfill", scheduled_at="2026-09-21T09:00:00Z") assert tuple(backfill["controls"]) == (8, 5, 10, 690, 425, 265) current_after_backfill = json.loads((root / "published" / "partition=2026-09-21.json").read_text(encoding="utf-8")) assert current_after_backfill["sha256"] == current_hash_before_backfill # 4) Replay current partition from immutable raw input; output must be identical. replay = run_partition(root, "2026-09-21", "RUN-20260921-0910-C", "2026-09-21T09:10:00Z", pool="current", scheduled_at="2026-09-21T09:10:00Z") assert replay["sha256"] == current_hash_before_backfill task_events = [json.loads(line) for line in (root / "task_events.jsonl").read_text(encoding="utf-8").splitlines() if line.strip()] transform_trace = [e for e in task_events if e["run_id"] == "RUN-20260921-0812-B" and e["task"] == "transform_partition"] assert [e["state"] for e in transform_trace] == ["FAILED", "RETRYING", "SUCCESS"] runbook = """# AtlasMart failed daily load runbook\n\n1. Freeze publication/certification for the affected partition; do not advance consumers on task completion alone.\n2. Identify whether the failure is source readiness, transform logic, quality/reconciliation, warehouse/storage, or publication.\n3. Preserve run ID, partition, immutable input hash, task trace, and operator notes before changing state.\n4. Repair the smallest scoped cause. Retry only idempotent tasks; otherwise restore/compensate first.\n5. Replay the affected partition in the appropriate resource pool; do not rebuild unrelated history by default.\n6. Re-run counts/sums/checksums and freshness/completeness SLO checks.\n7. Certify/publish only after evidence matches the contract; communicate impact and resolution to downstream owners.\n8. Record root cause, blast radius, prevention action, and rollback/backout evidence.\n""" (root / "runbook.md").write_text(runbook, encoding="utf-8") summary = { "python": sys.version.split()[0], "dag": DAG, "blocked_run": blocked, "current_run": current, "historical_before": broken_controls, "historical_after": backfill["controls"], "current_hash_unchanged_after_backfill": current_after_backfill["sha256"] == current_hash_before_backfill, "current_hash_unchanged_after_replay": replay["sha256"] == current_hash_before_backfill, "transform_retry_states": [e["state"] for e in transform_trace], "resource_pools": {"current":{"max_concurrency":2},"backfill":{"max_concurrency":1}}, "slo_policy": SLOS, } write_json(root / "acceptance_summary.json", summary) print("blocked run:", blocked["run_id"], "status:", blocked["status"]) print("current controls:", tuple(current["controls"])) print("current SLOs:", f"freshness_min={current['slo']['freshness_min']}", f"completeness_pct={int(current['slo']['completeness_pct'])}", f"duration_min={current['slo']['duration_min']}", f"availability_min={current['slo']['availability_min']}", f"all_met={current['slo']['all_met']}") print("transform retry states:", [e["state"] for e in transform_trace]) print("historical before:", broken_controls) print("historical after:", tuple(backfill["controls"])) print("current hash unchanged after backfill:", current_after_backfill["sha256"] == current_hash_before_backfill) print("current hash unchanged after replay:", replay["sha256"] == current_hash_before_backfill) print("resource pools: current=2 backfill=1") print("current_sha256:", current_hash_before_backfill) print("cleanup: remove", root)if __name__ == "__main__": main()
Run:
python atlasmart_ch16.py# Windows PowerShell works with the same Python command when python is on PATH.
Expected acceptance signals
blocked run: RUN-20260921-0800-A status: BLOCKED_UPSTREAMcurrent controls: (9, 7, 11, 740, 450, 290)current SLOs: freshness_min=6 completeness_pct=100 duration_min=6 availability_min=18 all_met=Truetransform retry states: ['FAILED', 'RETRYING', 'SUCCESS']historical before: (8, 5, 10, 690, 435, 255)historical after: (8, 5, 10, 690, 425, 265)current hash unchanged after backfill: Truecurrent hash unchanged after replay: Trueresource pools: current=2 backfill=1current_sha256: 02715e0034dcd82288d5db9ad711c74dc4253e3617acc7c2f0c6f9344f0ef332
The SHA-256 covers the current partition's rows and controls in this fixture. If you edit the deterministic source rows or canonical serialization, the hash should change; that is evidence, not an externally meaningful data identifier.
4. What each artifact proves—and does not prove
| Artifact/evidence | Proves in this fixture | Does not prove |
|---|---|---|
| dag.json | Explicit task dependencies and separate logical pool policy. | A real scheduler enforces them. |
| RUN-A BLOCKED_UPSTREAM | Readiness gate prevents premature downstream work. | The source manifest generator is correct. |
| FAILED→RETRYING→SUCCESS trace | Retry policy is observable for one idempotent transform. | All pipeline side effects are retry-safe. |
| Current controls 9/7/11/740/450/290 | Publication reconciles to Chapter 15 accepted current state. | Every historical/SCD/semantic consumer is correct. |
| Historical 435→425 cost repair | Backfill repaired the seeded historical corruption. | The correction scope is always one partition in production. |
| Current hash unchanged after backfill | Protected current partition did not change. | No shared external resource was contended. |
| SLO all_met=True | Fixture indicator calculations meet declared thresholds. | Future reliability or contractual SLA compliance. |
| runbook.md/operator_notes.jsonl | Response procedure and one operator action are recorded. | Human review quality or organizational governance. |
5. Communication template as operational data
A useful incident update should answer: what data product is affected, what time/partition is impacted, whether data is late/incomplete/wrong/unavailable, what last good state consumers can use, whether metrics may be restated, what repair is underway, and when the next evidence-based update will occur. Avoid messages such as “task 4 failed” without consumer meaning.
For the local delayed run, an example factual update would say: “The 2026-09-21 AtlasMart certified sales publication is delayed because the ERP source manifest is not yet ready. No incomplete partition has been published. The last certified state remains available. The workflow will resume only after the source readiness contract passes.”
6. Failure injection and rollback boundaries
Before testing recovery in production, define blast radius. Use synthetic/non-production partitions where possible. If testing a current production path, have a last-good version, restore/swap procedure, consumer freeze mechanism, and explicit authorization. Do not test backfill isolation by deliberately exhausting production compute.
A rollback must restore mutually consistent state surfaces. Restoring a table without the corresponding certification metadata, semantic cache, lineage/run state, or watermark can create a new inconsistency. Chapter 15 covered cursor/target atomicity; Chapter 16 extends that discipline to orchestration and publication.
7. Cleanup/reset
# Linux/macOS/Git Bashrm -rf atlasmart_ch16_labrm -f atlasmart_ch16.py# PowerShell# Remove-Item -Recurse -Force .\atlasmart_ch16_lab# Remove-Item -Force .\atlasmart_ch16.py
Cleanup removes only the synthetic fixture. In production, run/task history, manifests, operator notes, lineage, and certification records are audit/recovery evidence; apply a retention policy rather than deleting them as temporary clutter.
8. Bridge to Chapter 17
Chapters 14–16 establish layered, incremental, recoverable operational semantics. Chapter 17 shifts to physical warehouse design: schemas, tables, constraints, partitioning, clustering/sort/distribution choices, and environment/domain boundaries. The orchestration lesson remains relevant: a physical layout must support safe current loads and bounded backfills without changing the logical grain or metric meaning.
Knowledge check
Check your understanding
- What is the first action in the runbook?
- What does RUN-A prove?
- What state is accepted for Chapter 15 continuity?
- What proves historical backfill isolation in the fixture?
- Why communicate consumer state rather than task state alone?
Review the answers
1. Contain the affected publication/certification state and preserve evidence rather than blindly rerunning.
2. The workflow can block on source readiness rather than start downstream work merely because the clock fired.
3. 9 paid lines, 7 orders, 11 units, 740 GMV, 450 cost, 290 gross profit.
4. The current 2026-09-21 SHA-256 remains unchanged before/after repairing the 2026-09-20 partition.
5. Consumers need to know whether data is late, incomplete, wrong, unavailable, or restated—not which internal task name failed.
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.