Spring Boot 2.3.10 Kafka手动ACK模式下获取重试次数的问题
解决手动ACK下Kafka消息重试次数跟踪问题
在你使用的spring-kafka 2.5.12版本中,DELIVERY_ATTEMPT头确实只支持容器配置的自动重试(比如通过maxAttempts设置的重试),手动调用ack.nack()触发的消息重发,容器不会自动更新这个头的数值,所以每次都会显示为1。针对这个问题,有两种可行的解决思路:
方案一:自定义重试计数(无需升级框架)
通过本地缓存手动维护每个消息的重试次数,具体实现步骤如下:
- 定义一个线程安全的缓存,用
topic-partition-offset组合作为唯一键存储消息的重试次数:
private val retryCounter = ConcurrentHashMap<String, Int>()
- 在监听器方法中,生成消息唯一键并更新重试计数,处理完成或达到最大重试次数后清理缓存:
@KafkaListener( topics = ["my-topic"], groupId = "group-id", containerFactory = "containerFactoryManualAck", clientIdPrefix = "prefix" ) fun processUpdateDispatchStatusEvent( message: String, ack: Acknowledgment, @Header(KafkaHeaders.RECEIVED_TOPIC) topic: String, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) partition: Int, @Header(KafkaHeaders.OFFSET) offset: Long ) { val messageKey = "$topic-$partition-$offset" val currentAttempt = retryCounter.compute(messageKey) { _, count -> count?.plus(1) ?: 1 } println("attempt: $currentAttempt") try { // 执行业务逻辑 ack.acknowledge() retryCounter.remove(messageKey) // 成功消费后移除计数 } catch (e: Exception) { val maxRetry = 3 // 设定最大重试次数 if (currentAttempt >= maxRetry) { // 达到最大重试次数,可转发至死信队列或直接确认 ack.acknowledge() retryCounter.remove(messageKey) } else { ack.nack(10000) } } }
注意:本地缓存方案在服务重启后,未处理完的消息重试次数会重置为1。如果需要持久化计数,可以将缓存替换为Redis等外部存储,但会增加系统复杂度。
方案二:升级spring-kafka版本到2.8+
从spring-kafka 2.8版本开始(对应Spring Boot 2.6及以上版本),官方优化了手动ACK的重试逻辑,调用ack.nack()后,容器会自动更新DELIVERY_ATTEMPT头的数值,无需手动维护计数。但升级需要注意:
- 需同步升级Spring Boot到2.6.x及以上版本,可能涉及其他依赖的兼容性调整
- 升级前需做好充分测试,避免影响现有业务逻辑
内容的提问来源于stack exchange,提问作者Mauricio Avendaño
相关产品推荐
相关产品推荐

