Back to blog

Spark for Streaming Media Analytics · Part 8 of 8

Capstone: An End-to-End Spark Medallion Pipeline for Media Data

Connect bronze, silver, and gold layers into a pipeline you can test, rerun, and explain.

Chris Eberl

Chris Eberl

Founder • Engineering Leader, GenAI, Data, ML

Data EngineeringOct 8, 202610 min read
Capstone: An End-to-End Spark Medallion Pipeline for Media Data

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

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

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

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

LayerContractOwner checks
BronzeEvery parsed row, plus file lineageIngestion counts and quarantine volume
SilverValid, deduplicated, sessionizedQuality gate pass rate, unit tests
GoldDocumented business metricsMetric 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.

Read next

Newsletter

New posts, straight from Chris

A short note from me whenever a new article goes live — product engineering, AI workflows, IoT, indie apps, and engineering leadership. No spam, unsubscribe anytime.

By subscribing, you agree to our Privacy Policy. We do not share your email.