Spring Kafka消费者抛异常时消费lag未上涨的原因及处理
Spring Kafka 2.6.X 异常场景下消费Lag未上涨问题解析
问题背景
使用Spring Kafka 2.6.X版本,已将ENABLE_AUTO_COMMIT_CONFIG设置为false。当消费消息的消费者主动抛出异常时,预期因未手动提交offset,消费lag会上涨,但实际并未上涨。相关配置与消费者代码如下:
class KafkaConsumerConfig { @Bean fun multiTypeConsumerFactory(): ConsumerFactory<String, Any> { val props = HashMap<String, Any>() props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = "localhost:9092" props[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = StringDeserializer::class.java props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = JsonDeserializer::class.java props[ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG] = false return DefaultKafkaConsumerFactory(props) } @Bean fun multiTypeKafkaListenerContainerFactory(): ConcurrentKafkaListenerContainerFactory<String, Any> { val factory = ConcurrentKafkaListenerContainerFactory<String, Any>() factory.setConsumerFactory(multiTypeConsumerFactory()) return factory } } @Component @KafkaListener(topics = ["testTopic"], groupId = "group_sdas") class Consumer { @KafkaHandler fun handleFoo(message: Foo) { println("Received Message in group foo: ${message}") } @KafkaHandler fun handleBar(message: Bar) { throw RuntimeException("ssstest sss") } }
默认错误处理器的处理逻辑
Spring Kafka 2.6.X在未自定义错误处理器时,默认使用SeekToCurrentErrorHandler,核心处理逻辑如下:
- 当消费者抛出异常时,处理器会将当前消费分区的偏移量重置到本次消费前的位置,触发消息的重新消费
- 由于offset既未被手动提交,也没有被向前推进,Kafka集群记录的消费offset始终停留在异常消息之前的位置,因此消费lag不会出现预期的上涨(lag计算公式为
最新消息offset - 消费offset,消费offset未变动,lag自然不会因这条异常消息增加) - 默认配置下,处理器会对异常消息重试10次,重试失败后,若配置了死信队列(DLQ)则将消息转发至DLQ,否则会停止该分区的消费
问题根因
经排查,即使抛出异常消费lag仍未上涨的核心原因是:消息Header中未包含类型标识(默认key为__TypeId__)。
代码中使用@KafkaHandler实现多类型消息消费,Spring Kafka依赖消息Header中的类型标识来匹配对应的处理方法,当缺少该Header时,消息无法匹配到任何@KafkaHandler方法,会直接抛出类型匹配异常,此时SeekToCurrentErrorHandler立即重置offset,导致这条消息被反复重试,但消费offset始终未推进,因此lag不会上涨。
内容的提问来源于stack exchange,提问作者samsamsamsmasma
相关产品推荐
相关产品推荐

