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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 18:10:35