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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 01:23:12