Running the Platform

Running BookNest's Platform End to End on This Host

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_platform.py: the platform as one Airflow DAGPython
"""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):

Output of 67
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=success

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