Lesson reference: PySpark and Scala
Finalize with event-time progress
A watermark lets supported stateful operators reason about event-time progress and retire old state. withWatermark must precede aggregation on that timestamp. Append-mode window results appear when finalized, not after every event. The harness appends progress events at 01:00 and 01:10 to move event time forward, then excludes their later windows from the teaching result. Wall-clock waiting alone is not the mechanism.
PySpark example
from tempfile import TemporaryDirectory
from uuid import uuid4
from threading import Timer
from pyspark.sql.functions import col, lit, array, element_at, to_timestamp, window, count, sum, upper, lower, date_format
seed = orders.withColumn("event_time", to_timestamp(element_at(
array(lit("2026-09-01 00:00:00"), lit("2026-09-01 00:02:00"), lit("2026-09-01 00:07:00"), lit("2026-09-01 00:11:00"), lit("2026-09-01 00:14:00"), lit("2026-09-01 00:19:00")),
(col("order_id") - 1000).cast("int")
)))
name = "lesson_" + uuid4().hex
query = None
timer = None
with TemporaryDirectory(prefix="spark-stream-") as directory:
path = directory + "/input"
try:
seed.write.json(path)
stream = spark.readStream.schema(seed.schema).json(path)
transformed = stream.withWatermark("event_time", "10 minutes").groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = (transformed.writeStream.format("memory").queryName(name)
.outputMode("append").option("checkpointLocation", directory + "/checkpoint").start())
timer = Timer(45, query.stop)
timer.daemon = True
timer.start()
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:00:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:10:00"))).write.mode("append").json(path)
query.processAllAvailable()
if not query.isActive:
raise RuntimeError("Streaming example stopped before completion; retry the run.")
snapshot = spark.table(name).filter(col("window.start") < to_timestamp(lit("2026-09-01 00:30:00"))).select(date_format(col("window.start"), "yyyy-MM-dd HH:mm").alias("start"), date_format(col("window.end"), "yyyy-MM-dd HH:mm").alias("end"), col("metric"))
# Detach these few rows before stopping and removing the memory table.
result = spark.createDataFrame(snapshot.collect(), snapshot.schema).orderBy("start", "end", "metric")
result.show(truncate=False)
finally:
if timer is not None:
timer.cancel()
if query is not None:
query.stop()
spark.catalog.dropTempView(name)
Scala example
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.StreamingQuery
import java.nio.file.Files
import org.apache.hadoop.fs.Path
val seed = orders.withColumn("event_time", to_timestamp(element_at(
array(lit("2026-09-01 00:00:00"), lit("2026-09-01 00:02:00"), lit("2026-09-01 00:07:00"), lit("2026-09-01 00:11:00"), lit("2026-09-01 00:14:00"), lit("2026-09-01 00:19:00")),
(col("order_id") - 1000).cast("int")
)))
val name = "lesson_" + java.util.UUID.randomUUID().toString.replace("-", "")
val directory = Files.createTempDirectory("spark-stream-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("input").toString
var query: StreamingQuery = null
val timer = new java.util.Timer(true)
val result = try {
seed.write.json(path)
val stream = spark.readStream.schema(seed.schema).json(path)
val transformed = stream.withWatermark("event_time", "10 minutes").groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = transformed.writeStream.format("memory").queryName(name)
.outputMode("append").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:00:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:10:00"))).write.mode("append").json(path)
query.processAllAvailable()
require(query.isActive, "Streaming example stopped before completion; retry the run.")
val snapshot = spark.table(name).filter(col("window.start") < to_timestamp(lit("2026-09-01 00:30:00"))).select(date_format(col("window.start"), "yyyy-MM-dd HH:mm").alias("start"), date_format(col("window.end"), "yyyy-MM-dd HH:mm").alias("end"), col("metric"))
val output = spark.createDataFrame(spark.sparkContext.parallelize(snapshot.collect().toSeq), snapshot.schema).orderBy("start", "end", "metric")
output.show(false)
output
} finally {
timer.cancel()
if (query != null) query.stop()
spark.catalog.dropTempView(name)
fs.delete(root, true) // Only this example’s input and checkpoint directory.
}
Keep delay and window size separate
Watermark delay and window duration solve different problems. A delay describes tolerated event-time lateness; a window defines which events are aggregated together. Keep append mode and a watermark on event_time while changing the window width. Complete mode would retain all aggregate results rather than demonstrate this append-mode finalization.
PySpark example
from tempfile import TemporaryDirectory
from uuid import uuid4
from threading import Timer
from pyspark.sql.functions import col, lit, array, element_at, to_timestamp, window, count, sum, upper, lower, date_format
seed = orders.withColumn("event_time", to_timestamp(element_at(
array(lit("2026-09-01 00:00:00"), lit("2026-09-01 00:02:00"), lit("2026-09-01 00:07:00"), lit("2026-09-01 00:11:00"), lit("2026-09-01 00:14:00"), lit("2026-09-01 00:19:00")),
(col("order_id") - 1000).cast("int")
)))
name = "lesson_" + uuid4().hex
query = None
timer = None
with TemporaryDirectory(prefix="spark-stream-") as directory:
path = directory + "/input"
try:
seed.write.json(path)
stream = spark.readStream.schema(seed.schema).json(path)
transformed = stream.withWatermark("event_time", "10 minutes").groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = (transformed.writeStream.format("memory").queryName(name)
.outputMode("append").option("checkpointLocation", directory + "/checkpoint").start())
timer = Timer(45, query.stop)
timer.daemon = True
timer.start()
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:00:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:10:00"))).write.mode("append").json(path)
query.processAllAvailable()
if not query.isActive:
raise RuntimeError("Streaming example stopped before completion; retry the run.")
snapshot = spark.table(name).filter(col("window.start") < to_timestamp(lit("2026-09-01 00:30:00"))).select(date_format(col("window.start"), "yyyy-MM-dd HH:mm").alias("start"), date_format(col("window.end"), "yyyy-MM-dd HH:mm").alias("end"), col("metric"))
# Detach these few rows before stopping and removing the memory table.
result = spark.createDataFrame(snapshot.collect(), snapshot.schema).orderBy("start", "end", "metric")
result.show(truncate=False)
finally:
if timer is not None:
timer.cancel()
if query is not None:
query.stop()
spark.catalog.dropTempView(name)
Scala example
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.StreamingQuery
import java.nio.file.Files
import org.apache.hadoop.fs.Path
val seed = orders.withColumn("event_time", to_timestamp(element_at(
array(lit("2026-09-01 00:00:00"), lit("2026-09-01 00:02:00"), lit("2026-09-01 00:07:00"), lit("2026-09-01 00:11:00"), lit("2026-09-01 00:14:00"), lit("2026-09-01 00:19:00")),
(col("order_id") - 1000).cast("int")
)))
val name = "lesson_" + java.util.UUID.randomUUID().toString.replace("-", "")
val directory = Files.createTempDirectory("spark-stream-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("input").toString
var query: StreamingQuery = null
val timer = new java.util.Timer(true)
val result = try {
seed.write.json(path)
val stream = spark.readStream.schema(seed.schema).json(path)
val transformed = stream.withWatermark("event_time", "10 minutes").groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = transformed.writeStream.format("memory").queryName(name)
.outputMode("append").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:00:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:10:00"))).write.mode("append").json(path)
query.processAllAvailable()
require(query.isActive, "Streaming example stopped before completion; retry the run.")
val snapshot = spark.table(name).filter(col("window.start") < to_timestamp(lit("2026-09-01 00:30:00"))).select(date_format(col("window.start"), "yyyy-MM-dd HH:mm").alias("start"), date_format(col("window.end"), "yyyy-MM-dd HH:mm").alias("end"), col("metric"))
val output = spark.createDataFrame(spark.sparkContext.parallelize(snapshot.collect().toSeq), snapshot.schema).orderBy("start", "end", "metric")
output.show(false)
output
} finally {
timer.cancel()
if (query != null) query.stop()
spark.catalog.dropTempView(name)
fs.delete(root, true) // Only this example’s input and checkpoint directory.
}
Accept an out-of-order event
After the initial batch reaches event time 00:19, this example appends an event timestamped 00:12. It is out of order but within the ten-minute delay, so its still-open five-minute window can update. The exercise uses 00:17 instead. The later progress batches then finalize results. Data older than the watermark may be dropped; do not interpret the watermark as a guarantee that every older event will always be discarded.
PySpark example
from tempfile import TemporaryDirectory
from uuid import uuid4
from threading import Timer
from pyspark.sql.functions import col, lit, array, element_at, to_timestamp, window, count, sum, upper, lower, date_format
seed = orders.withColumn("event_time", to_timestamp(element_at(
array(lit("2026-09-01 00:00:00"), lit("2026-09-01 00:02:00"), lit("2026-09-01 00:07:00"), lit("2026-09-01 00:11:00"), lit("2026-09-01 00:14:00"), lit("2026-09-01 00:19:00")),
(col("order_id") - 1000).cast("int")
)))
name = "lesson_" + uuid4().hex
query = None
timer = None
with TemporaryDirectory(prefix="spark-stream-") as directory:
path = directory + "/input"
try:
seed.write.json(path)
stream = spark.readStream.schema(seed.schema).json(path)
transformed = stream.withWatermark("event_time", "10 minutes").groupBy(
window(col("event_time"), "5 minutes", "5 minutes")
).agg(count("*").alias("metric"))
query = (transformed.writeStream.format("memory").queryName(name)
.outputMode("append").option("checkpointLocation", directory + "/checkpoint").start())
timer = Timer(45, query.stop)
timer.daemon = True
timer.start()
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 00:12:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:00:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") == 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:10:00"))).write.mode("append").json(path)
query.processAllAvailable()
if not query.isActive:
raise RuntimeError("Streaming example stopped before completion; retry the run.")
snapshot = spark.table(name).filter(col("window.start") < to_timestamp(lit("2026-09-01 00:30:00"))).select(date_format(col("window.start"), "yyyy-MM-dd HH:mm").alias("start"), date_format(col("window.end"), "yyyy-MM-dd HH:mm").alias("end"), col("metric"))
# Detach these few rows before stopping and removing the memory table.
result = spark.createDataFrame(snapshot.collect(), snapshot.schema).orderBy("start", "end", "metric")
result.show(truncate=False)
finally:
if timer is not None:
timer.cancel()
if query is not None:
query.stop()
spark.catalog.dropTempView(name)
Scala example
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.StreamingQuery
import java.nio.file.Files
import org.apache.hadoop.fs.Path
val seed = orders.withColumn("event_time", to_timestamp(element_at(
array(lit("2026-09-01 00:00:00"), lit("2026-09-01 00:02:00"), lit("2026-09-01 00:07:00"), lit("2026-09-01 00:11:00"), lit("2026-09-01 00:14:00"), lit("2026-09-01 00:19:00")),
(col("order_id") - 1000).cast("int")
)))
val name = "lesson_" + java.util.UUID.randomUUID().toString.replace("-", "")
val directory = Files.createTempDirectory("spark-stream-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("input").toString
var query: StreamingQuery = null
val timer = new java.util.Timer(true)
val result = try {
seed.write.json(path)
val stream = spark.readStream.schema(seed.schema).json(path)
val transformed = stream.withWatermark("event_time", "10 minutes").groupBy(
window(col("event_time"), "5 minutes", "5 minutes")
).agg(count("*").alias("metric"))
query = transformed.writeStream.format("memory").queryName(name)
.outputMode("append").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 00:12:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:00:00"))).write.mode("append").json(path)
query.processAllAvailable()
seed.filter(col("order_id") === 1001).withColumn("event_time", to_timestamp(lit("2026-09-01 01:10:00"))).write.mode("append").json(path)
query.processAllAvailable()
require(query.isActive, "Streaming example stopped before completion; retry the run.")
val snapshot = spark.table(name).filter(col("window.start") < to_timestamp(lit("2026-09-01 00:30:00"))).select(date_format(col("window.start"), "yyyy-MM-dd HH:mm").alias("start"), date_format(col("window.end"), "yyyy-MM-dd HH:mm").alias("end"), col("metric"))
val output = spark.createDataFrame(spark.sparkContext.parallelize(snapshot.collect().toSeq), snapshot.schema).orderBy("start", "end", "metric")
output.show(false)
output
} finally {
timer.cancel()
if (query != null) query.stop()
spark.catalog.dropTempView(name)
fs.delete(root, true) // Only this example’s input and checkpoint directory.
}