如何修改spring-kafka中DeadLetterPublishingRecovererFactory的日志级别
解决Spring Kafka重试转DLT时的ERROR日志级别调整问题
你遇到的这条ERROR日志,来自kafka.retrytopic.DeadLetterPublishingRecovererFactory类,而非你之前调整的DefaultErrorHandler,所以修改后者的日志级别不会生效。以下两种方案可以解决这个问题:
方案一:直接修改日志框架配置(推荐快速解决)
通过日志配置文件(如Logback、Log4j2),将DeadLetterPublishingRecovererFactory的日志级别调低至INFO或WARN即可。
Logback配置示例(logback.xml)
<logger name="kafka.retrytopic.DeadLetterPublishingRecovererFactory" level="INFO" additivity="false"> <appender-ref ref="你的日志输出Appender"/> </logger>
Log4j2配置示例(log4j2.xml)
<Logger name="kafka.retrytopic.DeadLetterPublishingRecovererFactory" level="INFO" additivity="false"> <AppenderRef ref="你的日志输出Appender"/> </Logger>
方案二:自定义DeadLetterPublishingRecoverer(灵活定制)
如果需要更灵活的日志控制(比如自定义日志内容),可以自定义DeadLetterPublishingRecoverer并替换默认实现:
1. 创建自定义DLT Recoverer
@Bean fun customDltRecoverer(kafkaTemplate: KafkaTemplate<String, *>) = object : DeadLetterPublishingRecoverer(kafkaTemplate) { private val logger = LoggerFactory.getLogger(javaClass) override fun publishDeadLetter(record: ConsumerRecord<*, *>, exception: Exception, topicSuffixingStrategy: TopicSuffixingStrategy) { // 用INFO级别输出自定义日志内容 logger.info("消息转DLT处理:topic=${record.topic()}, partition=${record.partition()}, offset=${record.offset()},不再重试,发送至DLT主题") // 调用父类方法完成DLT发送逻辑 super.publishDeadLetter(record, exception, topicSuffixingStrategy) } }
2. 在重试配置中指定自定义Recoverer
修改你的webhookRetryTopic配置,添加.deadLetterPublishingRecoverer()指定自定义实现:
@Bean fun webhookRetryTopic( kafkaTemplate: KafkaTemplate<String, KafkaMessageWrapper>, customDltRecoverer: DeadLetterPublishingRecoverer ) = RetryTopicConfigurationBuilder .newInstance() .autoCreateTopics(true, partition, replicationFactor.toShort()) .maxAttempts(webhookBackOffProperties.attempts) .exponentialBackoff( webhookBackOffProperties.initialInterval * 1000L, webhookBackOffProperties.multiplier.toDouble(), webhookBackOffProperties.maxInterval * 1000L ) .sameIntervalTopicReuseStrategy(SameIntervalTopicReuseStrategy.SINGLE_TOPIC) .includeTopics(KafkaConfiguration.webhookTopics) .retryOn(KafkaRetryableException::class.java) .dltProcessingFailureStrategy(DltStrategy.FAIL_ON_ERROR) .dltHandlerMethod(EndpointHandlerMethod(KafkaWebhookDltHandler::class.java, "dlt")) // 加入自定义DLT Recoverer .deadLetterPublishingRecoverer(customDltRecoverer) .create(kafkaTemplate)
方案对比
- 方案一:配置简单,无需修改代码,适合仅调整日志级别的场景。
- 方案二:可以完全自定义日志内容、级别,甚至扩展DLT发送逻辑,适合需要定制化的场景。
内容的提问来源于stack exchange,提问作者Minuitdix
相关产品推荐
相关产品推荐

