如何在DefaultErrorHandler DLT中用SLF4J格式打印可重试错误日志
解决方案
要实现用SLF4J标准格式输出可重试错误日志,替代冗长的堆栈打印,核心是通过DefaultErrorHandler自定义错误日志逻辑,同时保留原有重试和死信队列(DLT)的功能。以下是具体实现步骤:
1. 配置自定义错误日志处理器
修改KafkaListenerConfig中的错误处理器配置,为DefaultErrorHandler设置自定义ErrorLogger,用SLF4J输出包含关键信息的简洁日志,避免打印完整堆栈:
package org.abc; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.boot.autoconfigure.kafka.ConcurrentKafkaListenerContainerFactoryConfigurer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.kafka.listener.ErrorLogger; import org.springframework.util.backoff.FixedBackOff; import com.azure.cosmos.CosmosException; import com.azure.spring.data.cosmos.exception.CosmosAccessException; import com.fasterxml.jackson.core.JsonProcessingException; @Configuration public class KafkaListenerConfig { private static final Logger LOGGER = LoggerFactory.getLogger(KafkaListenerConfig.class); @Bean(name = "kafkaListenerContainerFactory") public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactory( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory<Object, Object> kafkaConsumerFactory, KafkaTemplate<Object, Object> template) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); var recoverer = new DeadLetterPublishingRecoverer(template, (record, ex) -> new TopicPartition("abc_dlt", record.partition())); var errorHandler = getDLTDetails(recoverer); factory.setCommonErrorHandler(errorHandler); return factory; } @Bean(name = "kafkaListenerContainerFactoryForDLT") public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactoryForDLT( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory<Object, Object> kafkaConsumerFactory, KafkaTemplate<Object, Object> template) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); var errorHandler = new DefaultErrorHandler(new FixedBackOff(3, 2)); errorHandler.setErrorLogger(new CustomErrorLogger()); factory.setCommonErrorHandler(errorHandler); factory.getContainerProperties().setIdleEventInterval(70000L); return factory; } private DefaultErrorHandler getDLTDetails(DeadLetterPublishingRecoverer recoverer) { var errorHandler = new DefaultErrorHandler(recoverer, new FixedBackOff(3, 2)); errorHandler.addNotRetryableExceptions(JsonProcessingException.class); errorHandler.addRetryableExceptions(CosmosException.class, CosmosAccessException.class); errorHandler.setCommitRecovered(true); errorHandler.setErrorLogger(new CustomErrorLogger()); return errorHandler; } /** * 自定义错误日志器,输出SLF4J标准格式的错误信息 */ private static class CustomErrorLogger implements ErrorLogger { @Override public void log(Exception thrownException, ConsumerRecord<?, ?> record, String description) { if (record != null) { LOGGER.error("消费消息失败 | Topic: {} | Partition: {} | Offset: {} | 异常类型: {} | 异常信息: {}", record.topic(), record.partition(), record.offset(), thrownException.getClass().getSimpleName(), thrownException.getMessage()); } else { LOGGER.error("消费处理失败 | 描述: {} | 异常类型: {} | 异常信息: {}", description, thrownException.getClass().getSimpleName(), thrownException.getMessage()); } } } }
2. 替换自定义日志工具为SLF4J标准实现
修改Listener类,使用SLF4J标准Logger替代自定义LogUtil,遵循统一日志规范:
package org.abc; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; @Component public class Listener { private static final Logger LOGGER = LoggerFactory.getLogger(Listener.class); @KafkaListener(id = "hub_topic", topics = "hub_topic", groupId = "hub_topic_0", containerFactory = "kafkaListenerContainerFactory", clientIdPrefix = "hub_topic") public void listen(ConsumerRecord<String, String> consumerRecord, Acknowledgment ack) { LOGGER.info("收到消息: {}", consumerRecord.value()); try { saveToDB(consumerRecord.value()); ack.acknowledge(); } catch (Exception e) { // 异常交由错误处理器统一记录日志,无需手动打印堆栈 throw e; } } private void saveToDB(String payload) { // 数据库持久化逻辑实现 } }
关键说明
CustomErrorLogger实现ErrorLogger接口,自定义日志输出格式,仅记录异常类型、消息和消费记录关键元数据,避免冗余堆栈信息。- 保留
DefaultErrorHandler的重试和死信转发逻辑,日志处理与业务逻辑解耦。 - 统一使用SLF4J标准API,便于后续对接Logback、Log4j2等不同日志框架。
内容的提问来源于stack exchange,提问作者KGT
相关产品推荐
相关产品推荐

