You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot 2.3.10 Kafka手动ACK模式下获取重试次数的问题

解决手动ACK下Kafka消息重试次数跟踪问题

在你使用的spring-kafka 2.5.12版本中,DELIVERY_ATTEMPT头确实只支持容器配置的自动重试(比如通过maxAttempts设置的重试),手动调用ack.nack()触发的消息重发,容器不会自动更新这个头的数值,所以每次都会显示为1。针对这个问题,有两种可行的解决思路:

方案一:自定义重试计数(无需升级框架)

通过本地缓存手动维护每个消息的重试次数,具体实现步骤如下:

  1. 定义一个线程安全的缓存,用topic-partition-offset组合作为唯一键存储消息的重试次数:
private val retryCounter = ConcurrentHashMap<String, Int>()
  1. 在监听器方法中,生成消息唯一键并更新重试计数,处理完成或达到最大重试次数后清理缓存:
@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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 22:07:31