使用非阻塞重试时仅为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
相关产品推荐
相关产品推荐

