Back to blog

Spark for Streaming Media Analytics · Part 5 of 8

Spark Structured Streaming: Watermarks and Concurrent Viewers

Measure per-minute heartbeat activity on an unbounded stream and bound state with watermarks.

Chris Eberl

Chris Eberl

Founder • Engineering Leader, GenAI, Data, ML

Data EngineeringOct 8, 202610 min read
Spark Structured Streaming: Watermarks and Concurrent Viewers

Structured Streaming lets you write almost the same DataFrame code against data that never stops arriving. The new questions are about time: which clock you measure by, how late an event may be, and how long Spark must remember state.

A synthetic stream

The built-in rate source emits rows continuously, which makes it ideal for practice. We shape it into playback events and add artificial lateness.

stream_source.py
1stream = (2    spark.readStream.format("rate").option("rowsPerSecond", 200).load()3    .withColumn("viewer_id", (F.col("value") % 2_000).cast("int"))4    .withColumn("title_id", (F.col("value") % 300).cast("int"))5    .withColumn("event_type", F.when(F.col("value") % 5 == 0, "stop").otherwise("heartbeat"))6    # up to 90 seconds of simulated lateness7    .withColumn("event_ts", F.col("timestamp") - F.expr("make_interval(0,0,0,0,0,0, rand() * 90)"))8)

Windowed heartbeat activity with a watermark

Approximate distinct heartbeat viewers per title per minute is an activity proxy, not exact instantaneous concurrency. A viewer can appear anywhere in the minute; stop rows are not used to maintain active sessions. The example retains the concurrent_viewers column name for continuity, but that name does not change the metric’s limits.

concurrency.py
1concurrent = (2    stream.filter(F.col("event_type") == "heartbeat")3    .withWatermark("event_ts", "2 minutes")4    .groupBy(F.window("event_ts", "1 minute"), "title_id")5    .agg(F.approx_count_distinct("viewer_id").alias("concurrent_viewers"))6)7 8query = (9    concurrent.writeStream10    .outputMode("append")11    .format("delta")12    .option("checkpointLocation", "/tmp/media/_chk/concurrency")13    .trigger(processingTime="30 seconds")14    .start("/tmp/media/gold_concurrency")15)
ConceptMeaning here
Event timeWhen the heartbeat happened on the device (event_ts)
WatermarkDerived from maximum observed event time minus the delay; not from wall-clock time
Append modeA window is written once, after the watermark passes it
CheckpointOffsets and state that let the query restart without double counting

The watermark advances from the maximum observed event time. Records less late than the configured delay are guaranteed not to be dropped by that watermark; records more than the delay late may or may not be dropped. Append mode emits finalized windows only after the watermark passes their end. A short smoke run may produce no finalized rows: run the companion with python chapters/05_streaming.py --seconds 240 (at least 240 seconds) to observe output, allowing longer on a slow machine.

Operate the query

Inspect progress
1query.status2# Keep the stream running long enough for append windows to finalize.3query.awaitTermination(240)4progress = query.lastProgress5if progress is not None:6    print(progress["numInputRows"], progress.get("stateOperators", []))7else:8    print("No completed trigger progress yet")9query.stop()

Exercises

  • Shorten the watermark to 30 seconds and count how many late rows are dropped.
  • Switch to update output mode writing to the console sink. How does the output differ?
  • Delete the checkpoint and restart. What happens to previously written windows?
  • Add a sliding window of 5 minutes every 1 minute and compare state size.

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.