Back to blog

Spark for Streaming Media Analytics · Part 1 of 8

PySpark Mental Model: Lazy DataFrames for Media Analytics

Generate a synthetic playback dataset and learn why Spark waits until you ask for an answer.

Chris Eberl

Chris Eberl

Founder • Engineering Leader, GenAI, Data, ML

Data EngineeringOct 8, 20269 min read
PySpark Mental Model: Lazy DataFrames for Media Analytics

Most PySpark confusion comes from one assumption: that each line runs when you write it. It does not. Spark records what you want, builds a plan, and only executes when an action asks for a result. Once that clicks, performance, debugging, and even AI-generated code become much easier to reason about.

This series uses one running example: a fictional video service producing playback events. Over eight parts we go from a local DataFrame to a streaming pipeline, tested transformations, and a review workflow with Databricks Genie Code.

Generate the dataset

We start with a deterministic generator so every reader gets identical numbers. Each row is one playback event: a viewer starts, pauses, resumes, or completes a title on a device.

generate_events.py
1from pyspark.sql import SparkSession, functions as F2 3spark = SparkSession.builder.appName("media-analytics").getOrCreate()4 5N_EVENTS = 200_0006events = (7    spark.range(N_EVENTS)8    .withColumn("viewer_id", (F.col("id") % 5_000).cast("int"))9    .withColumn("title_id", (F.hash("id") % 300).cast("int"))10    .withColumn("title_id", F.abs("title_id"))11    .withColumn("device", F.element_at(12        F.array(*[F.lit(d) for d in ["tv", "mobile", "web", "tablet"]]),13        (F.col("id") % 4 + 1).cast("int")))14    .withColumn("event_type", F.element_at(15        F.array(*[F.lit(e) for e in ["start", "pause", "resume", "complete"]]),16        (F.abs(F.hash("id", F.lit("e"))) % 4 + 1).cast("int")))17    .withColumn("event_ts", F.timestamp_seconds(18        F.lit(1_790_000_000) + F.col("id") * 7))19    .withColumn("watch_seconds",20        (F.abs(F.hash("id", F.lit("w"))) % 3_600).cast("int"))21    .drop("id")22)

Notice that nothing has run yet. events is a description of work: a logical plan with a range source and a chain of projections.

This generator repeats a viewer every 5,000 events, with timestamps seven seconds apart: successive events for that viewer are 35,000 seconds apart. Most generated sessions therefore contain one event under the 30-minute rule; use handcrafted rows to explore session boundaries. watch_seconds is an independent per-event contribution, not cumulative player position or elapsed time between events.

Transformations versus actions

KindExamplesWhat happens
Transformationselect, filter, withColumn, join, groupByAdds a step to the plan. Returns a new DataFrame immediately.
Actioncount, show, collect, writeTriggers optimization and execution across the cluster.
Ask a question
1completions_by_device = (2    events.filter(F.col("event_type") == "complete")3    .groupBy("device")4    .agg(F.count("*").alias("completions"),5         F.round(F.avg("watch_seconds") / 60, 1).alias("avg_minutes"))6    .orderBy(F.desc("completions"))7)8 9completions_by_device.explain(mode="formatted")  # still lazy10completions_by_device.show()                      # action: runs now

explain is the habit worth building first. It shows the physical plan Spark intends to run, including where it will shuffle data between machines. Reading plans early makes later performance work feel familiar rather than mysterious.

Why laziness matters

  • The optimizer can push filters toward the source and drop unused columns before reading them.
  • Chained transformations fuse into fewer passes over the data.
  • Every action recomputes the plan unless you cache — so calling count() in a loop is expensive.
  • Errors in a transformation often surface only at the action, several lines later.

Exercises

  • Count events per event_type and confirm the split is roughly even.
  • Find the five titles with the most completions on tv.
  • Run explain() before and after adding a filter on device. Where does the filter appear in the plan?
  • Add a cache() to events, run two different actions, and compare timings in the Spark UI.

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.