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

底层问题修复后如何重新处理Kafka DLT消息?

Kafka DLT消息重处理方案优化问题

问题背景

我配置了一个带有@RetryableTopic注解的KafkaListener,会在将引发处理失败的数据发送至DLT前进行几次重试,示例代码如下:

@KafkaListener(
    groupId = "${kafka.stuff.groupId}",
    topics = ["${kafka.stuff.topic}"],
    concurrency = "${kafka.stuff.concurrency}",
    autoStartup = "${kafka.stuff.auto-start:true}"
)
@RetryableTopic(
    attempts = "#{'${kafka.stuff.attempts}'}",
    backoff = Backoff(
        delayExpression = "${kafka.stuff.backoff.delay}",
        multiplierExpression = "${kafka.stuff.backoff.multiplier}"
    ),
    autoCreateTopics = "#{'${kafka.stuff.auto-create}'}",
    topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
    dltStrategy = DltStrategy.FAIL_ON_ERROR
)
fun consumeStuff(
    @Payload stuff: StuffDto,
    @Header(KafkaHeaders.RECEIVED_KEY) key: String,
    @Header(KafkaHeaders.RECEIVED_PARTITION) partition: Int
) = runBlocking {
    log.info("{} stuff received {}", LogSymbols.START_SYMBOL, key)
    process(stuff, key)
} 

消息进入DLT后,我不希望在问题修复前自动重处理;而修复部署完成后,我希望能重新处理这些消息。因此我的初步方案是配置一个默认禁用的DLT监听器,示例代码如下:

@KafkaListener(
    groupId = "${kafka.stuff.groupId}",
    topics = ["${kafka.stuff.topic}-dlt"],
    concurrency = "${kafka.stuff.dlt.concurrency}",
    autoStartup = "${kafka.stuff.dlt.auto-start:false}" //disabled by default -> enable it on-demand
)
fun processStuffInDlt(
    @Payload stuff: StuffDto,
    @Header(KafkaHeaders.RECEIVED_KEY) key: String,
    @Header(KafkaHeaders.RECEIVED_PARTITION) partition: Int
) = runBlocking {
    log.info("{} Re-process stuff {} from DLT", LogSymbols.START_SYMBOL, key)
    process(stuff, key)
}

但我预见到该方案存在至少两个问题:

  1. 消息重处理完成后,processStuffInDlt监听器会保持活跃,导致后续新的问题消息被重复处理——能否无需重新部署即可按需启停该监听器?
  2. 如果某条消息仍引发处理失败,processStuffInDlt会重试几次后确认该消息,导致后续无法再重处理——能否在失败时阻止消息确认?

另外还有一个更宽泛的问题:有没有更优的方案来实现问题修复后重新处理DLT中引发失败的消息?


解决方案

1. 无需重新部署,按需启停DLT监听器

可以通过Spring Kafka提供的KafkaListenerEndpointRegistry实现动态控制:

  • 首先给DLT监听器添加唯一id标识:
@KafkaListener(
    id = "stuff-dlt-listener",
    groupId = "${kafka.stuff.groupId}",
    topics = ["${kafka.stuff.topic}-dlt"],
    concurrency = "${kafka.stuff.dlt.concurrency}",
    autoStartup = "${kafka.stuff.dlt.auto-start:false}"
)
  • 注入KafkaListenerEndpointRegistry,编写启停方法:
@Autowired
private lateinit var registry: KafkaListenerEndpointRegistry

// 启动DLT监听器
fun startDltListener() {
    registry.getListenerContainer("stuff-dlt-listener")?.start()
}

// 停止DLT监听器
fun stopDltListener() {
    registry.getListenerContainer("stuff-dlt-listener")?.stop()
}
  • 可以结合Spring Boot Actuator暴露自定义REST接口,或者通过运维工具调用这些方法,实现无需重启服务即可启停监听器。

2. 处理失败时阻止消息确认,保留重处理能力

需要关闭自动提交,改用手动确认机制:

  • 配置手动提交的容器工厂:
@Bean
fun manualAckKafkaListenerContainerFactory(): ConcurrentKafkaListenerContainerFactory<String, StuffDto> {
    val factory = ConcurrentKafkaListenerContainerFactory<String, StuffDto>()
    factory.consumerFactory = consumerFactory()
    factory.containerProperties.ackMode = ContainerProperties.AckMode.MANUAL
    return factory
}

@Bean
fun consumerFactory(): ConsumerFactory<String, StuffDto> {
    val props = HashMap<String, Any>()
    props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = "${kafka.bootstrap-servers}"
    props[ConsumerConfig.GROUP_ID_CONFIG] = "${kafka.stuff.groupId}"
    props[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = StringDeserializer::class.java
    props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = JsonDeserializer::class.java
    props[JsonDeserializer.VALUE_DEFAULT_TYPE] = StuffDto::class.java.name
    props[ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG] = false // 关闭自动提交
    return DefaultKafkaConsumerFactory(props)
}
  • 修改DLT监听器方法,注入Acknowledgment对象,仅在处理成功时提交:
@KafkaListener(
    id = "stuff-dlt-listener",
    groupId = "${kafka.stuff.groupId}",
    topics = ["${kafka.stuff.topic}-dlt"],
    concurrency = "${kafka.stuff.dlt.concurrency}",
    autoStartup = "${kafka.stuff.dlt.auto-start:false}",
    containerFactory = "manualAckKafkaListenerContainerFactory"
)
fun processStuffInDlt(
    @Payload stuff: StuffDto,
    @Header(KafkaHeaders.RECEIVED_KEY) key: String,
    @Header(KafkaHeaders.RECEIVED_PARTITION) partition: Int,
    acknowledgment: Acknowledgment
) = runBlocking {
    log.info("{} Re-process stuff {} from DLT", LogSymbols.START_SYMBOL, key)
    try {
        process(stuff, key)
        acknowledgment.acknowledge() // 处理成功才提交偏移量
    } catch (e: Exception) {
        log.error("Failed to process DLT message {}", key, e)
        // 不提交,消息会留在DLT中,待下次监听器启动后重新消费
    }
}

3. 更优的DLT消息重处理方案

推荐两种更灵活的方案:

  • 方案一:DLT消息重放至原主题
    问题修复后,编写一次性任务将DLT消息重新发送到原消费主题,复用原监听器的处理逻辑,无需维护单独的DLT监听器:
@Autowired
private lateinit var kafkaTemplate: KafkaTemplate<String, StuffDto>
@Value("\${kafka.stuff.topic}")
private lateinit var originalTopic: String
@Value("\${kafka.stuff.topic}-dlt")
private lateinit var dltTopic: String

fun replayDltMessages() {
    val consumer = consumerFactory().createConsumer()
    consumer.subscribe(listOf(dltTopic))
    try {
        while (true) {
            val records = consumer.poll(Duration.ofSeconds(1))
            if (records.isEmpty) break
            for (record in records) {
                kafkaTemplate.send(originalTopic, record.key(), record.value())
                consumer.commitSync() // 重放成功后提交DLT偏移量
            }
        }
    } finally {
        consumer.close()
    }
}

可以通过REST接口或定时任务触发该方法,重放完成后自动结束,不会持续监听DLT。

  • 方案二:使用外部工具管理DLT重处理
    借助Kafka生态工具(如Kafka Connect的MirrorMaker2,或自定义脚本),将DLT消息批量重放至原主题,适合大规模消息重处理场景,避免在应用中编写重放逻辑。

内容的提问来源于stack exchange,提问作者João Matos

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 03:39:53