Event Time & Watermarks — Allowed Lateness, Side Outputs & Idle Sources
In event-time processing the engine never knows a window is complete. A phone that was in airplane mode can upload yesterday's clicks right now. A watermark is the engine's declaration that it believes no more events older than W will arrive. It drives on-time window firing, timers, and state cleanup. Too aggressive and you drop or mis-count late data. Too conservative and every result is delayed. Allowed lateness keeps window state around a bit longer for corrections, a side output captures anything later than that, and idleness handling stops one silent partition from freezing the whole job.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What exactly is a watermark?
Answer
A timestamp W flowing through the dataflow that asserts no more records with event time at or before W are expected. It is a heuristic for unbounded sources. Records that violate it are late.
L2
How do you pick the out-of-orderness bound?
Answer
From the distribution of event-time delay per source. Set it near a high percentile, balancing latency against completeness, then handle the tail with allowed lateness and a side output.
L3
Why does one idle Kafka partition stall the job, and how do you fix it?
Answer
Operators take the min watermark across inputs. An idle partition never advances, so the min is frozen. Mark the split idle after a timeout, or emit heartbeat records upstream.
L4
What is the difference between allowed lateness and a bigger watermark delay?
Answer
A bigger delay postpones the first result for every window. Allowed lateness emits on time and then corrects with late firings, at the cost of keeping state longer and needing update-capable sinks.
L5
What happens to records later than allowed lateness?
Answer
They are dropped by default. Count the drops. Better: route them to a side output and reconcile in a correction job or a batch pass.
L6
How do watermarks interact with joins and multiple inputs?
Answer
Two-input operators take the min of both inputs, so the slower stream controls when interval-join state can be cleaned and when results are final.
L7
Why should watermarks be generated at the source, per partition?
Answer
Ordering is usually strongest inside a partition. Per-split watermarks capture that and are combined by min. Generating them after a shuffle mixes partitions and forces a much larger disorder bound.
Failure modes
Idle partition pins the minimum
Windows never fire, state grows, and the dashboard shows nothing, with zero errors.
A producer clock set to 2030
Max event time jumps into the future and everything else looks late.
Late firing into an append-only sink
The window already emitted. The correction becomes a second row instead of an update.
Watermarks after the shuffle
Per-partition order is gone, so the disorder bound has to cover the mix.
Misconceptions
A watermark means the window is complete.
It means the engine is willing to bet. Allowed lateness and a side output are the rest of the bet.
Idleness is free progress.
Excluding a quiet split lets the watermark move. A record from that split can then be late if the others raced ahead.
Processing time is fine for per-hour reports if you are usually caught up.
A backlog or a backfill stuffs a day of events into the current hour.
Interviewer traps
Picking the disorder bound by gut feel.
Ask for the delay distribution. Set the bound at a high percentile and measure the late-record counter.
Making the watermark delay huge so nothing is late.
That delays every on-time result. Allowed lateness corrects the tail without postponing the first pane.
Design scenario
Same prompt for every reader.
Requirements
A bounded-out-of-orderness watermark, allowed lateness with upserts, a side output for the tail, and idleness so the quiet partition does not freeze the job.
Traffic / scale
Three partitions. One region is near zero at night.
Latency
On-time panes within about ten seconds of the watermark. Corrections for another ten seconds of event time.
Consistency
A replay of the same log produces the same panes. Late updates overwrite by window key.
Availability
The job keeps firing when a partition is idle. Heartbeats or idleness are required, not optional.
Failure assumptions
- One device clock can be hours off.
- A partition can be silent for minutes.
Constraints
- Do not generate the watermark after the keyBy shuffle.
- Do not append late panes as new rows.
Prompt
Mobile clicks arrive up to a few minutes late, and one Kafka partition goes quiet overnight. Product wants per-minute counts that can update, and finance cannot lose the very late clicks.
API
What does a side-output record contain so a correction job can apply it?
Data
Which timestamp moves the watermark, and which one is clamped?
Architecture
Where is the watermark born, and which input's minimum holds a join?
Close the window, or wait longer
Prefer
Emit on time, then correct
The disorder bound covers the common delay. Allowed lateness keeps state for updates. The side output keeps the tail.
- The first pane is not postponed by the whole tail.
- Downstream upserts by window key.
- Drops are a metric, not a silent undercount.
Alternative
One huge delay, or processing time
A giant bound makes every result late. Processing time makes a replay rewrite the current minute.
- Users wait for completeness they did not ask for.
- A backlog looks like a traffic spike.
- An idle partition still pins the minimum if you forget idleness.
How a watermark moves
The slowest active input sets the pace. Idle is not the same as slow.
- 1
Each split tracks max event time
Ordering is strongest inside a partition. That is why the watermark starts there. - 2
Emit max minus the bound
Bounded out-of-orderness is the default. Flink emits it about every 200 milliseconds. - 3
Downstream takes the minimum
An operator with several inputs cannot pass the slowest one. An idle split is excluded after a timeout. - 4
Fire, then purge or side-output
Windows with end at or below the watermark fire. State stays until end plus allowed lateness. Later records leave through a side output.
Three clocks
| Time domain | Source | Deterministic on replay? | Late data possible? | Use when |
|---|---|---|---|---|
| Event time | Timestamp inside the record | Yes. Same input, same windows | Yes. Needs watermarks | Billing, analytics, sessions, anything per minute of reality |
| Processing time | The operator's wall clock | No. Depends on lag and speed | No, by definition | Ops rates, timeouts, rate limiting |
| Ingestion time | Assigned at the source operator | Mostly. Fixed after ingest | No | Sources without a trustworthy timestamp |
The interview line: event time makes results a function of the data. Processing time makes them a function of the pipeline's health. A backlog replay with processing-time windows stuffs a day of events into the current minute.
Partition lag that looks like idleness is also a consumer problem. The Kafka view of lag and rebalance is consumer lag and backpressure. A CDC source with its own timestamps is CDC. The watermark still belongs to this lesson.
How watermarks flow
Decisions
- 1
1. Split tracks max event time
- next2. Watermark is max minus bound
- 2
2. Watermark is max minus bound
- next3. Is the split idle?
- ?
3. Is the split idle?
- no4. Min across active inputs
- yes4b. Drop idle from the min
- 4
4. Min across active inputs
- next5. Fire at window end
- 5
4b. Drop idle from the min
- next5. Fire at window end
- 6
5. Fire at window end
- next6. Allowed lateness keeps state
- 7
6. Allowed lateness keeps state
- next7. Later rows to side output
- 8
7. Later rows to side output
Lesson map
Event Time & Watermarks — Allowed Lateness, Side Outputs & Idle Sources
In event-time processing the engine never knows a window is complete. A phone that was in airplane mode can upload yesterday's clicks right now. A watermark is the engine's declaration that it believes no more events older than W will arrive. It drives on-time window firing, timers, and state cleanup. Too aggressive and you drop or mis-count late data. Too conservative and every result is delayed. Allowed lateness keeps window state around a bit longer for corrections, a side output captures anything later than that, and idleness handling stops one silent partition from freezing the whole job.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB split["1. Split tracks max event time"] bound["2. Watermark is max minus bound"] mn["3. Min across inputs, skip idle"] late["4. Fire, lateness, side output"] split -->|1. Split tracks max event time| bound bound -->|2. Watermark is max minus bound| mn mn -->|3. Min across inputs, skip idle| late
- Each source split tracks the max event time it has seen.
- A bounded out-of-orderness strategy emits watermark = max timestamp minus the delay, periodically.
- Downstream operators take the minimum across inputs, so the slowest active partition sets the pace.
- Windows fire when the watermark reaches the window end.
- State is kept until the watermark reaches the window end plus allowed lateness.
- Anything later is too late: dropped by default, or routed to a side output.
Strategies
| Strategy | How it works | Latency | Late-data risk | Pick when |
|---|---|---|---|---|
| Ascending | Watermark equals the max seen. Assumes in-order per split | Lowest | High if anything is disordered | Ordered logs, single-partition CDC |
| Bounded out-of-orderness | Watermark equals max minus D | Plus D | Events delayed more than D are late | The default. Pick D from a high percentile of measured lateness |
| Percentile or heuristic | Track the arrival-delay distribution and set D dynamically | Adaptive | Tunable | Mobile or IoT feeds with variable delay |
| Punctuated | Special records carry the watermark | Source-defined | Low if the producer is honest | The producer knows completeness, such as end of file |
| Perfect | The source guarantees completeness | Varies | None | Bounded input, or ingestion time |
A ten-second window, walked
Bounded out-of-orderness of 5 seconds, allowed lateness of 10 seconds, side output after that, on 10 second tumbling windows.
Press Run. Snippets must be self-contained — no network, files, or native modules.
et=15 wm= 10: on-time fire [0,10) count=4
et=26 wm= 21: on-time fire [10,20) count=3
et=26 wm= 21: purge state for [0,10)
et= 3 wm= 21: TOO LATE -> side output (dead-letter / correction topic)
et=18 wm= 28: late-but-allowed -> re-fire [10,20) count=4
et=45 wm= 40: on-time fire [20,30) count=2
et=45 wm= 40: on-time fire [30,40) count=1
et=45 wm= 40: purge state for [10,20)
et=45 wm= 40: purge state for [20,30)Read it like an incident timeline. Window [0, 10) fires once the max event time reaches 15, so the watermark is 10. Event time 3 arrives when the watermark is already 21, past window end 10 plus allowed lateness 10. Its state is gone, so it goes to the side output instead of disappearing with no trace. Event time 18 is late but within lateness, so [10, 20) fires again with count 4. Downstream must upsert by window key.
Idle sources
An operator's watermark is the minimum of its inputs, so one partition with no traffic pins the watermark forever. Windows never fire, state grows, and dashboards show nothing. Flink's withIdleness marks a split idle after a timeout and drops it from the minimum until it produces again.
Press Run. Snippets must be self-contained — no network, files, or native modules.
--- no idleness handling (window [100,160) never closes) ---
pt=40 wm=94 active=[0,1,2]
pt=41 wm=94 active=[0,1,2]
pt=80 wm=94 active=[0,1,2]
pt=81 wm=94 active=[0,1,2]
--- withIdleness(30s) ---
pt=40 wm=165 active=[0]
pt=41 wm=165 active=[0,1]
pt=80 wm=205 active=[0]
pt=81 wm=205 active=[0,1]Without idleness the watermark sticks at 94 while data flows to 215. With idleness it advances. At processing time 40, partition 1 was also idle, 37 seconds since its last record, so its next record could have been late if the active partition had raced ahead. Idleness trades completeness for progress. Watermark alignment, from Flink 1.15, pauses fast splits that run too far ahead of slow ones and bounds state growth on a skewed backfill.
What to do with late records
| Option | Correctness | Cost | Downstream |
|---|---|---|---|
| Drop them | Silent undercount | None | Accept the loss, and count the drops |
| Larger disorder bound D | Better completeness | Every result waits D | None |
| Allowed lateness L | Corrections for up to L | Window state kept L longer, extra firings | Upsert or retraction |
| Side output | Nothing lost | A separate path | A reconcile job or a correction topic |
| Batch reconciliation | Exact, eventually | A daily job | Two systems, scoped to corrections |
Spark's withWatermark("ts", "10 minutes") sets one query-level watermark. By default it is the min across streams. spark.sql.streaming.multipleWatermarkPolicy=max is opt-in, and late rows beyond the watermark are dropped from stateful aggregations. Kafka Streams uses per-partition stream time plus a grace period on windows.
Pitfalls
In the Python log, name the event time that misses allowed lateness and the event time that only updates a fired window. Then say what the sink must do with the update.
Interview Q&A
What exactly is a watermark?
Answer
A timestamp W flowing through the dataflow that asserts no more records with event time at or before W are expected. It is a heuristic for unbounded sources. Records that violate it are late.
How do you pick the out-of-orderness bound?
Answer
From data: the distribution of event-time delay per source. Set it near a high percentile, balancing latency against completeness, then handle the tail with allowed lateness and a side output.
Why does one idle Kafka partition stall the job, and how do you fix it?
Answer
Operators take the min watermark across inputs. An idle partition never advances, so the min is frozen. Configure idleness so idle splits are temporarily excluded, or emit heartbeat records upstream.
Allowed lateness versus a bigger watermark delay. What is the difference?
Answer
A bigger delay postpones the first result for every window. Allowed lateness emits on time and then corrects with late firings, at the cost of keeping state longer and needing update-capable sinks.
What happens to records later than allowed lateness?
Answer
Dropped by default. Count them. Better: route them to a side output and reconcile in a correction job or a batch pass.
How do watermarks interact with joins and multiple inputs?
Answer
Two-input operators take the min of both inputs' watermarks, so the slower stream controls when interval-join state can be cleaned and when results are final.
Why should watermarks be generated at the source, per partition?
Answer
Ordering is usually strongest within a partition. Per-split watermarks capture that and are then combined by min. Generating them after a shuffle mixes partitions and forces a much larger disorder bound.