如何通过Spring Cloud Stream为DLT事件添加自定义Header
问题描述
我有一个消费者,希望在消费者处理事件失败时为DLT事件设置自定义Header。我尝试了以下实现方式,但未能生效,请问该如何正确实现?
@Component public class MyConsumer implements Consumer<Message<MyEvent>> { @Override public void accept(final Message<MyEvent> event) { // 处理事件的逻辑,会抛出异常 } @Bean public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(final KafkaOperations<?, ?> kafkaOperations) { final DeadLetterPublishingRecoverer deadLetterPublishingRecoverer = new DeadLetterPublishingRecoverer(kafkaOperations); deadLetterPublishingRecoverer.addHeadersFunction((consumerRecord, exception) -> { final Headers headers = consumerRecord.headers(); headers.add("foo", "bar".getBytes(StandardCharsets.UTF_8)); return headers; }); return deadLetterPublishingRecoverer; } @Bean public DefaultErrorHandler customDefaultErrorHandler(final DeadLetterPublishingRecoverer deadLetterPublishingRecoverer) { return new DefaultErrorHandler(deadLetterPublishingRecoverer); } @Bean public ListenerContainerCustomizer<AbstractMessageListenerContainer<byte[], byte[]>> customizer(final DefaultErrorHandler errorHandler) { return (container, destination, group) -> container.setCommonErrorHandler(errorHandler); } }
解决方案
直接修改consumerRecord.headers()无法生效的核心原因是:Kafka ConsumerRecord自带的Headers实例是不可修改的,修改操作无法被后续逻辑正确识别。正确的做法是创建全新的Headers对象,复制原消息的所有Header后再添加自定义Header。
修改后的DeadLetterPublishingRecoverer配置如下:
@Bean public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(final KafkaOperations<?, ?> kafkaOperations) { final DeadLetterPublishingRecoverer deadLetterPublishingRecoverer = new DeadLetterPublishingRecoverer(kafkaOperations); deadLetterPublishingRecoverer.addHeadersFunction((consumerRecord, exception) -> { // 创建新的Headers实例 Headers newHeaders = new RecordHeaders(); // 复制原消息的所有Header consumerRecord.headers().forEach(header -> newHeaders.add(header.key(), header.value())); // 添加自定义Header newHeaders.add("foo", "bar".getBytes(StandardCharsets.UTF_8)); return newHeaders; }); return deadLetterPublishingRecoverer; }
其他Bean配置无需改动,这样当消息进入DLT时,自定义Headerfoo: bar就会被正确添加到消息中。
内容的提问来源于stack exchange,提问作者Jonathan Henrique
相关产品推荐
相关产品推荐

