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
相关产品推荐
相关产品推荐

