如何通过Spring Integration的JdbcPollingChannelAdapter轮询DB并传数据至监听器
问题解答
1. 实现每120秒轮询一次、持续2小时,并将数据传递到监听器
你的代码存在几个关键问题,导致监听器未触发,以下是修正方案:
(1)修正@InboundChannelAdapter配置
你手动创建JdbcPollingChannelAdapter并调用receive()的方式错误,@InboundChannelAdapter会自动轮询调用目标方法/消息源,无需手动触发。正确写法是直接将JdbcPollingChannelAdapter作为MessageSource Bean,并指定查询SQL:
@Component class AccountConfiguration { @Bean @InboundChannelAdapter(value = "outChannel", poller = @Poller(fixedDelay = "120000", maxMessagesPerPoll = "100")) // fixedDelay单位是毫秒,120秒需写120000 public MessageSource<List<Account>> accountMessageSource(DataSource dataSource) { JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter( dataSource, "SELECT * FROM account" // 补充你的查询SQL ); adapter.setRowMapper(new AccountMapper()); adapter.setMaxRowsPerPoll(100); // 控制每次轮询返回的最大行数 return adapter; } }
(2)正确绑定监听器到通道
你在generateFile中动态创建的IntegrationFlow未注册到Spring上下文,因此不会生效。推荐用@ServiceActivator直接绑定监听器到outChannel:
@Component class AccountMessageListener { @ServiceActivator(inputChannel = "outChannel") public void onMessage(List<Account> list){ System.out.println("Message received, total accounts: " + list.size()); } }
或者在配置类中注册全局IntegrationFlow:
@Configuration class IntegrationFlowConfig { @Bean public IntegrationFlow accountProcessingFlow(AccountMessageListener listener) { return IntegrationFlows.from("outChannel") .handle(listener, "onMessage") .get(); } }
(3)实现持续2小时的轮询
2小时共包含60次轮询(7200秒 ÷ 120秒/次),可以通过计数器控制轮询停止:
@Bean @InboundChannelAdapter(value = "outChannel", poller = @Poller(fixedDelay = "120000", maxMessagesPerPoll = "100")) public MessageSource<List<Account>> accountMessageSource(DataSource dataSource) { JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter( dataSource, "SELECT * FROM account" ); adapter.setRowMapper(new AccountMapper()); adapter.setMaxRowsPerPoll(100); // 包装MessageSource,添加轮询次数控制 return new MessageSource<>() { private int pollCounter = 0; private static final int MAX_POLL_TIMES = 60; // 2小时内的轮询次数 @Override public Message<List<Account>> receive() { if (pollCounter >= MAX_POLL_TIMES) { return null; // 返回null会终止轮询 } pollCounter++; return adapter.receive(); } }; }
2. 查看通道中的消息数量
默认的<int:channel>是DirectChannel,它直接转发消息,没有队列缓存,因此无法统计消息数量。需要将通道改为QueueChannel:
修改XML配置
<int:queue-channel id="outChannel" capacity="1000"/> <!-- capacity为队列最大容量,可选 -->
代码中获取消息数量
注入QueueChannel并调用getQueueSize()方法:
@Component class ChannelMonitor { @Autowired @Qualifier("outChannel") private QueueChannel outChannel; public int getPendingMessageCount() { return outChannel.getQueueSize(); } }
内容的提问来源于stack exchange,提问作者Nichole
相关产品推荐
相关产品推荐

