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

Spring Cloud Stream Kafka绑定器:向DLQ消息添加自定义Header失败排查

问题原因

当你配置了enableDlq: true时,Spring Cloud Stream Kafka会自动创建一套完整的DLQ处理机制(包含DefaultErrorHandler和DeadLetterPublishingRecoverer),框架自动配置的优先级高于自定义Bean,导致你写的自定义错误处理代码不会被应用。

解决方案

提供两种可行方案,根据你的需求选择:


方案一:禁用自动DLQ,完全自定义错误处理流程

1. 修改配置文件

移除Stream的DLQ相关配置,让框架不再自动创建DLQ处理器:

spring:
  cloud:
    function:
      definition: numberConsumer
    stream:
      bindings:
        numberProducer-out-0:
          destination: first-topic
        numberConsumer-in-0:
          group: group
          destination: first-topic
      kafka:
        bindings:
          numberConsumer-in-0:
            consumer:
              standard-headers:

2. 修正自定义错误处理代码

完善DeadLetterPublishingRecoverer,指定DLQ主题并正确添加自定义Header,同时配置重试策略:

@Configuration
@Slf4j
public class KafkaConfiguration {

    // 获取根异常类型
    private Class<? extends Throwable> getRootCauseExceptionType(Throwable exception) {
        Throwable rootCause = exception;
        while (rootCause.getCause() != null) {
            rootCause = rootCause.getCause();
        }
        return rootCause.getClass();
    }

    @Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer<String, String>> customizer(DefaultErrorHandler errorHandler) {
        return (container, dest, group) -> container.setCommonErrorHandler(errorHandler);
    }

    @Bean
    public DefaultErrorHandler errorHandler(DeadLetterPublishingRecoverer deadLetterPublishingRecoverer) {
        // 配置0次重试,直接进入DLQ
        return new DefaultErrorHandler(deadLetterPublishingRecoverer, new FixedBackOff(0L, 0));
    }

    @Bean
    public DeadLetterPublishingRecoverer publisher(KafkaOperations<String, String> template) {
        // 指定DLQ主题和分区
        BiConsumer<ConsumerRecord<?, ?>, Exception> destinationResolver = (record, ex) -> 
                new TopicPartition("dlq", 0);

        DeadLetterPublishingRecoverer recover = new DeadLetterPublishingRecoverer(template, destinationResolver);

        recover.setExceptionHeadersCreator((kafkaHeaders, exception, isKey, headerNames) -> {
            // 添加异常类型Header
            String exceptionType = getRootCauseExceptionType(exception).getName();
            kafkaHeaders.add("exception-type", exceptionType.getBytes(StandardCharsets.UTF_8));
            // 添加异常栈信息Header
            String exceptionMsg = getExceptionStackTrace(exception);
            kafkaHeaders.add("exception", exceptionMsg.getBytes(StandardCharsets.UTF_8));
        });
        return recover;
    }

    // 手动获取异常栈信息(无需额外依赖)
    private String getExceptionStackTrace(Throwable exception) {
        StringWriter sw = new StringWriter();
        exception.printStackTrace(new PrintWriter(sw));
        return sw.toString();
    }
}

方案二:保留Stream自动DLQ配置,通过定制器添加自定义Header

这种方式无需改动原有DLQ配置,仅需添加定制器修改DLQ消息的Header:

1. 保留原有配置

继续使用你原来的application.yml配置(包含enableDlq: true等)。

2. 添加DLQ消息定制器

@Configuration
@Slf4j
public class KafkaConfiguration {

    @Bean
    public DlqMessageHandlerCustomizer dlqMessageHandlerCustomizer() {
        return (dlqHandler, destinationName, group) -> {
            dlqHandler.setHeaderMapper(new KafkaHeaderMapper() {
                @Override
                public void fromHeaders(Headers kafkaHeaders, org.springframework.messaging.MessageHeaders messageHeaders) {
                    // 先执行默认的Header映射逻辑
                    KafkaHeaderMapper.super.fromHeaders(kafkaHeaders, messageHeaders);
                    // 从消息头中获取异常对象
                    Throwable exception = (Throwable) messageHeaders.get(ErrorHeaders.ERROR_EXCEPTION);
                    if (exception != null) {
                        // 添加异常栈信息Header
                        StringWriter sw = new StringWriter();
                        exception.printStackTrace(new PrintWriter(sw));
                        kafkaHeaders.add(new RecordHeader("exception", sw.toString().getBytes(StandardCharsets.UTF_8)));
                        // 可选:添加异常类型Header
                        kafkaHeaders.add(new RecordHeader("exception-type", exception.getClass().getName().getBytes(StandardCharsets.UTF_8)));
                    }
                }

                @Override
                public void toHeaders(org.springframework.messaging.MessageHeaders messageHeaders, Headers kafkaHeaders) {
                    KafkaHeaderMapper.super.toHeaders(messageHeaders, kafkaHeaders);
                }
            });
        };
    }
}

验证方法

修改消费者代码主动抛出异常,触发DLQ:

@Bean
public Consumer<String> numberConsumer() {
    return message -> {
        log.info("receive message : {}", message);
        // 主动抛出异常测试DLQ
        throw new RuntimeException("测试消费异常");
    };
}

消费DLQ中的消息,检查Header是否包含exception(异常栈信息)和exception-type(异常类名)字段。

内容的提问来源于stack exchange,提问作者Yingjie Guan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 01:27:53