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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 11:00:55