Lesson reference: PySpark and Scala
Use tumbling windows
The setup gives the six orders explicit event timestamps at minutes 0, 2, 7, 11, 14, and 19. These are event times, not query execution times. Ten-minute tumbling windows do not overlap; boundaries include the start and exclude the end. Complete output mode exposes all current aggregates for this finite example and retains aggregate state.
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.groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = (transformed.writeStream.format("memory").queryName(name)
.outputMode("complete").option("checkpointLocation", directory + "/checkpoint").start())
timer = Timer(45, query.stop)
timer.daemon = True
timer.start()
query.processAllAvailable()
if not query.isActive:
raise RuntimeError("Streaming example stopped before completion; retry the run.")
snapshot = spark.table(name).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.groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = transformed.writeStream.format("memory").queryName(name)
.outputMode("complete").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
query.processAllAvailable()
require(query.isActive, "Streaming example stopped before completion; retry the run.")
val snapshot = spark.table(name).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.
}
Slide overlapping windows
A slide smaller than the duration creates overlapping windows. An event can contribute to multiple windows, so summing their counts can exceed the number of input events. Window alignment is based on the window duration and slide, not the first observed timestamp; a window can start before the first event.
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.groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = (transformed.writeStream.format("memory").queryName(name)
.outputMode("complete").option("checkpointLocation", directory + "/checkpoint").start())
timer = Timer(45, query.stop)
timer.daemon = True
timer.start()
query.processAllAvailable()
if not query.isActive:
raise RuntimeError("Streaming example stopped before completion; retry the run.")
snapshot = spark.table(name).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.groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = transformed.writeStream.format("memory").queryName(name)
.outputMode("complete").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
query.processAllAvailable()
require(query.isActive, "Streaming example stopped before completion; retry the run.")
val snapshot = spark.table(name).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.
}
Aggregate values, not only events
The window defines membership; the aggregate defines the measurement. Counting events differs from summing item_count. Both operate over the same event-time windows. Complete mode is used only for this bounded illustration; without state limits, an indefinitely running aggregation can retain growing state.
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.groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = (transformed.writeStream.format("memory").queryName(name)
.outputMode("complete").option("checkpointLocation", directory + "/checkpoint").start())
timer = Timer(45, query.stop)
timer.daemon = True
timer.start()
query.processAllAvailable()
if not query.isActive:
raise RuntimeError("Streaming example stopped before completion; retry the run.")
snapshot = spark.table(name).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.groupBy(
window(col("event_time"), "10 minutes", "10 minutes")
).agg(count("*").alias("metric"))
query = transformed.writeStream.format("memory").queryName(name)
.outputMode("complete").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
query.processAllAvailable()
require(query.isActive, "Streaming example stopped before completion; retry the run.")
val snapshot = spark.table(name).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.
}