Spring 6中JdbcPollingChannelAdapter实现按需轮询下批数据的方法
实现轮询器等待前一批数据处理完成后再轮询下一批
现有代码
@Bean @InboundChannelAdapter(value = "databaseInputChannel") public JdbcPollingChannelAdapter checkDbForInput() { JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(dataSource(), "select * from users"); adapter.setUpdateSql("update users set load_status = '3' where id in (:id)"); adapter.setMaxRows(100); return adapter; } @Bean public IntegrationFlow pollingFlow() { return IntegrationFlow.from("databaseInputChannel") .split() .channel("databaseRowInputChannel") .get(); } @Async @ServiceActivator(inputChannel = "databaseRowInputChannel") public void processMessage(Message<?> message) { }
需求
希望轮询器在前100行数据处理完成(受线程池限制)后,再轮询下100行数据。此前尝试的jdbcPollingChannelAdapter.setTrigger或jdbcPollingChannelAdapter.setTaskDecorator在Spring 6中已被移除,需要更简便的实现方式。
解决方案
利用Spring Integration的MessageTriggerAdvice配合聚合器,实现批次处理完成后触发下一次轮询:
- 配置带触发通知的轮询器,让轮询器等待批次完成信号
- 用线程池处理异步任务,替代
@Async以灵活控制线程资源 - 通过聚合器收集批次处理结果,确认100条数据全部处理完成后发送触发信号
完整代码示例:
@Bean @InboundChannelAdapter(value = "databaseInputChannel", poller = @Poller(advice = "messageTriggerAdvice")) public JdbcPollingChannelAdapter checkDbForInput() { // 过滤已处理数据,避免重复加载 JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(dataSource(), "select * from users where load_status != '3'"); adapter.setUpdateSql("update users set load_status = '3' where id in (:id)"); adapter.setMaxRows(100); return adapter; } @Bean public MessageTriggerAdvice messageTriggerAdvice() { MessageTriggerAdvice advice = new MessageTriggerAdvice(); // 指定触发轮询的通道 advice.setTriggerChannel("triggerChannel"); // 设置为等待触发信号再进行下一次轮询 advice.setWaitForTrigger(true); return advice; } @Bean public IntegrationFlow pollingFlow() { return IntegrationFlow.from("databaseInputChannel") // 拆分批次为单条消息 .split() // 使用线程池异步处理单条消息 .channel(MessageChannels.executor("databaseRowInputChannel", taskExecutor())) // 处理单条消息的业务逻辑 .handle(this::processMessage) // 聚合100条消息的处理结果,确认批次完成 .aggregate(aggregatorSpec -> aggregatorSpec .correlationStrategy(msg -> "batch") // 所有消息归为同一批次 .releaseStrategy(group -> group.size() == 100)) // 收集满100条时释放 // 发送批次完成信号,触发下一次轮询 .channel("triggerChannel") .get(); } @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(10); executor.setMaxPoolSize(20); executor.setQueueCapacity(50); executor.initialize(); return executor; } // 移除@Async,改用线程池处理异步 public void processMessage(Message<?> message) { // 你的业务处理逻辑 }
原理说明
MessageTriggerAdvice会拦截轮询器的执行,第一次轮询后进入等待状态,直到triggerChannel收到消息才会触发下一次轮询- 聚合器会收集拆分后的100条消息的处理结果,当所有消息处理完成后,向
triggerChannel发送消息,通知轮询器继续 - 使用
ExecutorChannel替代@Async,可以更直观地控制异步任务的线程池资源,避免@Async默认线程池的潜在问题
内容的提问来源于stack exchange,提问作者Andreas Laager
相关产品推荐
相关产品推荐

