How to fix Kafka consumer lag that looks wrong for a Flink job
FlinkThe Kafka dashboard says a Flink job is falling behind, yet the job looks healthy. The consumer group’s lag keeps growing, or it sits flat for minutes and then drops in one step, or it jumps back up after a restart. Sometimes the group shows no committed offsets at all, or Kafka tooling lists it apart from the normal consumer groups.
Consumer lag is the distance between the newest offset in a partition and the offset a consumer group has committed for it. For a Flink job, those committed offsets are written by the Flink Kafka source only when a checkpoint completes, so the lag Kafka reports is only as fresh as the last completed checkpoint.
The Flink Kafka connector documentation says the source commits the current consuming offset when checkpoints are completed, and that it does not rely on committed offsets for fault tolerance. Committing exists only to expose the job’s progress for monitoring. Flink restores from the offsets stored in its own checkpoints, so Kafka-side lag can be wrong while the job is reading exactly where it should.
This page shows how to tell real lag from checkpoint lag, which Flink metrics show the job’s true position, and what to change for each cause. If the lagging consumer is an ordinary Kafka consumer group and not a Flink job, the guide to Kafka consumer lag is the better starting point.
For how this fits into running Flink in production more widely, see the complete Flink guide.
How the Flink Kafka source commits offsets
First, offsets are committed on checkpoint completion. With checkpointing enabled, nothing reaches Kafka between checkpoints. With checkpointing disabled, the source falls back to the Kafka consumer’s own periodic commits, configured by two consumer properties.
enable.auto.commit
auto.commit.interval.ms
Second, committing can be switched off. The source has a property that controls whether consuming offsets are committed to the brokers on checkpoint at all.
commit.offsets.on.checkpoint
Third, a commit failure does not affect the job. The connector documentation says a failed commit does not affect the integrity of Flink’s checkpointed partition offsets, because committing back to Kafka is only a means to expose consumer progress. The connector keeps a count of each outcome.
KafkaSourceReader.commitsSucceeded
KafkaSourceReader.commitsFailed
The source also does not join a consumer group in the usual way. The Flink split enumerator assigns Kafka partitions to source readers itself, in round-robin style, and the reader’s KafkaPartitionSplitReader source code calls the Kafka consumer’s assign method rather than subscribe. The group id is used to store committed offsets, but the group has no active members from Kafka’s point of view. Lag tooling therefore has to monitor these consumers differently from a normal consumer group, which is why some tools list them separately or show the group as empty.
Tell real lag from checkpoint lag
Start by ruling out the case where the job really is behind. If the Flink-side metrics show the source is close to the end of each partition, the Kafka-side figure is stale and the cause is in checkpointing or configuration. If the Flink-side metrics also show a growing backlog, the lag is real.
Compare the current and committed offsets
The Kafka source reports its own position in Flink’s metric system for each partition, next to the offset it last committed. The connector documentation lists the two offset gauges with the topic and partition as variables, and the metrics class in the connector source code gives the names a metric reporter shows, in the singular and scoped by topic and partition.
KafkaSourceReader.topic.<topic_name>.partition.<partition_id>.currentOffset
KafkaSourceReader.topic.<topic_name>.partition.<partition_id>.committedOffset
The first is the consumer’s current read offset and the second is the last offset successfully committed to Kafka. The gap between them is the progress the job has made since the last successful commit. A large gap with a healthy job points at checkpoints, not at throughput.
Measure the backlog from Flink
Two further metrics measure the backlog from the job’s side. The connector documentation describes pendingRecords as the number of records the source has not fetched yet, for example the records available after the consumer offset in a Kafka partition. The Kafka consumer’s own metrics are registered under the KafkaSourceReader.KafkaConsumer group when register.consumer.metrics is on, which the documentation says is the default.
pendingRecords
KafkaSourceReader.KafkaConsumer.records-lag-max
The Apache Kafka monitoring documentation defines records-lag-max as the maximum lag in records for any partition in the window, and notes that it is based on the current offset and not the committed offset. The Flink-side figures therefore follow what the job has read, while the consumer group lag in Kafka follows what was last committed.
Measure lag in time
For lag in time rather than records, the connector documentation lists two gauges. The first is the time from a record’s event timestamp to when the source emits it, and the second is how far the watermark trails the wall clock.
currentEmitEventTimeLag
watermarkLag
A growing watermarkLag with a small record backlog usually points at event time and idle partitions, not at offsets, which the guide to Flink windows that never fire or drop late events covers.
Fix each cause of Kafka lag that looks wrong
The job is really behind
The Flink backlog metrics grow along with the Kafka lag, and upstream tasks report backpressure. The back pressure monitoring documentation says every subtask reports the time it spends busy, back pressured and idle, adding up to about 1000 ms per second.
busyTimeMsPerSecond
backPressuredTimeMsPerSecond
idleTimeMsPerSecond
The guide to the bottleneck behind Flink backpressure finds the operator that is holding the job back, and the guide to rescaling a Flink job that cannot keep up covers adding parallelism once that operator is known.
Checkpoints are failing
The Kafka lag grows or stays flat while the job keeps running and its own metrics show it close to the end of each partition. No checkpoint has completed recently, so no offsets have been committed.
The checkpoint monitoring documentation describes the web UI’s Checkpoints tab, which shows the latest completed checkpoint, the latest failed checkpoint and the count of failed checkpoints since the job started. The same statistics come from the REST API.
GET /jobs/<job-id>/checkpoints
A job can keep running through failed checkpoints when it is configured to tolerate them. The checkpointing documentation says the tolerable failure number defines how many consecutive checkpoint failures are tolerated before the job fails over, and that its default of 0 means no checkpoint failures are tolerated and the job fails on the first reported failure. The documentation limits this to three failure reasons: an IOException on the Job Manager, a failure in the async phase on the Task Managers and a checkpoint expiring on a timeout.
execution.checkpointing.tolerable-failed-checkpoints
With a non-zero value, checkpoints can fail for a long time while the job stays up and the Kafka lag stays frozen. Fix the checkpoint failure itself. The guide to checkpoints that fail or time out walks through which stage is slow and the change to try first. Long checkpoint duration and backpressure are the per-job signals to watch together. The checkpointing under backpressure documentation says that under heavy backpressure the time to propagate checkpoint barriers can dominate a checkpoint’s end-to-end time.
Commits fail while checkpoints succeed
The Kafka lag stays frozen even though checkpoints complete. The connector documentation says a commit failure does not affect the integrity of Flink’s checkpointed offsets, so the job keeps its true position and only the Kafka view goes stale. The two commit counters show which case applies. A rising failure count with a flat success count means the commits to Kafka are failing and the checkpoints are not.
KafkaSourceReader.commitsSucceeded
KafkaSourceReader.commitsFailed
The documentation counts these only when offset committing is turned on and checkpointing is enabled. It does not list the causes of a failed commit, so check the task manager logs for the commit error.
The checkpoint interval is long
The Kafka lag climbs steadily and drops each time a checkpoint completes, in a sawtooth. Flink’s own metrics stay low. This is expected behaviour, because commits happen only at checkpoint completion.
The checkpointing documentation sets the base interval with one key and the minimum pause after a completed checkpoint with another. The minimum pause also means the effective interval is never smaller than that pause.
execution.checkpointing.interval
execution.checkpointing.min-pause
A shorter interval makes the Kafka-side figure fresher at the cost of more checkpoint work. If the sawtooth is only a monitoring problem, the cheaper fix is to alert on pendingRecords or records-lag-max from Flink and treat the Kafka figure as a delayed view.
The job was restored from an older savepoint or started fresh
The Kafka lag jumps up after a restart and then falls. The connector documentation describes how the source’s split state stores the current consuming offset of each partition and turns it into the split’s starting offset when the reader is snapshotted. A restored job therefore resumes from the offsets in the checkpoint or savepoint it was restored from. If that savepoint is older than the last commit, the job re-reads records and the next completed checkpoint commits the older offsets back to Kafka.
A job started without any state takes its starting position from the offsets initializer. The connector documentation says that if no initializer is specified, OffsetsInitializer.earliest() is used, so a fresh start with default settings reads each partition from the beginning. To continue from the group’s committed offsets on a fresh start, set the initializer explicitly.
OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST)
The guide to upgrading a Flink job from a savepoint covers restoring from the right savepoint so the job does not fall back to the initializer by accident.
The group shows no committed offsets
Three causes are common. No checkpoint has completed since the job started, so nothing has been committed. Committing on checkpoint is switched off with commit.offsets.on.checkpoint set to false. Or the offsets were committed once and have since expired.
Kafka expires offsets for consumers like the Flink source on a different rule from normal groups. The Kafka broker configuration documentation says that for standalone consumers using manual assignment, offsets expire once the retention period has passed since the last commit. The default is 10080 minutes, which is seven days.
offsets.retention.minutes
A job whose checkpoints have failed for longer than that period loses its committed offsets in Kafka, even though Flink still holds its position in its own checkpoints. The committed offsets themselves outlive the job’s consumers, so while they exist, lag for a group with no active members can still be measured from them.
Two jobs share one group.id
The committed offsets move back and forth, and the lag looks erratic. The Kafka consumer configuration documentation describes group.id as the string that identifies the consumer group, and Kafka-based offset management stores committed offsets under that group. Two Flink jobs that read the same topic with the same group.id write to the same offsets, so the group shows whichever job committed last.
group.id
Give each job its own group.id. Each job’s own per-partition currentOffset gauge will show it advancing normally, which confirms the cause.
Why resetting the Kafka group does not move the Flink job
Changing committed offsets in Kafka does not change where a running or restored Flink job reads, because the job resumes from its own checkpointed state. The reset only takes effect for a job started without state and with the committed offsets initializer above.
The Kafka operations documentation also says to make sure the consumer instances are inactive before running an offset reset with the consumer groups tool.
kafka-consumer-groups.sh --reset-offsets
For a Flink source that means stopping the job, not just checking the group’s member list, because the group shows no members while the job is reading. To move a Flink job to a different position, stop it, then start it without state from a chosen initializer, or restore from a savepoint taken at the right point.
Where Kpow and Flex fit
This page is published by Factor House, which makes Kpow for Kafka and Flex, its Flink job management and monitoring product, so weigh this section with that in mind. For other options, see the comparison of tools for self-managed Apache Flink and the list of tools to monitor Kafka consumer lag.
What Kpow shows on the Kafka side. The Kpow groups documentation says Kpow identifies a consumer group as simple when it uses manual partition assignment with no group coordination, names some Flink jobs as one of the workloads that do this, and shows such groups in their own Simple consumers tab. The same page says Kpow tracks lag for EMPTY consumer groups, calculated from the start and end offsets of each assignment and combined at group, broker or topic level. The Kpow Prometheus metrics glossary lists group_offset_lag as the total lag of all assignments of a group. The groups documentation also says offset changes to simple consumers apply immediately in Kafka’s offset store but Kafka does not notify running simple consumers, and recommends scaling them to zero before a reset, which matches the Flink advice above to stop the job first.
What Flex shows on the Flink side. The Flex jobs documentation describes a Checkpoints tab with counts of Triggered, Completed, In progress, Restored and Failed checkpoints, a history table with each checkpoint’s status and duration, and the job’s checkpoint configuration including interval, min_pause and timeout. The job overview shows the last checkpoint, and the task views show backpressure level and watermarks. These are the screens that prove or rule out the checkpoint branches above. The Flex Flink cluster configuration documentation says Flex connects to the Flink REST API endpoint set in FLINK_REST_URL, the API that the Flink REST API reference says is used by Flink’s own dashboard and is designed for custom monitoring tools as well.
What stays in Flink. The Flex documentation does not describe showing pendingRecords, the offset gauges or the Kafka consumer’s records-lag-max, so the job’s true read position comes from Flink’s metric reporters, as the per-partition currentOffset and committedOffset gauges and the commit counters above. Neither product changes how the Flink source commits offsets.
To check the latest completed checkpoint and the backpressure of a job when Kafka lag looks wrong, open the job in Flex. Flex reads both from the Flink REST API endpoint.
FAQ
Why is Kafka consumer lag higher than Flink’s own lag metrics?
The Kafka figure is measured against the offsets the Flink source last committed, and the source commits only when a checkpoint completes. By the Kafka documentation’s own definition, records-lag-max is based on the current offset and not the committed one. The difference between the two is the progress made since the last completed checkpoint.
Does Flink use the offsets committed to Kafka when it restarts?
Not when it restores from a checkpoint or savepoint. The Kafka connector documentation says the source does not rely on committed offsets for fault tolerance and restores from the offsets in its own state. Committed offsets are used only by a job started without state and configured with the committed offsets initializer.
Why does the Flink job’s consumer group show no members?
The Flink split enumerator assigns partitions to readers itself and the reader uses the Kafka consumer’s assign method, so the job never joins the group as a member. The group id is used only to store committed offsets. Tools that list groups by membership show it as empty or as a simple consumer.
How often does a Flink job commit offsets to Kafka?
Once per completed checkpoint, when checkpointing is enabled and commit.offsets.on.checkpoint is not set to false. Without checkpointing, the source falls back to the Kafka consumer’s automatic commits configured by enable.auto.commit and auto.commit.interval.ms.
Related reading
- Apache Flink: the complete guide
- How to fix Flink checkpoints that fail or time out
- How to find the bottleneck behind Flink backpressure
- How to fix Flink windows that never fire or drop late events
- How to upgrade a Flink job from a savepoint without losing state
- Best tools to monitor Kafka consumer lag