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

如何修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 23:12:32