In Flink, what do allowedLateness and sideOutputLateData do with records behind the watermark?
answer
- two deadlines, not one
- the first emits, the second cleans up
- a straggler makes it fire twice
- the residue can be caught, not dropped
- every extra minute is state
basics
~20 sallowedLateness keeps a window's state alive past its watermark deadline, so a straggler triggers an extra firing with an updated result. sideOutputLateData routes records too late even for that grace period to a separate stream instead of dropping them silently.
solid answer
~50 sOnce the watermark passes a window's end, the window fires. Without extra configuration its state is then purged and any later-arriving record for that window is dropped, counted only by the `numLateRecordsDropped` metric. `allowedLateness(Duration.ofMinutes(1))` changes the purge deadline to `windowEnd + 1 minute`. Between those two points a straggler is still accepted: the window fires **again**, emitting a corrected result for the same window. Downstream must therefore tolerate multiple results per window key — either an idempotent upsert sink, or a consumer that treats later emissions as replacements. `sideOutputLateData(new OutputTag<>("late"))` catches records that arrive even after the allowed lateness has expired. They never enter the window; they are routed to a side stream you retrieve with `result.getSideOutput(tag)`, so you can log them, alert on the volume, or divert them to a slower reconciliation path. The cost of lateness is state: every window stays in the state backend for the extra duration, multiplied by the number of live keys.
code
java · 11 linesfinal OutputTag<Order> lateTag = new OutputTag<Order>("late-orders") {};
SingleOutputStreamOperator<Revenue> result = orders
.keyBy(Order::getUserId)
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
.allowedLateness(Duration.ofMinutes(1))
.sideOutputLateData(lateTag)
.aggregate(new RevenueAggregate());
result.sinkTo(upsertSink);
result.getSideOutput(lateTag).sinkTo(auditSink);go deeper
Recall that records arriving after the watermark has passed a window are late, and that by default they are simply dropped unless the job is configured otherwise.
Explain the two deadlines — the firing timer at window end and the cleanup timer at window end plus allowed lateness — and that a straggler in between causes an extra firing with an updated result.
Show the downstream consequences you have had to handle: idempotent upsert sinks for repeated firings, state growth from retained windows, and a late side output wired to metrics so silent data loss becomes visible.
Own the completeness policy end to end: what freshness downstream consumers are promised, whether corrections are acceptable or a single final answer is required, and what the reconciliation path is for records that miss every deadline.
## Three fates for a record For a given event-time window, a record can be: 1. **On time** — its timestamp falls in the window and it arrives while the watermark is still before the window's end. It participates normally. 2. **Late but within the grace period** — the watermark has passed the window end, but not by more than `allowedLateness`. 3. **Too late** — the watermark has passed `windowEnd + allowedLateness`. The window's state is gone; the record cannot be incorporated. The default configuration has `allowedLateness = 0`, which collapses cases 2 and 3 into one and makes the outcome a silent drop. ## allowedLateness ```java stream.keyBy(Order::getUserId) .window(TumblingEventTimeWindows.of(Duration.ofMinutes(5))) .allowedLateness(Duration.ofMinutes(1)) .sideOutputLateData(lateTag) .aggregate(new RevenueAggregate()); ``` The window's *firing* timer still sits at `windowEnd - 1`, so the first result is emitted as promptly as the watermark allows. What `allowedLateness` moves is the **cleanup** timer, to `windowEnd + allowedLateness`. Between the two, the window's accumulated state is retained. When a late-but-within-grace record arrives, the event-time trigger fires the window again immediately for that record. The window function is re-evaluated over the full accumulated contents and emits a new result. Critically, this is an *additional* emission, not a retraction — the DataStream API has no retraction concept, so the earlier, incomplete result has already left the job. That forces a downstream contract. Common answers: - Write to an upsert sink keyed by `(window start, key)` — JDBC upsert, search-index document id, Kafka compacted topic with the window key as the message key. Later firings overwrite earlier ones. - Emit the window key alongside the result so consumers can deduplicate. - Accept the duplicates in an append-only sink and resolve them at query time by taking the last write per window. If you cannot tolerate multiple emissions, the alternative is to *not* use allowed lateness and instead widen the out-of-orderness bound so the first firing is already complete — paying uniform latency instead of emitting corrections. ## sideOutputLateData ```java final OutputTag<Order> lateTag = new OutputTag<>("late-orders") {}; // ... DataStream<Order> late = result.getSideOutput(lateTag); late.sinkTo(auditSink); ``` Records that miss even the grace period go here. Note the anonymous-subclass syntax: `OutputTag` needs its generic type at runtime, so `new OutputTag<Order>("late") {}` (with braces) is required. The value of a late side output is mostly operational. A pipeline that silently discards records has no way to answer "are we losing data, and how much?" Routing them out gives you a countable stream: alert on its rate, sample it to see which producers or partitions are responsible, and — if the business cares — feed it into a nightly reconciliation job that patches the affected aggregates. ## The state cost Allowed lateness is not free. Window state for every live key is retained for the extra duration. The rough shape is: ``` retained windows per key ≈ (allowedLateness / window slide) + 1 ``` So a one-minute window with a one-hour lateness holds about sixty windows per key instead of one, and with millions of keys that is the difference between a comfortable RocksDB footprint and an incident. Sliding windows multiply this again, since each record already belongs to several windows. Size lateness against your state backend's capacity, not just against the data's tail. ## Choosing between the two dials There are two knobs that both look like "wait longer": the watermark's out-of-orderness bound, and `allowedLateness`. - The **bound** delays *every* result uniformly and keeps the output single-shot. Use it to cover the common case — the bulk of the out-of-orderness distribution. - **allowedLateness** keeps latency low for the common case and pays only for the tail, at the price of corrections and retained state. Use it to cover stragglers beyond the bound. A reasonable default: bound at a high percentile of measured lateness, a modest allowed lateness sized by what your state backend can hold, and a late side output so the residue is visible rather than lost. ## Related but distinct Do not confuse this with **state TTL**, which expires keyed state in `ValueState`/`MapState` and is a `StateTtlConfig` concern, or with the Table/SQL API's watermark and `TABLE_EXEC_EMIT` late-firing settings, which express the same ideas with different syntax. And note that in `RuntimeExecutionMode.BATCH` there is no lateness at all — input is sorted by time, so nothing can be late. ## Version note This answer assumes Flink 2.3, where `allowedLateness` and the window assigners take `java.time.Duration`. Older code passes Flink's own `Time` class instead; it was deprecated in 1.19 and removed in 2.0, so that code no longer compiles.
- If allowed lateness makes a window fire again, what must the sink support?Idempotent overwrite. The DataStream API emits an additional, corrected result rather than retracting the earlier one, so the sink should upsert on `(window start, key)` — a JDBC upsert, a search-index document id, or a compacted Kafka topic keyed by the window key. An append-only sink will accumulate duplicate rows that must be resolved at read time.
- How is allowedLateness different from simply increasing the out-of-orderness bound?The bound delays every result uniformly and keeps output single-shot; allowed lateness keeps the first result fast and pays only for stragglers, at the cost of extra firings and retained state. Cover the bulk of the lateness distribution with the bound and the tail with allowed lateness plus a side output.
- What is the state cost of a large allowed lateness?Roughly `(allowedLateness / slide) + 1` windows retained per key instead of one. A one-minute window with an hour of lateness holds about sixty windows per key; multiplied by millions of keys, that is a state-backend capacity decision, and sliding windows multiply it again because each record already belongs to several windows.
- How would you know in production that the lateness configuration is wrong?Watch two signals: the window operator's `numLateRecordsDropped` counter, and the throughput of the late side-output stream. A rising drop count with no side output configured means data is vanishing; a heavy side-output stream means the out-of-orderness bound no longer matches the real distribution and should be re-measured.
The out-of-orderness bound is arriving fashionably late to the meeting; allowed lateness is the chair agreeing to reopen the minutes for anyone who turns up in the next ten minutes. After that, your comment goes in the appendix.
saying these in an interview costs you the question
- Thinks allowedLateness delays the first firing of the window
- Expects Flink to retract the earlier result automatically
- Sets hours of lateness without mentioning state growth
- Confuses allowedLateness with the out-of-orderness bound
- Believes late records are logged by default without a side output