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

使用非阻塞重试时仅为DLT Topic添加自定义Header的实现方案

Spring Kafka 2.8.8:重试/DLT Topic自定义Header重复问题的解决方案

问题描述

我通过以下代码为发送至重试Topic和DLT Topic的事件添加自定义Header:

@Bean(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME)
public DeadLetterPublishingRecovererFactory factory(DestinationTopicResolver destinationTopicResolver) {
    DeadLetterPublishingRecovererFactory factory = new DeadLetterPublishingRecovererFactory(destinationTopicResolver);
    factory.setDeadLetterPublishingRecovererCustomizer(dlpr -> dlpr.addHeadersFunction((consumerRecord, exception) -> {
        return new RecordHeaders().add(EXTERNAL_SYSTEM_EXCEPTION_MESSAGE_HEADER, exception.getMessage().getBytes(StandardCharsets.UTF_8));
    }));
    return factory;
}

但出现了Header重复的问题:配置2个重试Topic和1个DLT Topic后,抛出异常后主Topic添加该Header,重试-0 Topic接收时已存在该Header,再次添加后重试-1 Topic出现2个同名Header,最终DLT Topic中存在3个同名Header。这是因为Apache Kafka支持多个同名Header。

我想到两种解决方案,但不确定如何实现:

  • 每次添加Header前先检查是否存在,若存在则删除后再添加;
  • 仅为DLT Topic添加该Header。

请问这两种方案是否可行?我使用的是spring-kafka 2.8.8版本。


方案可行性及实现方法

两种方案均完全可行,针对spring-kafka 2.8.8版本的具体实现如下:

方案1:检查并覆盖同名Header

实现思路是复制原记录的Header集合,移除已存在的同名Header后再添加新值,避免重复创建:

@Bean(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME)
public DeadLetterPublishingRecovererFactory factory(DestinationTopicResolver destinationTopicResolver) {
    DeadLetterPublishingRecovererFactory factory = new DeadLetterPublishingRecovererFactory(destinationTopicResolver);
    factory.setDeadLetterPublishingRecovererCustomizer(dlpr -> dlpr.addHeadersFunction((consumerRecord, exception) -> {
        // 复制原记录的所有Header到新集合
        RecordHeaders newHeaders = new RecordHeaders(consumerRecord.headers());
        // 移除已存在的同名Header
        newHeaders.remove(EXTERNAL_SYSTEM_EXCEPTION_MESSAGE_HEADER);
        // 添加最新的Header值
        newHeaders.add(EXTERNAL_SYSTEM_EXCEPTION_MESSAGE_HEADER, exception.getMessage().getBytes(StandardCharsets.UTF_8));
        return newHeaders;
    }));
    return factory;
}

注:RecordHeaders的构造方法传入原headers会创建副本,不会修改原记录的Header,确保流程无副作用。

方案2:仅为DLT Topic添加Header

通过DestinationTopicResolver解析目标Topic的类型,判断是否为DLT后再添加Header:

@Bean(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME)
public DeadLetterPublishingRecovererFactory factory(DestinationTopicResolver destinationTopicResolver) {
    DeadLetterPublishingRecovererFactory factory = new DeadLetterPublishingRecovererFactory(destinationTopicResolver);
    factory.setDeadLetterPublishingRecovererCustomizer(dlpr -> dlpr.addHeadersFunction((consumerRecord, exception) -> {
        // 获取当前记录的目标Topic元数据
        DestinationTopic destinationTopic = destinationTopicResolver.resolveDestinationTopic(consumerRecord);
        // 仅为DLT Topic添加自定义Header
        if (destinationTopic.isDlt()) {
            return new RecordHeaders().add(EXTERNAL_SYSTEM_EXCEPTION_MESSAGE_HEADER, exception.getMessage().getBytes(StandardCharsets.UTF_8));
        }
        // 非DLT Topic返回空Header集合
        return new RecordHeaders();
    }));
    return factory;
}

如果已知DLT Topic的具体名称,也可以直接通过名称对比判断,逻辑更直观:

@Bean(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME)
public DeadLetterPublishingRecovererFactory factory(DestinationTopicResolver destinationTopicResolver) {
    DeadLetterPublishingRecovererFactory factory = new DeadLetterPublishingRecovererFactory(destinationTopicResolver);
    String dltTopicName = "your-dlt-topic-name"; // 替换为实际DLT Topic名称
    
    factory.setDeadLetterPublishingRecovererCustomizer(dlpr -> dlpr.addHeadersFunction((consumerRecord, exception) -> {
        DestinationTopic destinationTopic = destinationTopicResolver.resolveDestinationTopic(consumerRecord);
        if (dltTopicName.equals(destinationTopic.getDestinationName())) {
            return new RecordHeaders().add(EXTERNAL_SYSTEM_EXCEPTION_MESSAGE_HEADER, exception.getMessage().getBytes(StandardCharsets.UTF_8));
        }
        return new RecordHeaders();
    }));
    return factory;
}

注:spring-kafka 2.8.8版本中DestinationTopic.isDlt()方法已稳定支持,可直接使用。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:21:03