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.
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.
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)| Concept | Meaning here |
|---|---|
| Event time | When the heartbeat happened on the device (event_ts) |
| Watermark | Derived from maximum observed event time minus the delay; not from wall-clock time |
| Append mode | A window is written once, after the watermark passes it |
| Checkpoint | Offsets 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
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.
Official documentation
