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
相关产品推荐
相关产品推荐

