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

Spring Kafka DLT消息发布失败问题排查

问题原因及解决方法

你的代码里的DefaultErrorHandler仅配置了一个自定义恢复回调函数,这个回调只会打印日志,并没有实现将消息发送到DLT主题的逻辑。要启用DLT发布功能,必须使用DeadLetterPublishingRecoverer作为恢复策略,它依赖KafkaTemplate完成消息转发。

修改步骤:

1. 调整DefaultErrorHandler的Bean定义

注入KafkaTemplate,创建DeadLetterPublishingRecoverer实例并作为恢复回调传入DefaultErrorHandler:

@Component
public class MyConsumer {

    private final KafkaTemplate<String, String> kafkaTemplate;

    // 构造注入KafkaTemplate
    public MyConsumer(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    @KafkaListener(topics = "my-topic")
    public void receive(ConsumerRecord<String, String> record) {
        System.out.println(record);
        throw new RuntimeException();
    }

    @Bean
    public DefaultErrorHandler defaultErrorHandler() {
        // 创建DeadLetterPublishingRecoverer,负责将消息转发到DLT
        DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate);
        // 配置重试策略:0次重试,直接进入DLT
        return new DefaultErrorHandler(recoverer, new FixedBackOff(0L, 0L));
    }
}

2. 确保application.properties的关键配置

需要配置消费者组ID,可选开启主题自动创建(方便测试):

# 消费者配置
spring.kafka.consumer.group-id=my-consumer-group
spring.kafka.consumer.auto-offset-reset=earliest

# 生产者配置(DLT转发需要用到)
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer

# 可选:开启自动创建主题(测试环境用)
spring.kafka.admin.auto-create=true

3. 验证DLT主题

默认情况下,DeadLetterPublishingRecoverer会将失败消息发送到原主题名.DLT(即my-topic.DLT)。如果Kafka集群未开启自动创建主题功能,需要手动创建该主题。

修改后,当消费抛出异常时,消息会被转发到my-topic.DLT,不会再循环执行消费逻辑。

内容的提问来源于stack exchange,提问作者Vin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 22:33:14