自定义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
相关产品推荐
相关产品推荐

