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

Spring Kafka基于自定义消息属性使用RecordFilterStrategy过滤消息方案

Spring Kafka基于自定义消息属性的过滤实现

报错的核心原因是你已经在消费者中配置了JSON反序列化规则,消息消费时已经自动将JSON字符串转为com.local.spring.kafka.Shipment类实例,因此record.value()返回的是自定义实体对象,不是原始JSON字符串,调用String类型的contains方法自然会抛出类型转换异常。

核心过滤逻辑

RecordFilterStrategy返回true代表丢弃当前消息,返回false代表保留消息进入业务监听逻辑,你只需要直接读取实体类的shipmentType属性判断即可。

完整配置示例

消费者工厂配置

import com.local.spring.kafka.Shipment;
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.listener.adapter.RecordFilterStrategy;

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Shipment> shipmentKafkaListenerFactory(
            ConsumerFactory<String, Shipment> consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<String, Shipment> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);

        // 配置自定义过滤规则
        factory.setRecordFilterStrategy((RecordFilterStrategy<String, Shipment>) record -> {
            Shipment shipment = record.value();
            // 空值防护,可根据业务调整逻辑
            if (shipment == null) {
                return true;
            }
            // shipmentType为YAW0时返回true,丢弃该消息
            return "YAW0".equals(shipment.getShipmentType());
        });

        // 可选配置:设置为true时,被过滤的消息会自动提交offset,避免重复消费
        factory.setAckDiscarded(true);

        return factory;
    }
}

业务监听代码

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class ShipmentMessageListener {

    // 指定使用我们配置好的带过滤规则的containerFactory
    @KafkaListener(topics = "你的业务Topic名称", containerFactory = "shipmentKafkaListenerFactory")
    public void handleShipmentMessage(Shipment shipment) {
        // 此处收到的消息已经完成过滤,不存在shipmentType为YAW0的记录
        // 编写你的业务处理逻辑即可
    }
}

注意事项

  • 请确保Shipment实体类中已经定义了shipmentType属性对应的getter方法,否则无法正常读取属性值
  • 如果你的项目使用手动提交offset模式,不需要额外调整过滤配置,过滤逻辑会在业务方法执行前触发,被丢弃的消息不会进入handleShipmentMessage方法

内容的提问来源于stack exchange,提问作者Sanchelz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 05:39:02