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

Spring-Kafka自定义DeadLetterRecoverer失效,如何在DLT前修改消息数据

解决自定义DeadLetterPublishingRecoverer不生效的问题

核心问题定位

你遇到的createProducerRecord未被调用的情况,大概率是因为Spring Boot自动配置逻辑中,DeadLetterRecovererFactory默认创建了原生的DeadLetterPublishingRecoverer,而非你的自定义子类。如果代码依赖工厂的自动创建逻辑,手动注入的自定义恢复器会被覆盖。

两种可行解决方案

方案一:手动配置DefaultErrorHandler,直接传入自定义恢复器

绕过DeadLetterRecovererFactory的自动创建,直接在DefaultErrorHandler中指定自定义恢复器,确保重试逻辑和死信处理都生效:

  1. 实现自定义DeadLetterPublishingRecoverer子类
public class MyDeadLetterRecoverer extends DeadLetterPublishingRecoverer {
    public MyDeadLetterRecoverer(KafkaTemplate<?, ?> kafkaTemplate) {
        super(kafkaTemplate);
    }

    @Override
    protected ProducerRecord<?, ?> createProducerRecord(ConsumerRecord<?, ?> original, TopicPartition topicPartition, Exception exception) {
        // 在这里修改死信消息,比如添加自定义头、修改消息体
        ProducerRecord<?, ?> baseRecord = super.createProducerRecord(original, topicPartition, exception);
        // 示例:为消息体添加前缀,同时保留原头部信息
        return new ProducerRecord<>(
                baseRecord.topic(),
                baseRecord.partition(),
                baseRecord.timestamp(),
                baseRecord.key(),
                "dlt-modified:" + baseRecord.value(),
                baseRecord.headers()
        );
    }
}
  1. 配置DefaultErrorHandler Bean,注入自定义恢复器
@Configuration
public class KafkaErrorConfig {

    @Bean
    public MyDeadLetterRecoverer myDeadLetterRecoverer(KafkaTemplate<Object, Object> kafkaTemplate) {
        return new MyDeadLetterRecoverer(kafkaTemplate);
    }

    @Bean
    public DefaultErrorHandler defaultErrorHandler(MyDeadLetterRecoverer myDeadLetterRecoverer) {
        // 配置重试策略:3次重试,间隔1秒
        FixedBackOff backOff = new FixedBackOff(1000L, 3);
        DefaultErrorHandler errorHandler = new DefaultErrorHandler(myDeadLetterRecoverer, backOff);
        
        // 可选:指定哪些异常需要重试/直接进入死信
        errorHandler.addNotRetryableExceptions(IllegalArgumentException.class);
        
        return errorHandler;
    }
}

方案二:自定义DeadLetterRecovererFactory,返回子类实例

如果场景必须依赖DeadLetterRecovererFactory(比如动态创建恢复器的场景),可以自定义工厂实现,让它返回你的自定义恢复器:

@Configuration
public class KafkaErrorConfig {

    @Bean
    public DeadLetterRecovererFactory deadLetterRecovererFactory(KafkaTemplate<Object, Object> kafkaTemplate) {
        return new DeadLetterRecovererFactory() {
            @Override
            public DeadLetterRecoverer createDeadLetterRecoverer(ConsumerFactory<?, ?> consumerFactory) {
                return new MyDeadLetterRecoverer(kafkaTemplate);
            }
        };
    }

    @Bean
    public DefaultErrorHandler defaultErrorHandler(DeadLetterRecovererFactory recovererFactory) {
        FixedBackOff backOff = new FixedBackOff(1000L, 3);
        return new DefaultErrorHandler(recovererFactory.createDeadLetterRecoverer(null), backOff);
    }
}

常见排查点

  • 确认重试次数已耗尽:只有当重试次数用完后,才会触发死信恢复器,可临时将重试次数设为0测试是否触发createProducerRecord。
  • 检查Bean优先级:确保没有其他DefaultErrorHandler或DeadLetterRecoverer Bean覆盖你的自定义实例。
  • 验证恢复器实例化:在自定义恢复器的构造方法中添加日志或断点,确认Spring是否正确创建了实例。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 03:07:19