ML Engineer MasterClass (October) | 4 seats left

Spark · Event-time windows
AmazonAmazon Analytics
Amazon · Structured Streaming

Event-time windows

Group events by their timestamps using tumbling and sliding windows.

Step 1 of 6 · Learn

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.

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.
}
example.pyPySpark
1. Event timestampsUse explicit timestamps on the six source orders.
2. Window membershipDuration and slide determine each event’s window memberships.
3. Current aggregatesComplete output shows the finite example’s current window results.
Source order_events6 rows
order_idcustomer_idstatustotalitem_countevent_time
1001101Delivered89.522026-09-01 00:00:00
1002102Shipped14932026-09-01 00:02:00
1003101Cancelled3512026-09-01 00:07:00
1004103Delivered219.9942026-09-01 00:11:00
1005104Delivered49.9912026-09-01 00:14:00
1006105Shipped12022026-09-01 00:19:00
Use tumbling windows
Result2 rows
startendmetric
2026-09-01 00:002026-09-01 00:103
2026-09-01 00:102026-09-01 00:203
The Learn example has two windows with three events each. No watermark-based eviction is being demonstrated here.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.