How to fix Flink windows that never fire or drop late events
FlinkAn event-time Flink job is RUNNING, the source is reading records from Kafka, and the window operator emits nothing. Or the job does emit results, but the counts come up short against the source system, and nothing in the logs explains where the missing events went. Both symptoms usually trace back to the watermark.
A watermark is Flink’s signal of progress in event time: a declaration that, by that point in the stream, all events up to a certain timestamp should have arrived, as the timely stream processing documentation puts it. An event-time window fires only when the watermark passes the end of the window, and an event that arrives after that point is late. A window that never fires means the watermark is not advancing far enough. Events that vanish mean they arrived behind the watermark and were dropped or routed somewhere nobody reads.
This page covers event-time windows in DataStream jobs, and in Flink SQL jobs where it says so. It was checked against the Apache Flink 2.3 documentation, the stable release at the time of writing, including the Kafka connector page.
Not this problem?
- Windows fire, but later and later as load rises. The job is more likely backpressured, and the guide to finding the bottleneck behind Flink backpressure is the better starting point.
- The question is how far a Kafka consumer is behind the topic, not how far event time has moved. See Kafka consumer lag: how to find the cause and monitor it.
For how this fits into running Flink in production more widely, see the complete Flink guide.
Rule out a window that is never meant to fire
Before reading any watermark, check the window assigner. The windows documentation says every event-time window assigner uses an EventTimeTrigger by default, which fires once the watermark passes the end of the window. The default trigger of the GlobalWindow is the NeverTrigger, which never fires.
.window(GlobalWindows.create()) // fires only with a custom trigger
A job built on GlobalWindows with no custom trigger produces no output whatever the watermark does. The same documentation notes that no data is ever considered late in a global window, because its end timestamp is Long.MAX_VALUE. If the job uses tumbling, sliding or session event-time windows, the cause is in the watermark, and the next section reads it.
A related check is the notion of time itself. Processing-time windows fire on the wall clock and are unaffected by everything below. The timely stream processing documentation states that processing time gives the best performance and the lowest latency but does not provide determinism, while event time produces correct and consistent results with out-of-order or late events, or when historic data is reprocessed. Switching a stuck job to processing time makes the windows fire, but it changes what the results mean.
Read the watermark
Watermarks are generated at or directly after the sources, and each parallel source subtask generates its own. An operator that consumes more than one input, which includes every operator after a keyBy, takes the minimum of its inputs’ event times as its own. The window operator’s watermark therefore can only be as far along as the slowest source subtask. The diagnosis starts by finding that subtask.
Three metrics from the Flink metrics documentation carry it. The first is the last watermark the operator has received, read on the window operator. The second is the last watermark an operator has emitted, read on each source subtask. Both are in milliseconds. The third counts records the window operator has dropped because they arrived late.
currentInputWatermark
currentOutputWatermark
numLateRecordsDropped
Flink’s monitoring REST API is the same API that Flink’s own dashboard uses, and the REST API documentation states that it is designed for custom monitoring tools as well. It returns the watermark of every subtask of a task in one call, which is the quickest way to see which subtask is lowest.
GET /jobs/:jobid/vertices/:vertexid/watermarks
The vertex IDs come from the job details call, whose response lists every vertex of the job with its ID and name.
GET /jobs/:jobid
Call the watermarks endpoint for the source vertex and for the window vertex, a minute apart. If every source subtask advances and the window vertex matches the lowest of them, watermarks are flowing and the question is lateness. If one source subtask is stuck, or never reports a value, the question is idleness or timestamps. The metrics documentation also lists per-split metrics for sources, which show a single stuck Kafka partition directly where the source reports them. They give each split’s last watermark and the time it has been marked idle by idleness detection.
watermark.currentWatermark
watermark.idleTimeMsPerSecond
Diagnose the cause from the metrics
The table maps what the three metrics show to the usual cause and its fix, in the order the sections below take them. Change one thing at a time and read the same metric again after each change.
Each cause, its proof and its fix
One Kafka partition has no traffic
The generating watermarks documentation calls this an idle input. If one partition carries no events for a while, its watermark generator gets no new information, and the watermark downstream is held back because it is computed as the minimum over all the parallel watermarks. The other partitions keep carrying events, so records keep flowing while no window ever fires.
The proof is a source subtask whose currentOutputWatermark stops at one value while its peers move on, and whose idle time is high. The back pressure documentation defines three per subtask metrics that add up to approximately 1000 ms at any point in time.
idleTimeMsPerSecond
backPressuredTimeMsPerSecond
busyTimeMsPerSecond
Idle time is the time a subtask spent waiting for something to process. A source subtask that is idle rather than back pressured is waiting for data, which points at the partition and not at the job’s throughput.
The fix is idleness detection on the watermark strategy. An input that has seen no records for the timeout is marked idle and stops holding back the watermark.
WatermarkStrategy
.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withIdleness(Duration.ofMinutes(1));
For Flink SQL, the table configuration documentation provides the same behavior as an option, with a default of 0, which leaves idleness detection off.
table.exec.source.idle-timeout: 1 min
Why one partition is empty in the first place is a Kafka question. A producer keyed on a field with few distinct values sends most records to a few partitions and leaves the rest without traffic, which the partition skew section of how to diagnose and fix an unbalanced Kafka cluster covers. Idleness detection keeps the job moving, but the empty partition stays empty until the key changes.
More source subtasks than Kafka partitions
A second form of idleness has nothing to do with traffic. The Flink Kafka connector documentation describes the source’s split enumerator assigning partitions to readers, distributed across subtasks round-robin. Flink, not Kafka’s consumer group protocol, decides which subtask reads which partition. When the source parallelism is higher than the partition count, some subtasks get no partition at all.
The same documentation states that the Kafka source does not go automatically into an idle state in that case, and that the parallelism has to be lowered or an idle timeout added to the watermark strategy. Compare the source parallelism with the number of partitions in the topic, then check the watermarks endpoint for source subtasks that never report a value. A source parallelism above the partition count is the condition the connector documentation describes.
The connector documentation also states that the source commits offsets to Kafka only to expose progress for monitoring and does not rely on them for fault tolerance. A Kafka consumer lag view of the job’s group can therefore look healthy while the job’s watermark is stuck. Lag shows how far the reads are behind the topic, and says nothing about event time.
Timestamps or the watermark strategy are wrong
If the watermark moves but its value makes no sense, the timestamps are the problem. The generating watermarks documentation states that both timestamps and watermarks are milliseconds since the Java epoch. A TimestampAssigner that returns seconds moves event time forward by one millisecond for every real second, so a five-minute window needs about three and a half days of data before the watermark passes its end. Convert the value to milliseconds and compare currentOutputWatermark on the source with the current time in milliseconds.
.withTimestampAssigner((event, recordTs) -> event.getEventTimeMillis())
The same documentation describes two places to set a strategy. Setting it directly on the source lets the source track watermarks per partition, and with the Kafka source the watermarks are generated per Kafka partition and merged as they are on a shuffle. Setting it later, after arbitrary operations, through assignTimestampsAndWatermarks, should be reserved for when the source option is not available. A strategy applied after the source sees records from several partitions interleaved, so a pattern that holds within each partition no longer holds in the merged stream.
env.fromSource(kafkaSource, strategy, "events"); // per-partition watermarks
For a Kafka source with no TimestampAssigner, the connector documentation states that the timestamp embedded in the Kafka record is used as the event time. If producers set that timestamp at send time rather than from the event, the windows group events by when they were produced, not when they happened.
The out-of-orderness bound is too small or too large
The built-in watermark generators documentation describes forBoundedOutOfOrderness as a watermark that lags the highest timestamp seen by a fixed time. That time is the maximum an element can be late before it is ignored when the result for its window is computed.
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10))
A bound that is too small shows up as numLateRecordsDropped climbing on the window operator, because events that arrive more out of order than the bound reach the window after the watermark has passed it. A bound that is too large shows up as a watermark that advances steadily but trails the newest events by the whole bound, so each window fires that much later than expected. In a test with a handful of events and a bound of several minutes, the watermark may never pass the end of the first window, which looks like a window that never fires. The timely stream processing documentation names the trade-off: delaying the watermark by too much causes too much delay in evaluating event-time windows.
The bound belongs at the out-of-orderness the data really has, and allowed lateness, in the next section, handles the tail beyond it. The watermark emission interval sits in the Flink configuration and defaults to 200 ms. It rarely needs changing for this problem.
pipeline.auto-watermark-interval: 200 ms
Late events are dropped because allowed lateness is 0
The windows documentation states that by default late elements are dropped when the watermark is past the end of the window, and that allowed lateness defaults to 0. An event that arrives after the watermark passes the end of its window, but before it passes the end plus the allowed lateness, is still added to the window. With the default EventTimeTrigger, that late event makes the window fire again.
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
.allowedLateness(Duration.ofMinutes(2))
Two consequences come with it, both from the same documentation. Flink keeps each window’s state until its allowed lateness expires, so a long lateness across many keys grows window state, and with it checkpoint size. If checkpoints start to slow after the change, the guides to Flink checkpoints that fail or time out and Flink state that keeps growing pick it up from there. Late firings also emit updated results for a window that already fired, so the stream carries more than one result for the same window, and the sink has to upsert or deduplicate them.
Late events go to a side output nobody reads
The windows documentation describes sideOutputLateData, which turns the late elements that would be discarded into a separate stream. That stream exists only if the job reads it with getSideOutput and sends it somewhere.
OutputTag<Event> late =
new OutputTag<Event>("late-data") {};
// keyBy, window and allowedLateness as above
.sideOutputLateData(late)
.aggregate(new CountAgg());
result.getSideOutput(late).sinkTo(lateSink);
In the WindowOperator source code at the 2.3.0 release, a late element goes to the side output when a late-data tag is set, and numLateRecordsDropped is incremented only when no tag is set. A job with the tag set and no consumer of the side output reports zero dropped records while late events disappear. Check the job code. Search for sideOutputLateData and confirm that each tag has a matching getSideOutput with a sink.
Apply the fix without losing window state
Changing the watermark strategy, the bound, the allowed lateness or the side output is a code change, and this applies to every cause above. The job keeps its window state only if it is redeployed from a savepoint. The procedure is in how to upgrade a Flink job from a savepoint without losing state.
Where Flex fits
This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. The comparison of tools for self-managed Apache Flink sets Flex beside the Flink web UI and the REST API.
Diagnosis. The Flex job documentation describes a Topology tab that shows the job graph. For a selected task it lists a Watermarks view with the current watermark value for each subtask, subtask metrics for each parallel instance, and a backpressure status of OK, LOW or HIGH with a backpressure level percentage. Those are the per-subtask watermark comparison from the section above and the idle versus back pressured check for a source subtask, in one screen. The per-subtask watermark view is part of Flex.
Redeploying the fix. The same job documentation describes Stop, Cancel, Savepoint and Checkpoint actions on a selected job, and a Submit form for a JAR file that accepts a savepoint path and a claim mode, so a stopped job can be resubmitted from its savepoint.
What Flex does not show. The Flex Prometheus metrics glossary lists job, JobManager and TaskManager metrics such as read and write records, job state and minutes inactive. As of Flex 96.5, it has no watermark metric and no late-record metric. Alerting on a stalled watermark or on numLateRecordsDropped needs Flink’s own metric reporters. Flex also does not change a watermark strategy, an out-of-orderness bound or an allowed lateness, which live in the job’s code.
FAQ
Why does a Flink window not fire when data is flowing?
Because the window fires on the watermark, not on arriving records. The watermark of the window operator is the minimum over all its inputs, so one idle Kafka partition, or a source subtask with no partition assigned, holds it back. Compare the per-subtask watermarks of the source and add withIdleness if one is stuck.
What does numLateRecordsDropped count?
The number of records the window operator dropped because they arrived after the watermark passed the end of their window plus the allowed lateness. It stays at 0 when a late-data side output is set, because those records go to the side output instead.
Is it better to raise the out-of-orderness bound or set allowed lateness?
The bound delays every window by its full length, while allowed lateness lets the window fire on time and fire again for each late event within the lateness. A bound sized for the usual disorder, with allowed lateness for the rare late tail and a side output beyond that, keeps results timely without losing events silently.
Does withIdleness lose data from the idle partition?
No. The partition only stops holding back the watermark while it carries no records. When it carries events again it takes part in the watermark once more, and any of its events that are already behind the watermark are late by the same rule as any other event.
Related reading
- Apache Flink: the complete guide
- How to find the bottleneck behind Flink backpressure
- How to fix Flink checkpoints that fail or time out
- How to upgrade a Flink job from a savepoint without losing state
- How to find and fix Flink state that keeps growing
- How to diagnose and fix an unbalanced Kafka cluster