Spring-Kafka自定义DeadLetterRecoverer失效,如何在DLT前修改消息数据
解决自定义DeadLetterPublishingRecoverer不生效的问题
核心问题定位
你遇到的createProducerRecord未被调用的情况,大概率是因为Spring Boot自动配置逻辑中,DeadLetterRecovererFactory默认创建了原生的DeadLetterPublishingRecoverer,而非你的自定义子类。如果代码依赖工厂的自动创建逻辑,手动注入的自定义恢复器会被覆盖。
两种可行解决方案
方案一:手动配置DefaultErrorHandler,直接传入自定义恢复器
绕过DeadLetterRecovererFactory的自动创建,直接在DefaultErrorHandler中指定自定义恢复器,确保重试逻辑和死信处理都生效:
- 实现自定义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() ); } }
- 配置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或DeadLetterRecovererBean覆盖你的自定义实例。 - 验证恢复器实例化:在自定义恢复器的构造方法中添加日志或断点,确认Spring是否正确创建了实例。
内容的提问来源于stack exchange,提问作者Rostyslav Kholodnytskyi
相关产品推荐
相关产品推荐

