Kafka topic分区偏移量缺失、消费组持续滞后问题排查咨询
问题原因解答
固定消费滞后、控制台查不到对应偏移量记录的原因
这个现象是Kafka事务机制下的正常表现,不存在真实的消费堆积,核心逻辑如下:
- Kafka 0.11版本引入事务与幂等生产者机制后,会在分区中写入事务控制消息(事务提交标记、事务中止标记),这类消息是Broker内部使用的元数据消息,会占用正式的分区偏移量,同时会被计入分区的Log End Offset(LEO),但不会作为业务消息返回给普通消费者。
- Flink流应用开启Checkpoint、配合Kafka Sink使用Exactly-Once语义时,Kafka Source默认的消费隔离级别为
read_committed,该级别下消费者只会读取已提交的业务消息,会主动跳过控制消息、未提交的事务消息,因此提交的消费位点(CURRENT-OFFSET)永远不会覆盖控制消息占用的偏移量。 - 你执行
kafka-consumer-groups.sh查询得到的LAG值,是直接用LOG-END-OFFSET - CURRENT-OFFSET做算术计算得到的,没有过滤掉控制消息占用的偏移量,因此会显示存在少量滞后。你观测到新消息可以被正常处理、滞后数值长期不上涨也不消失,完全符合这个特征——Flink已经消费完了所有可读的业务消息,剩下的偏移量位置存的是不可见的控制消息,自然不会被消费,也不会提交对应位点。 - 你用
kafka-console-consumer.sh读不到9566、9567偏移量的内容,也是因为这两个位置存的是控制消息,控制台消费者默认不打印内部消息,所以会直接跳过,看起来就像偏移量断档。如果要验证,可以在启动控制台消费者时加上--property print.internal=true参数,就能看到这些内部控制消息的内容。
偏移量跳变的原因
偏移量跳过单个数值(比如2345直接跳到2347)和生产者写入逻辑直接相关:
- 从Kafka 2.8版本开始,生产者默认开启幂等性(
enable.idempotence=true),如果显式开启事务配置,或者使用Flink Kafka Sink的Exactly-Once语义(默认基于Kafka事务实现两阶段提交),每次事务提交/中止时,Broker都会写入对应的控制消息,这类消息占用偏移量但对业务消费者不可见,就会出现消费时偏移量跳号的现象。 - 如果生产者开启事务后异常崩溃,没有主动提交或回滚事务,等Broker端
transaction.timeout.ms配置的超时时间到达后,会自动中止该事务,同时写入中止标记,这个标记同样会占用偏移量,也会造成偏移量跳变。
注意:如果你的业务监控依赖Kafka消费组的LAG指标做告警,建议针对这种长期存在的固定小数值伪滞后做规则过滤,避免误报。这种控制消息造成的偏移量跳变、固定小lag不会影响业务正常消费,不需要额外处理。
内容的提问来源于stack exchange,提问作者davyjones
相关产品推荐
相关产品推荐

