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

如何实现@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 09:51:17