A Flink event-time job stops emitting window results when one Kafka partition goes quiet. Why?
answer
- records flow, results do not
- the slowest input sets the pace
- silence produces no timestamps
- one builder method fixes it
- and it trades stalls for late records
basics
~20 sWatermarks are generated per input and combined by taking the minimum, so a silent partition's watermark never advances and pins the whole operator's event time. Configure withIdleness on the WatermarkStrategy so idle inputs are excluded from that minimum.
solid answer
~60 sWatermarks are generated per source split and per input channel, and every downstream operator sets its own watermark to the **minimum** across its inputs. A Kafka partition with no traffic produces no records, so its generator never advances its maximum timestamp and its watermark stays where it was. That minimum then pins the operator's event time, no window's end timer is ever reached, and results stop coming — while the job looks perfectly healthy: no backpressure, no failures, checkpoints succeeding. The fix is `.withIdleness(Duration.ofMinutes(1))` on the `WatermarkStrategy`. After that period of silence the stream is marked idle and downstream operators exclude it from the watermark minimum; the moment a record appears it becomes active again. The same shape appears when partition count exceeds source parallelism unevenly, when a `union` combines a high-rate and a low-rate stream, and overnight when a business-hours topic simply has no events. Diagnose it by looking at the per-subtask low watermark in the Web UI or the `currentOutputWatermark` gauge — one subtask will be stuck.
code
java · 6 linesWatermarkStrategy<Order> strategy = WatermarkStrategy
.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((o, ts) -> o.getPlacedAtMillis())
.withIdleness(Duration.ofMinutes(1));
DataStream<Order> orders = env.fromSource(kafkaSource, strategy, "orders");go deeper
Remember that event time only moves when records carry it forward, so a stream with no records cannot advance the watermark and everything downstream waits.
Explain the minimum-across-inputs rule and name withIdleness as the setting that removes a silent stream from that minimum until it produces data again.
Diagnose it from metrics — per-subtask watermark lag with healthy throughput and checkpoints — and discuss the correctness trade idleness makes, including pairing it with a late side output.
Make it a platform default rather than a per-job discovery: standard watermark strategy with idleness, source parallelism aligned to partition count, watermark-lag alerting shipped with every event-time pipeline.
## The symptom A job that has been emitting five-minute aggregates every five minutes stops emitting. Records are still flowing in — Kafka consumer lag is fine, throughput is normal, checkpoints complete, no exception is thrown, no backpressure is reported. Only the sink goes silent, and state size grows steadily because windows keep opening and none close. That combination — data flowing, results not — is nearly always a watermark problem, and the most common cause is an input that has gone quiet. ## Why silence is fatal Event time in Flink advances only through watermarks, and watermarks are derived from record timestamps. No records means no timestamp progress means no watermark progress on that path. Combine that with the minimum rule. Every operator tracks the last watermark received per input channel and adopts the minimum. The `KafkaSource` does the same thing internally across the splits a subtask owns: it runs one `WatermarkGenerator` per split and emits the minimum across them. So a single quiet partition — one split, out of possibly hundreds — holds the entire job's event time at whatever value it last reached. Downstream, that stalled watermark propagates through the shuffle to every window operator subtask, because each of them receives a channel from the stalled source subtask. This is not a bug. Advancing past the silent input would declare `[12:00, 12:05)` complete while that partition might still deliver a 12:03 record, and that record would then be dropped as late. The engine's conservatism is correct; it just needs to be told that the input is genuinely idle rather than merely slow. ## The fix: withIdleness ```java WatermarkStrategy<Order> strategy = WatermarkStrategy .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((o, ts) -> o.getPlacedAtMillis()) .withIdleness(Duration.ofMinutes(1)); ``` After one minute without a record, that stream (or split) is marked **idle**, which is a distinct in-band signal — Flink's `WatermarkStatus`. Downstream operators drop idle channels out of the watermark minimum entirely, so the remaining active channels are free to advance event time. When a record finally arrives, the channel goes active again and rejoins the minimum. Since Flink 2.0 the idleness timer counts only time in which the input could have made progress: time spent backpressured or paused by watermark alignment no longer counts towards the timeout, so a blocked but busy split is not wrongly declared idle. Understand what you have traded. Idleness does not manufacture progress on the quiet stream; it declares that the stream will not object to the rest of the job moving on. If that partition then delivers an old record, it will be behind the advanced watermark and treated as late. Idleness converts a stalled pipeline into a pipeline that may drop a few records — usually the right trade, but a trade nonetheless. Pair it with a late side output if the records matter. Choose the timeout to be comfortably longer than the stream's natural inter-record gap. A one-minute idleness on a topic that legitimately has three-minute gaps will oscillate between idle and active and can produce jittery watermark behaviour. ## The other half of the fix: generate watermarks in the source Where you attach the strategy matters: ```java // per-split watermarking — correct env.fromSource(kafkaSource, strategy, "orders"); // single generator over the interleaved stream — hides the problem env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "orders") .assignTimestampsAndWatermarks(strategy); ``` With `assignTimestampsAndWatermarks` after the source, a subtask reading four partitions computes one maximum over the interleaved records. A quiet partition no longer stalls anything — but that is worse, not better: the watermark now races ahead on the strength of the busy partitions, and when the quiet partition wakes up, its records are silently late. Per-split watermarking inside the source keeps the semantics honest and lets `withIdleness` make the idleness decision explicitly. ## Related shapes of the same failure - **Over-parallelised source.** If source parallelism exceeds the partition count, the extra subtasks own no splits and produce no watermarks at all. The Kafka source does not mark such a reader idle on its own, so it pins the minimum exactly like a silent partition: keep source parallelism at or below the partition count, or rely on `withIdleness`. - **Union of uneven streams.** `streamA.union(streamB)` makes the operator watermark the minimum of both. A low-rate reference stream will throttle a high-rate event stream unless it also carries an idleness setting. - **Business-hours topics.** A topic with no traffic between midnight and 08:00 stalls the pipeline every night and catches up in a burst every morning. - **Broadcast and connected streams.** A `BroadcastConnectedStream` or `connect` has two logical inputs, and the rule is the same minimum. ## Diagnosing it Open the Web UI and look at the low watermark per subtask on the window operator, or scrape the `currentInputWatermark` / `currentOutputWatermark` gauges. One subtask (or one source subtask) will be pinned far behind the others, and the gap will equal the wall-clock time since that input went quiet. Watermark lag per subtask is the metric worth alerting on: it catches this failure and the general "job is falling behind" case with one dashboard.
- What does withIdleness actually cost you?Correctness at the margins. Idleness does not advance the quiet stream's event time; it removes that stream from the watermark minimum so the rest of the job can proceed. If the stream then delivers an old record, it is behind the now-advanced watermark and is late — dropped by default. Pair idleness with `allowedLateness` or a late side output when those records matter.
- Why is generating watermarks inside env.fromSource better than assignTimestampsAndWatermarks after the source?The connector runs one generator per split and takes the minimum across the splits a subtask owns, so a lagging partition is visible and can be handled with idleness or alignment. A generator placed after the source sees one interleaved stream and computes a single maximum, letting busy partitions push the watermark past a slow partition's records, which then arrive late and are dropped silently.
- What happens if source parallelism is set higher than the number of Kafka partitions?The surplus subtasks own no splits, read nothing, and therefore produce no watermark progress. The Kafka source does not mark a split-less reader idle by itself, so without `withIdleness` those subtasks stall downstream event time exactly like a silent partition. The robust configuration is source parallelism at or below the partition count, with idleness as the safety net.
- Which metric would you alert on to catch this before a user does?Watermark lag per subtask — wall clock minus `currentOutputWatermark` — on the window operator and the source. A stall shows up as one subtask's lag growing linearly with time while the others stay flat, which distinguishes it from a whole-job backlog where every subtask's lag grows together.
A convoy moves at the speed of its slowest vehicle. If one vehicle has parked for good, you either take it off the roster or the convoy never arrives.
saying these in an interview costs you the question
- Blames backpressure or consumer lag without checking the watermark
- Says Flink automatically skips idle inputs by default
- Suggests switching to processing time as the fix
- Proposes to advance the watermark from wall clock without mentioning late records
- Thinks a single slow partition only delays its own keys