The previous seven chapters each solved one problem. This capstone assembles them into one pipeline with clear layers, idempotent writes, and outputs a product team could actually use.
The layers
1raw JSON ──▶ bronze_playback typed, append-only, with lineage2 │3 ▼4 silver_sessions valid rows, sessionized, catalog joined5 │6 ▼7 gold_genre_daily + gold_concurrency (streaming)Silver: idempotent merge
This lab rereads the complete bronze history on each run. Validate it with the strict <1% quality gate, deduplicate the complete synthetic event payload before sessionization, and then upsert silver by a deterministic key. watch_seconds belongs in both the payload and key: two events that differ in their watch contribution are not duplicates.
1from delta.tables import DeltaTable2from transforms import sessionize3from quality import apply_checks4 5bronze = spark.read.format("delta").load("/tmp/media/bronze_playback")6valid, rejected = apply_checks(bronze)7failure_rate = rejected.count() / max(bronze.count(), 1)8assert failure_rate < 0.01, f"Quality gate failed: {failure_rate:.2%}"9 10PAYLOAD = ["viewer_id", "title_id", "event_type", "event_ts", "watch_seconds"]11deduplicated = valid.dropDuplicates(PAYLOAD)12 13silver_batch = (14 sessionize(deduplicated)15 .withColumn("event_key", F.sha2(F.concat_ws("|",16 "viewer_id", "title_id", "event_type", "event_ts", "watch_seconds"), 256))17 .join(F.broadcast(catalog), "title_id")18)19 20SILVER = "/tmp/media/silver_sessions"21if not DeltaTable.isDeltaTable(spark, SILVER):22 silver_batch.limit(0).write.format("delta").save(SILVER)23 24(DeltaTable.forPath(spark, SILVER).alias("t")25 .merge(silver_batch.alias("s"), "t.event_key = s.event_key")26 .whenMatchedUpdateAll()27 .whenNotMatchedInsertAll()28 .execute())Full-history recalculation can change existing session assignments when older events arrive, so matched rows must be updated, not just skipped. This append-only synthetic lab assumes required payload fields are validated and valid rows are not later deleted. It is a teaching rebuild, not incremental late-arrival sessionization: an incremental design must identify affected viewers and time ranges, reconcile session boundaries, and handle changed or removed assignments explicitly.
Gold: daily genre metrics
1gold = (2 spark.read.format("delta").load(SILVER)3 .withColumn("day", F.to_date("event_ts"))4 .groupBy("day", "genre")5 .agg(F.countDistinct("viewer_id").alias("viewers"),6 F.round(F.sum("watch_seconds") / 3600, 1).alias("hours"),7 F.round(F.avg((F.col("event_type") == "complete").cast("double")), 3)8 .alias("completion_event_share"))9)10 11(gold.write.format("delta").mode("overwrite")12 .save("/tmp/media/gold_genre_daily"))Gold is fully overwritten from the recalculated silver history for this lab. An arbitrary replaceWhere date range would leave stale aggregates outside that range. Production incremental updates need explicit affected-partition logic rather than this full-history teaching pipeline.
| Layer | Contract | Owner checks |
|---|---|---|
| Bronze | Every parsed row, plus file lineage | Ingestion counts and quarantine volume |
| Silver | Valid, deduplicated, sessionized | Quality gate pass rate, unit tests |
| Gold | Documented business metrics | Metric definitions reviewed with consumers |
What to take forward
- Read plans before tuning.
- Define schemas and metrics explicitly.
- Make transformations testable functions.
- Choose watermarks from observed lateness.
- Make every write safe to rerun.
Exercises
- Run silver.py twice and confirm the row count does not change.
- Schedule the batch layers as a job and the concurrency stream separately.
- Add a gold table for device mix per day and write its metric definition first.
- Write a short runbook: what to check when gold numbers look wrong.
Official documentation
