如何实现@KafkaListener按Event-Type消息头值过滤消费消息
结论
完全可以实现三个独立@KafkaListener按Event-Type消息头过滤消费的效果,你之前使用RecordFilterStrategy未生效,核心是踩了Kafka消费者组的机制坑。
核心原理说明
Kafka消费者组的默认语义是:同一条消息只会被同一个消费者组内的一个消费者实例消费。如果你给三个监听器配置了相同的group-id,即使给每个监听器绑定了不同的过滤规则,消息只会被随机分配给组内一个监听器,一旦被分配到的监听器过滤规则不匹配,消息会被直接丢弃,其他两个监听器永远无法收到这条消息,这就是你之前实现失败的根本原因。
要实现需求,只需要给三个监听器配置完全独立的消费者组ID,此时Kafka会将Topic内的每条消息全量投递给三个独立的消费者组,每个监听器只需要过滤掉不符合自己处理规则的消息即可,不会出现消息丢失。
基于RecordFilterStrategy的实现代码
1. 定义三个对应事件类型的过滤器
每个过滤器判断消息头的Event-Type值,不符合规则的消息直接过滤(RecordFilterStrategy返回true代表丢弃当前消息):
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.listener.RecordFilterStrategy; import java.nio.charset.StandardCharsets; // 类上记得加@Configuration注解 @Bean public RecordFilterStrategy<String, String> createEventFilter() { return record -> { byte[] headerValue = record.headers().lastHeader("Event-Type").value(); String eventType = new String(headerValue, StandardCharsets.UTF_8); return !"create".equals(eventType); }; } @Bean public RecordFilterStrategy<String, String> updateEventFilter() { return record -> { byte[] headerValue = record.headers().lastHeader("Event-Type").value(); String eventType = new String(headerValue, StandardCharsets.UTF_8); return !"update".equals(eventType); }; } @Bean public RecordFilterStrategy<String, String> deleteEventFilter() { return record -> { byte[] headerValue = record.headers().lastHeader("Event-Type").value(); String eventType = new String(headerValue, StandardCharsets.UTF_8); return !"delete".equals(eventType); }; }
2. 为每个过滤器绑定独立的监听器容器工厂
每个工厂配置和你原有逻辑一致的手动确认模式,绑定对应过滤器:
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.listener.ContainerProperties; @Bean public ConcurrentKafkaListenerContainerFactory<String, String> createListenerFactory( ConsumerFactory<String, String> consumerFactory, RecordFilterStrategy<String, String> createEventFilter) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setRecordFilterStrategy(createEventFilter); return factory; } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> updateListenerFactory( ConsumerFactory<String, String> consumerFactory, RecordFilterStrategy<String, String> updateEventFilter) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setRecordFilterStrategy(updateEventFilter); return factory; } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> deleteListenerFactory( ConsumerFactory<String, String> consumerFactory, RecordFilterStrategy<String, String> deleteEventFilter) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setRecordFilterStrategy(deleteEventFilter); return factory; }
3. 编写三个独立的监听器方法
每个监听器指定独立的group-id、绑定对应的容器工厂即可:
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.messaging.Message; @KafkaListener( topicPattern = "SameTopic", groupId = "biz-same-topic-create", containerFactory = "createListenerFactory" ) public void onCreate(Message<String> message, Acknowledgment acknowledgment) { doCreate(message); acknowledgment.acknowledge(); } @KafkaListener( topicPattern = "SameTopic", groupId = "biz-same-topic-update", containerFactory = "updateListenerFactory" ) public void onUpdate(Message<String> message, Acknowledgment acknowledgment) { doUpdate(message); acknowledgment.acknowledge(); } @KafkaListener( topicPattern = "SameTopic", groupId = "biz-same-topic-delete", containerFactory = "deleteListenerFactory" ) public void onDelete(Message<String> message, Acknowledgment acknowledgment) { doDelete(message); acknowledgment.acknowledge(); }
其他说明
- 被过滤器丢弃的消息,Spring Kafka容器会自动提交对应位移,不需要你手动执行ack操作,不会出现消息堆积或者重复消费的问题。
- 如果觉得定义三个容器工厂过于冗余,也可以不使用
RecordFilterStrategy,直接在监听器方法内部做类型判断,不符合规则的消息直接ack后返回即可,最终效果完全一致,只是代码耦合度稍高。
内容的提问来源于stack exchange,提问作者hideburn
相关产品推荐
相关产品推荐

