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

自定义DeadLetterPublishingRecoverer处理无效消息时未生效的问题求助

解决自定义DeadLetterPublishingRecoverer在RetryTopic场景下不生效的问题

问题根源在于Spring Kafka的RetryTopic功能会自动装饰你指定的监听器工厂,替换其中的错误处理器为内部创建的默认DeadLetterPublishingRecoverer,忽略你原本在listenerFactory中配置的自定义恢复器。

解决方案:自定义DeadLetterPublishingRecovererFactory

RetryTopic内部通过DeadLetterPublishingRecovererFactory创建恢复器,因此我们可以提供自定义的工厂Bean,让RetryTopic使用我们的DeadLetterPublishingRecoverer实现。

修改后的完整配置

@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
    var listenerFactory = new ConcurrentKafkaListenerContainerFactory<String, Object>();
    listenerFactory.setConsumerFactory(consumerFactory());
    listenerFactory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD);
    return listenerFactory;
}

@Bean
public ConsumerFactory<String, Object> consumerFactory() {
    return new DefaultKafkaConsumerFactory<>(
            kafkaConsumerProps,
            StringDeserializer::new,
            () -> new ErrorHandlingDeserializer<>(new KafkaAvroDeserializer()),
            true
    );
}

@Bean
public ProducerFactory<String, Object> producerFactory() {
    return new DefaultKafkaProducerFactory<>(
            kafkaProducerProps,
            StringSerializer::new,
            KafkaAvroSerializer::new,
            true
    );
}

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

// 关键:自定义DeadLetterPublishingRecovererFactory,让RetryTopic使用我们的恢复器
@Bean
public DeadLetterPublishingRecovererFactory deadLetterPublishingRecovererFactory(KafkaTemplate<String, Object> kafkaTemplate) {
    return new DeadLetterPublishingRecovererFactory() {
        @Override
        public DeadLetterPublishingRecoverer create() {
            DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate, (r, e) -> {
                log.error("Custom error resolving {}", e);
                // 可根据异常类型自定义DLT目标,例如返回new TopicPartition("custom-dlt", 0)
                return null;
            });
            recoverer.setFailIfSendResultIsError(false);
            return recoverer;
        }
    };
}

@Bean
public RetryTopicConfiguration retryTopicConfiguration(KafkaTemplate<String, Object> template,
                                                       ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory) {
    return RetryTopicConfigurationBuilder
            .newInstance()
            .fixedBackOff(2000)
            .maxAttempts(3)
            .dltHandlerMethod("deadLetterHandler", "handle")
            .dltSuffix("-test-dlt")
            .retryTopicSuffix("-test-retry")
            .doNotRetryOnDltFailure()
            .listenerFactory(kafkaListenerContainerFactory)
            .create(template);
}

原理说明

当RetryTopic初始化时,会自动查找容器中的DeadLetterPublishingRecovererFactory Bean。如果存在自定义实现,就会用它创建恢复器,替代默认的工厂逻辑,这样就能保证你的自定义DeadLetterPublishingRecoverer被正确使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 14:37:20