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.

Intermediate → Advanced150–180 minutesOperational recovery/runbook labPython 3 stdlib · local/syntheticLast reviewed: September 2026

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.

01

Design a failed-load runbook that starts from containment and evidence preservation rather than blind reruns.

02

Classify root cause across source readiness, transformation, data quality/reconciliation, storage/warehouse, and publication/consumer surfaces.

03

Repair and replay the smallest safe partition while retaining run IDs, hashes, task traces, SLO measurements, and operator notes.

04

Communicate impact using consumer-visible states instead of internal task jargon alone.

05

Run the complete local acceptance harness and prove blocked scheduling, safe retry, historical backfill isolation, SLO attainment, and deterministic replay.

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. 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

  1. Contain. Freeze certification/publication for the affected partition or withdraw certification if a bad version is already visible.
  2. Preserve evidence. Capture run ID, partition, input hash/version, task trace, failure text, current publication hash, SLO state, and operator notes.
  3. 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.
  4. Choose the smallest repair. Change data/code/config only where needed. If semantics change, version/approve the correction rather than silently editing history.
  5. Replay safely. Retry only tasks with known idempotency/compensation semantics. For historical repair, scope a backfill partition/range and isolate resources.
  6. Reconcile. Re-run grain-level counts, sums, keys, checksums, quality rules, history/incremental invariants, and SLO measurements.
  7. Certify and communicate. Publish/recertify only after evidence passes; state affected consumer data, time range, last-good/current version, and whether metrics were restated.
  8. 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.

atlasmart_ch16.py
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:

run.sh
python atlasmart_ch16.py# Windows PowerShell works with the same Python command when python is on PATH.

Expected acceptance signals

expected-output.txt
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

cleanup.sh
# 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

  1. What is the first action in the runbook?
  2. What does RUN-A prove?
  3. What state is accepted for Chapter 15 continuity?
  4. What proves historical backfill isolation in the fixture?
  5. 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

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.