The orchestration is an Airflow 3.3.2 129 DAG of nine BashOperator tasks, each calling one short script in demos/ch08/capstone/steps/, so you can also run any step by hand. Note the trailing space after each command: without it, a bash_command ending in .sh is read as a template file name (question 6 of Test Yourself!).
"""BookNest's data platform end to end: one run from the XML catalog to a published mart."""
from datetime import datetime, timezone
from airflow.providers.standard.operators.bash import BashOperator
from airflow.sdk import DAG, chain
STEPS = "/mnt/d/Books/Data Engineering/demos/ch08/capstone/steps"
with DAG("booknest_platform", schedule=None, max_active_tasks=1,
start_date=datetime(2026, 7, 1, tzinfo=timezone.utc)):
step = {name: BashOperator(task_id=name, bash_command=f"bash '{STEPS}/{name}.sh' ")
for name in ["catalog", "generate", "load_lake", "stream_events",
"events_contract", "build_mart", "pii_gate", "reconcile", "publish"]}
chain(step["catalog"], step["generate"], # Chapters 2 and 3: sources
step["load_lake"], # Chapter 5 Spark -> Iceberg (8.4)
step["stream_events"], # Chapter 6 Kafka -> Iceberg (8.10)
step["events_contract"], # gate: data contract (8.12)
step["build_mart"], # Chapter 4 dbt on Trino, with tests
[step["pii_gate"], step["reconcile"]], # gates: PII (8.13), totals
step["publish"]) # Iceberg tag = publish (8.6)A full Airflow stack holds about 1.5 GiB (Choosing an Executor), so run_platform.sh runs the DAG with airflow dags test and a SQLite 4,756 metadata database (dag.test() and Local DAG Runs): the same task code in one process. It prints each step's report with the time it was logged (UTC):
07:25:45 booknest-catalog.xml validates
07:25:46 catalog: 6 books -> landing/books.jsonl
07:25:54 generate: 100000 orders, 390737 events, checksums match
07:26:28 spark: 6 books, 100000 orders, 137944 lines, gross 3303427.30
07:27:25 kafka: 390737 events on booknest.order-events -> 390737 rows in platform.order_events
07:27:32 contract: 18 checks: {'passed': 18}, contract result: passed
07:28:23 dbt: 3 models, 7 tests passed, 0 failed
07:28:26 reconcile: non-cancelled gross
07:28:26 source files, Python 3,303,427.30
07:28:26 Iceberg tables, Spark 3,303,427.30
07:28:26 mart, Trino 3,303,427.30
07:28:26 mart, DuckDB 3,303,427.30
07:28:46 pii: 0 of 15 published columns tagged
07:28:47 publish: tag published -> snapshot 1490863382816522832 (HTTP 200)
07:28:49 publish: 2911 rows, gross 3303427.30 readable as of tag published
run_duration=190.290075, state=successThe platform ran in just over three minutes from an empty namespace, the heavy services taking turns: Spark 129 lived only inside load_lake, stream_events started and removed the l1-kafka broker and l1-connect worker, and l1-trino ran from build_mart to publish. Those three steps took three quarters of the run (on this shared 4-CPU host an earlier run beside two other Kafka 129 brokers took 274 seconds).
The reconciliation is the real test: four independent paths agree to the cent. The mart's path is the longest: each order's status comes from its latest Kafka event, its lines join books from the XML catalog, and dbt 37,942 aggregates them in Trino 403,499 . The test status_matches_orders confirms row by row that the stream and the order service agree on all 100,000 orders, including the four pending ones with only an order_placed event.