底层问题修复后如何重新处理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) }
但我预见到该方案存在至少两个问题:
- 消息重处理完成后,processStuffInDlt监听器会保持活跃,导致后续新的问题消息被重复处理——能否无需重新部署即可按需启停该监听器?
- 如果某条消息仍引发处理失败,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
相关产品推荐
相关产品推荐

