Spring Batch远程分区Worker消息消费控制问题求助
Spring Batch远程分区Worker消费控制问题解决
问题描述
开发基于manager-worker架构的Spring Batch应用,通过ActiveMQ Artemis实现远程分区。期望Worker完成当前远程步骤后,再从队列消费下一条消息,但当前Listener会一次性消费队列中所有消息。已尝试将Listener的concurrency设置为1-1,但无效:启动Worker1后,artemis queue stat显示DELIVERING_COUNT变为10,Worker1仅处理第一个分区;启动Worker2后无法处理剩余分区,不符合预期。
原配置代码
Manager端配置
@Bean public DirectChannel arequests() { return new DirectChannel(); } @Bean public IntegrationFlow aoutboundFlow(ActiveMQConnectionFactory connectionFactory) { return IntegrationFlows.from(arequests()) .handle(Jms.outboundAdapter(connectionFactory).destination("testingReq")) .get(); } @Bean public DirectChannel breplies() { return new DirectChannel(); } @Bean public IntegrationFlow binboundFlow(ActiveMQConnectionFactory connectionFactory) { return IntegrationFlows.from(Jms.messageDrivenChannelAdapter(connectionFactory).destination("testingRes")) .channel(breplies()) .get(); } @Bean public Step managerStep() { return this.managerStepBuilderFactory.get("managerStep") .partitioner("workerStep", new BasicPartitioner()) .gridSize(GRID_SIZE) .outputChannel(arequests()) .inputChannel(breplies()) .build(); } @Bean public Job remotePartitioningJob() { return jobBuilderFactory.get("remotePartitioningJob") .incrementer(new RunIdIncrementer()) .start(managerStep()) .build(); }
Worker端配置
@Bean public DirectChannel crequests() { return new DirectChannel(); } @Bean public IntegrationFlow cinboundFlow(ActiveMQConnectionFactory connectionFactory) { return IntegrationFlows.from( Jms.messageDrivenChannelAdapter(connectionFactory) .destination("testingReq") .configureListenerContainer(c -> c.concurrency("1-1"))) .log(LoggingHandler.Level.INFO, "Received Message", m -> "Received message: " + m.getPayload()) .channel(crequests()) .get(); } @Bean public DirectChannel dreplies() { return new DirectChannel(); } @Bean public IntegrationFlow doutboundFlow(ActiveMQConnectionFactory connectionFactory) { return IntegrationFlows.from(dreplies()) .handle(Jms.outboundAdapter(connectionFactory).destination("testingRes")) .get(); } @Bean public Step workerStep() { return this.workerStepBuilderFactory.get("workerStep") .inputChannel(crequests()) .outputChannel(dreplies()) .tasklet(tasklet(null)) .build(); } @Bean @StepScope public Tasklet tasklet(@Value("#{stepExecutionContext['partition']}") String partition) { return (contribution, chunkContext) -> { log.info("Started"); Thread.sleep(20000); log.info("finished " + partition); return RepeatStatus.FINISHED; }; }
解决方案
核心问题是ActiveMQ Artemis默认的消息预取机制:消费者会一次性预取大量消息到本地,即使限制concurrency也会导致所有消息被标记为DELIVERING状态,其他Worker无法获取。需从以下3个方面调整:
1. 限制Worker的JMS消息预取数量
在Worker的cinboundFlow配置中,为ListenerContainer添加预取策略,设置每次仅预取1条消息:
@Bean public IntegrationFlow cinboundFlow(ActiveMQConnectionFactory connectionFactory) { return IntegrationFlows.from( Jms.messageDrivenChannelAdapter(connectionFactory) .destination("testingReq") .configureListenerContainer(c -> { c.concurrency("1-1"); // 开启事务,配合预取控制确保消息确认时机 c.getContainer().setSessionTransacted(true); // 设置接收超时,避免无消息时阻塞 c.getContainer().setReceiveTimeout(10000); // 针对ActiveMQ Artemis设置预取数为1 ((DefaultMessageListenerContainer) c.getContainer()).setDestinationResolver((session, destinationName, pubSubDomain) -> { Queue queue = session.createQueue(destinationName); ((ActiveMQSession) session).setPrefetchSize(queue, 1); return queue; }); })) .log(LoggingHandler.Level.INFO, "Received Message", m -> "Received message: " + m.getPayload()) .channel(crequests()) .get(); }
2. 调整消息确认机制
原AUTO确认模式会在消息被监听器接收后立即确认,导致队列认为消息已处理。改为SESSION_TRANSACTED模式,确保Worker完成Step任务后才提交事务、确认消息:
- 上述代码中已通过
setSessionTransacted(true)开启事务,无需额外修改。
3. 替换请求通道类型
将Worker的crequests通道从DirectChannel改为QueueChannel,确保消息按顺序处理,避免预取堆积:
@Bean public QueueChannel crequests() { return new QueueChannel(); }
验证效果
修改后启动Worker,通过artemis queue stat查看:
- DELIVERING_COUNT会保持为1(每个Worker对应1条消息)
- Worker完成当前任务后,才会从队列获取下一条消息
- 多个Worker可并行处理不同分区任务
内容的提问来源于stack exchange,提问作者KCJJ
相关产品推荐
相关产品推荐

