Watermarks & Late Data
Deciding a window is complete when events can still arrive, and what to do when they do.
5 to work through
-
advanced
A daily reconciliation shows the streaming aggregation is 2% below the batch equivalent. Both read the same source. Why?
2 min answer -
advanced
A dispatch analytics pipeline computes per-minute aggregates and the numbers keep changing after the fact. What is happening, and what decisions must be made explicit?
2 min answer -
advanced
A streaming aggregation must produce results promptly while some events arrive minutes late. How do watermarks, windowing and allowed lateness interact, and what is the business decision underneath?
2 min answer -
advanced
Mobile clients buffer events offline and deliver them hours later. How should watermarks and late data be handled?
2 min answer -
advanced
Your streaming aggregate reports 2% lower daily revenue than the batch reconciliation. Both are "correct". Explain what is happening and how you resolve it.
2 min answer
5 terms in this topic
Late-Data Policy
The explicit decision about what happens to an event that arrives after its window has closed - drop it, route it aside, or update the emitted result.
conceptProvisional Result
An aggregate emitted before its window is definitively complete, which may be corrected when late events arrive - correct behaviour that destroys tru…
conceptWatermark
A heuristic assertion that no further events earlier than a given event time are expected, which is how a stream processor decides when a window is c…
conceptWatermark
The stream processor's assertion that no further events older than a given event time are expected, which is what allows a window to close.
metricWatermark Lag
How far behind the newest event a watermark is held, which trades result latency against how much late data is included.
Neighbouring topics
Streaming & Real-Time Data
General material on continuous processing of unbounded data.
Streaming vs Batch
The freshness requirement that actually justifies streaming, and the cost of assuming one.
Exactly-Once Semantics
What the phrase really means, where it holds, and the idempotent sink underneath it.
Stream Processing Frameworks
Flink, Kafka Streams, Spark Structured Streaming — state, checkpointing and recovery.
Windowing
Tumbling, sliding and session windows, and the aggregation each one answers.
Stateful Stream Processing
Keyed state, state backends, checkpoint size, and the restore time that follows.
Stream-Table Duality
A changelog and a table as two views of the same thing, and materialising between them.
Kappa vs Lambda
One pipeline replayed versus two pipelines reconciled, and the maintenance each carries.
Streaming Schema Evolution
Changing an event's shape while a retained log still holds every older version of it.
Streaming Joins
Joining two unbounded streams, the buffering it needs, and the enrichment alternative.
Backfill & Reprocessing
Replaying history through changed logic without double-counting the live output.
CDC to Stream
Turning database changes into an event log, and how that differs from a domain event.
Real-Time Serving Layer
Where a low-latency read of a streaming aggregate actually lands.
Feature Freshness
How stale a feature can be before the model degrades, and the pipeline that follows.
Streaming SLOs
End-to-end latency, consumer lag and completeness as commitments rather than dashboards.
Partition Keys & Ordering
Ordering guaranteed only within a partition, and choosing the key that makes that enough.
Dead Letter Handling
The poison message that blocks a partition, and the queue nobody reads.
Streaming Cost
Always-on compute, retention and cross-zone traffic as the three bills that surprise.
Real-Time Analytical Stores
Druid, Pinot and ClickHouse — ingest-and-query engines for sub-second aggregation.