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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 16:13:36