启用Kafka事务后ConsumerInterceptor的onCommit方法未触发问题
在Spring Boot应用中使用Spring Kafka时,未启用生产者事务时,ConsumerInterceptor的onCommit方法能正常触发;但配置transaction-id-prefix启用生产者事务后,onCommit不再被调用,仅onConsume正常执行,@KafkaListener也能接收消息。
未启用事务的正常场景代码
应用启动类
@SpringBootApplication @EnableKafka class Application { @KafkaListener(topics = ["test"]) fun onMessage(message: String) { log.warn("onMessage: $message") } }
拦截器实现
class Interceptor : ConsumerInterceptor<String, String> { override fun onCommit(offsets: MutableMap<TopicPartition, OffsetAndMetadata>) { log.warn("onCommit: $offsets") } override fun onConsume(records: ConsumerRecords<String, String>): ConsumerRecords<String, String> { log.warn("onConsume: $records") return records } }
应用配置
spring: kafka: consumer: enable-auto-commit: false auto-offset-reset: earliest properties: "interceptor.classes": com.example.Interceptor group-id: test-group listener: ack-mode: record
测试代码
@Test fun sendMessage() { kafkaTemplate.send("test", "id", "sent message").get() // 阻塞等待消费者处理 }
此时日志输出符合预期:
onConsume: org.apache.kafka.clients.consumer.ConsumerRecords@6a646f3c onMessage: sent message onCommit: {test-0=OffsetAndMetadata{offset=1, leaderEpoch=null, metadata=''}}
启用生产者事务后的异常场景
添加生产者事务配置:
spring: kafka: producer: transaction-id-prefix: tx-id-
修改测试代码为事务发送:
@Test fun sendMessage() { kafkaTemplate.executeInTransaction { kafkaTemplate.send("test", "a", "sent message").get() } }
此时日志仅输出:
onConsume: org.apache.kafka.clients.consumer.ConsumerRecords@738b5968 onMessage: sent message
原因分析
当启用生产者事务时,Spring Kafka默认会将消费者偏移量提交纳入生产者事务,也就是通过KafkaTransactionManager调用sendOffsetsToTransaction方法提交偏移量,而非调用消费者自身的commitSync/commitAsync方法。
而Kafka原生的ConsumerInterceptor#onCommit方法仅在调用消费者本地的commit*系列方法时才会触发,sendOffsetsToTransaction属于事务API的提交路径,不会触发该拦截器方法。
解决方案
方案1:禁用消费者事务绑定(简单但有一致性风险)
如果不需要将消费者偏移量与生产者事务强绑定,可以配置spring.kafka.listener.transactional=false,让消费者偏移量回到常规的提交路径,从而触发ConsumerInterceptor#onCommit:
spring: kafka: listener: transactional: false
⚠️ 注意:此配置会导致偏移量提交与生产者事务分离,可能出现消息已消费但偏移量未提交(重复消费),或者偏移量已提交但生产者事务回滚(消息丢失)的情况,需根据业务场景权衡。
方案2:使用Spring Kafka的偏移量提交回调替代拦截器
通过配置OffsetCommitCallback来实现原拦截器onCommit中的逻辑,该回调会在偏移量提交(包括事务内的sendOffsetsToTransaction)时触发:
@Bean fun kafkaListenerContainerFactory( consumerFactory: ConsumerFactory<String, String> ): ConcurrentKafkaListenerContainerFactory<String, String> { val factory = ConcurrentKafkaListenerContainerFactory<String, String>() factory.consumerFactory = consumerFactory // 设置偏移量提交回调 factory.setOffsetCommitCallback { offsets, exception -> if (exception != null) { log.error("偏移量提交失败", exception) } else { log.warn("onCommit (回调实现): $offsets") // 这里执行原Interceptor#onCommit中的业务逻辑 } } return factory }
此方案既保留了生产者事务与偏移量的一致性,又能在偏移量提交时执行自定义逻辑。
内容的提问来源于stack exchange,提问作者francis

