ML Engineer MasterClass (October) | 4 seats left

Spark · Watermarks
AmazonAmazon Analytics
Amazon · Structured Streaming

Watermarks

Relate event-time progress, accepted late events, and finalized append-mode windows.

Step 1 of 6 · Learn

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.

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.
}
example.pyPySpark
1. Initial event-time batchEvents reach 00:19; the declared delay governs event-time progress.
2. Optional late event, then progressA late example arrives before 01:00 and 01:10 progress batches.
3. Finalized early windowsAppend-mode output is snapshotted after processing; future windows are excluded.
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
Finalize with event-time progress
Result2 rows
startendmetric
2026-09-01 00:002026-09-01 00:103
2026-09-01 00:102026-09-01 00:203
Both delays eventually finalize these initial windows after the progress batches. Their output timing and state retention can differ.

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.