Taming Late Data in Spark Structured Streaming: Practical Event‑Time Watermarking
When IoT devices send logs late, dashboards can become stale. Spark’s event‑time watermarks let you drop or keep late data automatically. This guide shows how to set up watermarks, a working example, and the trade‑offs you need to monitor.
26 Sept 2026, 06:14 UTC

The Late‑Data Problem in Real‑Time IoT
IoT sensors often send data out of order. A temperature reading generated at 14:03 may arrive at the stream at 14:07 due to network jitter or device buffering. If you aggregate on event time without any guard, that late record will be included in the 14:00‑14:05 window, skewing the average and potentially triggering false alerts.
Without a safety net, the only way to keep aggregates accurate is to wait indefinitely for all late data, which is impractical for dashboards that need near‑real‑time updates.
Event‑Time Watermarking in Spark Structured Streaming
Spark Structured Streaming introduces watermarks to bound how late data may be considered for aggregations. A watermark is a timestamp that represents the maximum event time seen so far minus a user‑defined delay. When the watermark passes a record’s event time, that record is deemed too late and is dropped for windowed operators.
Key points:
- Watermarks are defined per
groupByorwindowoperation. - They work only with event‑time columns, not processing time.
- They automatically clean state for windows that have passed the watermark.
- State size and processing latency are bounded by the watermark delay.
Practical Example: Temperature Sensor Dashboard
Below is a minimal Scala snippet that demonstrates reading from Kafka, applying a 5‑minute tumbling window on event time, and setting a 2‑minute watermark. Late records arriving after the watermark are discarded.
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
val spark = SparkSession.builder
.appName("IoT Watermark Demo")
.getOrCreate()
// Kafka source: replace placeholders with your config.
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka-broker:9092")
.option("subscribe", "sensor-readings")
.option("startingOffsets", "earliest")
.load()
// Assume the value is JSON with fields: sensorId, temperature, eventTime.
val sensorDF = df.selectExpr("CAST(value AS STRING)")
.select(from_json(col("value"), schema_of_json(lit("{\"sensorId\":\"s1\",\"temperature\":23.5,\"eventTime\":\"2026-10-10T14:03:00Z\"}")))
.as[SensorRecord]
.withColumn("eventTime", to_timestamp(col("eventTime")))
val windowed = sensorDF
.withWatermark("eventTime", "2 minutes") // watermark
.groupBy(window(col("eventTime"), "5 minutes"), col("sensorId"))
.agg(avg("temperature").alias("avgTemp"))
val query = windowed.writeStream
.outputMode("update")
.format("console")
.option("truncate", "false")
.start()
query.awaitTermination()
Explanation of placeholders:
kafka.bootstrap.servers– comma‑separated list of broker hosts.subscribe– Kafka topic containing sensor JSON.- Schema inference is shown via
schema_of_jsonfor clarity; in production define a static schema.
When you run this, you should see console output every few seconds. Events whose eventTime is older than the current watermark will not contribute to the next window’s average.
Choosing the Right Watermark: Trade‑offs and Monitoring
Setting the watermark is a balancing act:
- Short watermark (e.g., 30 s) – drops more late data, reducing state size but risking loss of legitimate events that arrive slightly late.
- Long watermark (e.g., 10 min) – keeps more data, ensuring accuracy but increasing state size and potentially causing higher latency or out‑of‑memory (OOM) errors in stateful operators.
How to decide:
- Collect latency statistics from your ingestion pipeline. If most records arrive within 2 min of their event time, a 2‑minute watermark is reasonable.
- Monitor Spark UI
Streamingtab: look atWatermark LagandState Sizemetrics. If state size grows steadily, consider tightening the watermark. - Enable
spark.sql.streaming.stateStore.stateRetentionDurationto automatically purge old state and mitigate OOM. - Periodically run a small test job that injects intentionally late records and verify they’re dropped when expected.
Remember that watermarks are only evaluated at the point of the aggregation. If you write data to a sink (e.g., Parquet) before the watermark passes, those late records will still be persisted unless you add a filter.
Take‑Away: Tune, Monitor, and Avoid State Bloat
1. Define a realistic watermark based on observed event‑time delays. Start with a conservative value and tighten as you gain confidence.
2. Instrument your pipeline – enable Spark UI metrics, log watermark lag, and set alerts if state size exceeds a threshold.
3. Use state TTL and checkpointing to keep memory usage in check, especially for long‑running jobs.
4. Validate with tests – inject late data, run the job, and confirm aggregates ignore late events after the watermark.
With these practices, Spark Structured Streaming can deliver accurate, real‑time dashboards even when IoT devices send data out of order.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.