如何修改spring-kafka批监听器过滤逻辑实现整批仅1次数据库查询
解决方案
可以通过继承FilteringBatchMessageListenerAdapter重写onMessage方法的方式覆写默认过滤逻辑,实现整批消息仅调用一次数据库查询,具体实现如下:
1. 自定义批过滤监听器适配器
核心逻辑是先批量提取整批消息的业务过滤维度字段,仅调用一次数据库查询得到需要过滤的记录集合,再遍历执行过滤操作:
public class BatchQueryFilteringListenerAdapter<K, V> extends FilteringBatchMessageListenerAdapter<K, V> { // 注入自己的批量查询过滤服务 private final YourBatchFilterService filterService; public BatchQueryFilteringListenerAdapter(BatchMessageListener<K, V> delegate, RecordFilterStrategy<K, V> fallbackFilterStrategy, boolean ackDiscarded, YourBatchFilterService filterService) { super(delegate, fallbackFilterStrategy, ackDiscarded); this.filterService = filterService; } @Override public void onMessage(List<ConsumerRecord<K, V>> records, @Nullable Acknowledgment acknowledgment, Consumer<?, ?> consumer) { // 批量提取所有需要过滤的业务标识,比如订单ID、用户ID等 Set<Long> businessIdSet = records.stream() .map(record -> record.value().getYourBusinessId()) // 替换为实际的业务字段提取逻辑 .collect(Collectors.toSet()); // 仅调用一次数据库查询,拿到所有需要丢弃的记录标识集合 Set<Long> discardIdSet = filterService.batchQueryDiscardIds(businessIdSet); // 遍历执行过滤,无需再查询数据库 Iterator<ConsumerRecord<K, V>> iterator = records.iterator(); while (iterator.hasNext()) { ConsumerRecord<K, V> record = iterator.next(); // 若需要保留原有单条过滤策略,可在这里叠加fallbackFilterStrategy的判断 if (discardIdSet.contains(record.value().getYourBusinessId())) { iterator.remove(); } } // 后续逻辑与原有父类保持一致 if (!records.isEmpty()) { getDelegate().onMessage(records, acknowledgment, consumer); } else if (acknowledgment != null && isAckDiscarded()) { // 过滤后无剩余消息时,按配置决定是否提交ack acknowledgment.acknowledge(); } } }
2. 替换默认的批监听器适配器
在Kafka监听器工厂配置中,使用自定义的适配器替换默认实现即可:
@Bean public ConcurrentKafkaListenerContainerFactory<String, YourMessageDTO> batchKafkaListenerFactory( ConsumerFactory<String, YourMessageDTO> consumerFactory, YourBatchFilterService filterService) { ConcurrentKafkaListenerContainerFactory<String, YourMessageDTO> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 开启批监听 factory.setBatchListener(true); // 注入自定义的批过滤适配器 factory.setBatchMessageListenerAdapter(delegate -> new BatchQueryFilteringListenerAdapter( delegate, null, // 若不需要原有单条过滤策略可传null,否则传入自定义的RecordFilterStrategy实现 true, // 按业务需求配置是否自动ack过滤丢弃的消息 filterService )); return factory; }
注意事项
- 批量查询数据库时,注意控制单次查询的参数数量,避免超过数据库IN查询的参数上限,批次过大时可拆分分批查询
- 空消息列表的ack逻辑需要和业务的消费确认配置保持一致,避免出现消息丢失或重复消费问题
- 若原有业务中存在不需要走批量查询的单条过滤规则,可保留原有
RecordFilterStrategy的实现,在过滤逻辑中叠加判断即可
内容的提问来源于stack exchange,提问作者rloeffel
相关产品推荐
相关产品推荐

