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

Spring Boot @KafkaListener批量处理消息时无法获取自定义Header

问题

我有一个接收批量消息的Kafka监听器,需要从监听器中获取自定义Header的列表,但系统提示未找到该Header。

原代码如下:

@KafkaListener(id = KAFKA_LISTENER_ID,
    topics = "${kafka.event-topics.someTopic}",
    properties = {"spring.json.value.default.type=com.id.somegateway.domain.dto.consumerevent.SomeEventConsumedV1"})
public void consumeMessages(@Payload
                List<SomeEventConsumedV1> messages,
                @Header(KafkaHeaders.RECEIVED_TIMESTAMP) List<Long> timestamps,
                @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
                @Header(KafkaHeaders.OFFSET) List<Long> offsets,
                @Header("CUSTOM_HEADER") List<Long> sendTimeHeaders,
                Acknowledgment acknowledgment) {
    // 业务逻辑
}

上述代码中无法读取@Header("CUSTOM_HEADER") List<Long> sendTimeHeaders。

目前已有一种解决方案,即将方法参数改为接收List<Message<SomeEventConsumedV1>>:

@KafkaListener(id = KAFKA_LISTENER_ID,
    topics = "${kafka.event-topics.someTopic}",
    properties = {"spring.json.value.default.type=com.id.somegateway.domain.dto.consumerevent.SomeEventConsumedV1"})
public void consumeMessages(List<Message<SomeEventConsumedV1>> messages) {
    messages.stream().map(msg -> msg.getHeaders().get("CUSTOM_HEADER"));
    // 后续处理
}

但想了解是否存在使用@Header注解的解决方案?

解决方案(使用@Header注解实现)
  • 要在批量消费场景下通过@Header获取自定义Header列表,核心前提是生产者发送每条消息时都正确携带了CUSTOM_HEADER,在此基础上可以通过两种方式实现:

方式一:通过@Headers批量获取后提取

直接注入MessageHeaders对象,从中提取自定义Header的列表:

@KafkaListener(id = KAFKA_LISTENER_ID,
    topics = "${kafka.event-topics.someTopic}",
    properties = {"spring.json.value.default.type=com.id.somegateway.domain.dto.consumerevent.SomeEventConsumedV1"})
public void consumeMessages(@Payload List<SomeEventConsumedV1> messages,
                            @Headers MessageHeaders headers,
                            Acknowledgment acknowledgment) {
    // 提取批量消息的CUSTOM_HEADER列表
    List<Long> sendTimeHeaders = headers.get("CUSTOM_HEADER", List.class);
    if (sendTimeHeaders != null) {
        // 处理自定义Header逻辑
    }
    // 其他业务逻辑
}

方式二:直接使用@Header注解(指定required=false)

修改原代码的@Header注解,添加required = false避免找不到Header时报错,同时确保生产者端的Header格式匹配:

@KafkaListener(id = KAFKA_LISTENER_ID,
    topics = "${kafka.event-topics.someTopic}",
    properties = {"spring.json.value.default.type=com.id.somegateway.domain.dto.consumerevent.SomeEventConsumedV1"})
public void consumeMessages(@Payload List<SomeEventConsumedV1> messages,
                            @Header(KafkaHeaders.RECEIVED_TIMESTAMP) List<Long> timestamps,
                            @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
                            @Header(KafkaHeaders.OFFSET) List<Long> offsets,
                            @Header(value = "CUSTOM_HEADER", required = false) List<Long> sendTimeHeaders,
                            Acknowledgment acknowledgment) {
    if (sendTimeHeaders != null) {
        // 处理自定义Header列表
    }
    // 其他业务逻辑
}

关键注意事项

  • Spring Kafka在批量消费时,会自动将每条消息的同名Header收集为List,但如果生产者发送消息时未携带该Header,或者Header类型与声明的List<Long>不匹配(比如是字符串),则需要手动转换:
    // 若Header是字符串类型,转换为Long列表
    List<String> headerStrList = headers.get("CUSTOM_HEADER", List.class);
    List<Long> sendTimeHeaders = headerStrList.stream().map(Long::valueOf).collect(Collectors.toList());
    
  • 若始终无法获取到Header,需排查生产者端是否正确设置了Header,比如使用ProducerRecord时通过headers()方法添加:
    ProducerRecord<String, SomeEvent> record = new ProducerRecord<>(topic, key, value);
    record.headers().add("CUSTOM_HEADER", String.valueOf(System.currentTimeMillis()).getBytes());
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 03:53:13