JdbcPollingChannelAdapter maxRows与Poller maxMessagesPerPoll的差异及并发疑问
我有多个复用相同逻辑的Polling流,通过数据库列id_channel区分不同的目标数据行。当前JdbcPollingChannelAdapter的setMaxRows固定为1,我理解每次数据库往返只会获取1行数据。现在有5条Polling流、10个线程,想了解这些流之间是如何竞争资源的?另外,在setMaxRows始终为1的情况下,设置Pollers.maxMessagesPerPoll对并发处理有影响吗?
我的application.properties(自定义数据源)配置如下:
spring.task.scheduling.pool.size=10 spring.pgsql.hikari.maximum-pool-size=10
流逻辑代码:
private MessageSource<Object> buildJdbcMessageSource(final int channel) { JdbcPollingChannelAdapter adapter = new JdbcPollingChannelAdapter(dataSource, FETCH_QUERY); adapter.setMaxRows(1); adapter.setUpdatePerRow(true); adapter.setSelectSqlParameterSource(new MapSqlParameterSource(Map.of("idCanal", channel))); adapter.setRowMapper((RowMapper<IntControle>) (rs, i) -> new IntControle(rs.getLong(1), rs.getInt(2), rs.getString(3))); adapter.setUpdateSql(UPDATE_QUERY); return adapter; } private IntegrationFlow buildIntegrationFlow(final int channel, final long rate, final int maxMessages) { return IntegrationFlows.from(buildJdbcMessageSource(channel), c -> c.poller(Pollers.fixedDelay(rate) .transactional(transactionInterceptor()) .maxMessagesPerPoll(maxMessages))) .split() .enrichHeaders(h -> h.header(MessageHeaders.ERROR_CHANNEL, ERROR_CHANNEL)) .channel(SybaseFlowConfiguration.SYBASE_SINK) .get(); } public IntegrationFlow pollingFlowChannel1() { return buildIntegrationFlow(1, properties.getChan1RateMs(), properties.getChan1MaxMessages()); } public IntegrationFlow pollingFlowChannel2() { return buildIntegrationFlow(2, properties.getChan2RateMs(), properties.getChan2MaxMessages()); } ...
1. 多Polling流的资源竞争逻辑
你的5条Polling流各自对应独立的JdbcPollingChannelAdapter,每个适配器通过idCanal参数查询不同id_channel的数据行,流之间不存在数据层面的竞争——因为它们查询的是数据库中完全隔离的数据集。
资源竞争主要体现在两个层面:
- 调度线程池:你配置的
spring.task.scheduling.pool.size=10是Spring调度器的线程池大小,用于触发Poller的轮询任务。5条流的轮询任务会从这10个线程中获取执行线程,只要线程池有空闲线程,所有流的轮询任务都能并行触发,不会互相阻塞。 - 数据库连接池:
spring.pgsql.hikari.maximum-pool-size=10是Hikari连接池的最大连接数。每条流的轮询操作(查询+更新)都会占用一个数据库连接,5条流同时执行时最多占用5个连接,远小于连接池上限,所以数据库连接不会成为瓶颈。只有当后续业务处理(比如发送到SYBASE_SINK的操作)也需要占用连接时,才可能出现连接竞争,但当前场景下轮询阶段的连接压力不大。
另外,每个适配器配置了transactional,每次轮询的查询和更新操作会在同一个事务中执行,配合updatePerRow=true,能确保获取到的行被立即标记(比如更新状态),不会被同流的后续轮询重复获取——不过因为流之间是按id_channel隔离的,这一点更多是保障单流内部的幂等性,而非解决流间竞争。
2. maxMessagesPerPoll在setMaxRows=1时的影响
setMaxRows=1控制的是单次数据库查询返回的行数,而maxMessagesPerPoll控制的是单次轮询任务中从MessageSource获取并发送的消息总数。
当setMaxRows=1时:
- 如果
maxMessagesPerPoll=1:每次轮询只会执行一次数据库查询,获取1行数据并发送成1条消息,轮询任务结束。 - 如果
maxMessagesPerPoll>1(比如设置为5):单次轮询任务会循环调用MessageSource的receive()方法,最多执行maxMessagesPerPoll次数据库查询——每次查询依然只返回1行数据(因为setMaxRows=1),直到没有数据可查或者达到次数上限。
这会带来两个直接影响:
- 并发效率:单次轮询任务会占用调度线程多次执行数据库操作,相当于在一个线程内串行处理多行同
id_channel的数据。如果你的SYBASE_SINK是支持并发的通道(比如ExecutorChannel),这些消息可以被并行处理,但轮询阶段的查询操作是串行的。 - 资源占用:单次轮询任务会多次占用数据库连接,每次查询都需要获取、释放连接(或复用连接池中的空闲连接),如果
maxMessagesPerPoll设置过大,可能会导致调度线程长时间被占用,影响其他流的轮询触发。
总结来说,maxMessagesPerPoll在setMaxRows=1时依然会影响并发效率,它决定了单次轮询任务能处理的消息数量,但所有查询操作都是串行执行的。如果想真正提升单流的并发处理能力,建议调整setMaxRows返回多行数据,配合split()处理器拆分消息,再通过并行通道处理。
内容的提问来源于stack exchange,提问作者Thiago Sayão

