SpringBoot中如何在监听器级别基于Kafka Header值过滤消息?
问题:基于Kafka Header过滤消息时应用启动失败
我希望在监听器级别根据Kafka Header的值过滤消息,编写了自定义的KafkaConsumerConfig.java配置类,代码如下:
package com.example.test.consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.header.Header; import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaConsumerConfig { @Value("${spring.kafka.consumer.group-id}") private String group_id; @Value("${spring.kafka.bootstrap-servers}") private String bootstrap_server; @Bean ConsumerFactory<String, String> consumerFactory() { Map<String, Object> config = new HashMap<>(); config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap_server); config.put(ConsumerConfig.GROUP_ID_CONFIG, group_id); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return new DefaultKafkaConsumerFactory<>(config); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setRecordFilterStrategy(consumerRecord -> { if (consumerRecord.headers().lastHeader("CUSTOM_HEADER").key().equals("VALID")) { return false; } return true; }); return factory; } }
但应用启动失败,报错信息如下:
Cancelled in-flight API_VERSIONS request with correlation id 20 due to node -1 being disconnected
问题分析与解决
1. 启动连接错误排查
报错node -1 being disconnected本质是Kafka客户端无法连接指定的bootstrap服务器,常见原因:
- 配置错误:
spring.kafka.bootstrap-servers的地址/端口写错,或Kafka服务未启动 - 网络限制:应用所在机器无法访问Kafka服务器(防火墙、网络策略拦截)
- 注入异常:检查配置文件中
spring.kafka.bootstrap-servers的键名是否正确,确保@Value能正常读取值
可先用Kafka命令行工具验证连接:
kafka-topics.sh --list --bootstrap-server <你的bootstrap地址>
2. 过滤逻辑的两处关键错误
即使连接问题解决,原代码的过滤逻辑存在致命问题:
- 空指针风险:如果消息没有
CUSTOM_HEADER,lastHeader("CUSTOM_HEADER")会返回null,调用key()会直接抛出空指针 - 逻辑错误:
Header.key()返回的是Header的名称(即CUSTOM_HEADER),而非Header存储的实际值,应该用Header.value()获取值后转成字符串比较
3. 修正后的完整配置类
package com.example.test.consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.header.Header; import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaConsumerConfig { @Value("${spring.kafka.consumer.group-id}") private String groupId; @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; @Bean ConsumerFactory<String, String> consumerFactory() { Map<String, Object> config = new HashMap<>(); config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return new DefaultKafkaConsumerFactory<>(config); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setRecordFilterStrategy(consumerRecord -> { Header customHeader = consumerRecord.headers().lastHeader("CUSTOM_HEADER"); if (customHeader != null) { String headerValue = new String(customHeader.value()); // 返回true=过滤消息,返回false=保留消息 return !"VALID".equals(headerValue); } // 无目标Header时过滤消息 return true; }); return factory; } }
内容的提问来源于stack exchange,提问作者Sayan Saha
相关产品推荐
相关产品推荐

