Spring Boot:FilteringMessageListenerAdapter正确使用及过滤失效问题
嘿,我之前也踩过这个坑!你说的RecordFilterStrategy没被调用,大概率是因为没正确把FilteringMessageListenerAdapter和消费者容器绑定,或者用@KafkaListener时没配置过滤策略。下面给你拆解正确的用法:
正确使用FilteringMessageListenerAdapter的两种方式
1. 手动配置容器(适合自定义需求较多的场景)
这种方式需要你自己把实际的消息监听器包装进FilteringMessageListenerAdapter,再绑定到Kafka消费者容器:
第一步:实现自定义RecordFilterStrategy
先写好你的过滤逻辑,返回true表示丢弃这条消息,false表示保留:
@Component public class CustomKafkaFilter implements RecordFilterStrategy<String, String> { @Override public boolean filter(ConsumerRecord<String, String> record) { // 示例:过滤掉包含"test"的消息 return record.value().contains("test"); } }
第二步:配置FilteringMessageListenerAdapter和容器
把你的实际业务监听器包装进适配器,再设置到容器中:
@Configuration public class KafkaConsumerConfig { @Autowired private CustomKafkaFilter filterStrategy; @Autowired private KafkaProperties kafkaProperties; // 你的实际业务消息监听器 @Bean public MessageListener<String, String> businessMessageListener() { return record -> { // 这里只处理通过过滤的消息 System.out.println("收到有效消息: " + record.value()); }; } // 包装成带过滤的适配器 @Bean public FilteringMessageListenerAdapter<String, String> filteringListenerAdapter() { return new FilteringMessageListenerAdapter<>(businessMessageListener(), filterStrategy); } // 配置消费者容器并绑定适配器 @Bean public ConcurrentMessageListenerContainer<String, String> kafkaListenerContainer() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(defaultConsumerFactory()); ConcurrentMessageListenerContainer<String, String> container = factory.createContainer("your-topic"); // 关键:把适配器设置为容器的监听器 container.setupMessageListener(filteringListenerAdapter()); return container; } private ConsumerFactory<String, String> defaultConsumerFactory() { return new DefaultKafkaConsumerFactory<>(kafkaProperties.buildConsumerProperties()); } }
2. 结合@KafkaListener注解(更常用的简化方式)
如果习惯用@KafkaListener注解,不需要手动创建适配器,只要在容器工厂中配置过滤策略即可,框架会自动用FilteringMessageListenerAdapter包装你的监听器:
第一步:配置容器工厂
@Configuration public class KafkaListenerConfig { @Autowired private CustomKafkaFilter filterStrategy; @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(ConsumerFactory<String, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 关键:给容器工厂设置过滤策略 factory.setRecordFilterStrategy(filterStrategy); return factory; } }
第二步:编写@KafkaListener监听器
这时候你的监听器就只会收到通过过滤的消息了:
@Component public class KafkaBusinessListener { @KafkaListener(topics = "your-topic") public void handleValidMessage(ConsumerRecord<String, String> record) { System.out.println("处理业务消息: " + record.value()); } }
常见坑点提醒
- 别搞反过滤逻辑:
filter()返回true是丢弃消息,false是保留消息 - 用
@KafkaListener时,必须在容器工厂中设置recordFilterStrategy,否则框架不会启用过滤 - 不要直接把
FilteringMessageListenerAdapter当成普通MessageListener使用,必须让容器/容器工厂感知到它的存在
内容的提问来源于stack exchange,提问作者user1052610
相关产品推荐
相关产品推荐

