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

Spring Batch远程分片失败任务无法处理的问题排查

Spring Batch远程分片重启失败任务时管理器无法连接ActiveMQ Artemis队列问题

我在Spring Batch中实现远程分片(Remote Chunking)功能,采用ActiveMQ Artemis作为JMS消息代理,基础配置完成后正常任务可顺利运行并生成预期结果,管理器能正常记录任务完成状态。但测试失败任务处理机制时发现,重启失败任务后,Worker可正常连接Artemis服务器,管理器无法连接队列导致程序停滞。

连接工厂配置

@Configuration
public class ConnectionFactory {

    @Value("${broker.url}")
    private String brokerUrl;

    @Bean
    public ActiveMQConnectionFactory ActiveMQconnectionFactory() throws JMSException {
        ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory();
        connectionFactory.setBrokerURL(this.brokerUrl);
        return connectionFactory;
    }
}

管理器配置类

@Configuration
@EnableBatchIntegration
@EnableIntegration
@Profile("manager")
@Slf4j
public class ManagerConfiguration {
    @Autowired
    private DataSource dataSource;
    @Autowired
    private JobRepository jobRepository;
    @Autowired
    private JobListener jobListener;
    @Autowired
    private StepListener stepListener;
    @Autowired
    private ApplicationContext applicationContext;
    @Autowired
    private PlatformTransactionManager transactionManager;
    @Autowired
    private ManagerItems managerItems;
    @Autowired
    private ManagerChannel channel;
    final SimpleJobOperator jobOp = new SimpleJobOperator();
    private List<Long> failedJobs; 

    @Bean
    public TaskletStep managerStep() {
        return new RemoteChunkingManagerStepBuilder<CustomersInfoDTO,CustomersInfoDTO>("managerStep",jobRepository)
                .chunk(50)
                .outputChannel(channel.requests()) 
                .inputChannel(channel.replies())
                .listener(stepListener)
                .transactionManager(transactionManager)
                .reader(managerItems.infoReader()) // 通用Reader
                .allowStartIfComplete(Boolean.TRUE)
                .build();
    }

    @Bean
    public Job remoteChunkingJob() {
        log.info("Outside -> JOB");
        return new JobBuilder("remoteChunkingJob", jobRepository)
                .incrementer(new RunIdIncrementer())
                .listener(jobListener)
                .start(managerStep())
                .build();
    }

    @Bean
    public JobLauncher Launcher() throws Exception {
        TaskExecutorJobLauncher jobLauncher = new TaskExecutorJobLauncher();
        jobLauncher.setJobRepository(jobRepository);
        jobLauncher.afterPropertiesSet();
        return jobLauncher;
    }

    @Bean
    public JobOperator jobOp(final JobRegistry jobRegistry) throws Exception {
        jobOp.setJobLauncher(Launcher());
        jobOp.setJobRepository(jobRepository);
        jobOp.setJobRegistry(jobRegistry);
        jobOp.setJobExplorer(jobExp());
        return jobOp;
    }

    @Bean
    public JobExplorer jobExp() throws Exception {
        final JobExplorerFactoryBean bean = new JobExplorerFactoryBean();
        bean.setDataSource(dataSource);
        bean.setTransactionManager(transactionManager);
        bean.setTablePrefix("batch_");
        bean.setJdbcOperations(new JdbcTemplate(dataSource));
        bean.afterPropertiesSet();
        return bean.getObject();
    }

    @Bean
    @BeforeStep
    public Long getFailedInstance() throws JobInstanceAlreadyCompleteException, NoSuchJobException, NoSuchJobExecutionException, JobParametersInvalidException, JobRestartException {
        String sql = "SELECT job_execution_id from batch_job_execution where status='FAILED'";
        failedJobs =new JdbcTemplate(dataSource).query(sql, (rs, rowNum) -> {
            return rs.getLong("job_execution_id");
        });
        if(!failedJobs.isEmpty()) {
            jobOp.restart((Long) failedJobs.get(0));
        }
        return 1L;
    }
}

管理器通道配置类

@Configuration
@Slf4j
@Profile("manager")
public class ManagerChannel {

    @Bean
    public DirectChannel requests() {
        return new DirectChannel();
    }
    @Bean
    public QueueChannel replies() {
        return new QueueChannel();
    }

    @Bean
    public IntegrationFlow managerOutboundFlow(@Autowired ActiveMQConnectionFactory connectionFactory) {
        return IntegrationFlow.from(requests())
                .handle(Jms.outboundAdapter(connectionFactory).destination("requests"))
                .get();
    }

    @Bean
    public IntegrationFlow managerInboundFlow(@Autowired ActiveMQConnectionFactory connectionFactory) {
        return IntegrationFlow.from(Jms.messageDrivenChannelAdapter(connectionFactory).destination("replies"))
                .channel(replies())
                .get();
    }

}

Worker配置类

@Configuration
@Slf4j
@EnableBatchIntegration
@EnableIntegration
@Profile("worker")
public class WorkerConfiguration {

    @Autowired
    private RemoteChunkingWorkerBuilder<CustomersInfoDTO, CustomersInfoDTO> remoteChunkingWorkerBuilder;
    @Autowired
    private WorkerChannel channel;
    @Autowired
    private WorkerItems workerItems;


    @Bean
    public IntegrationFlow workerIntegrationFlow(){
        return new RemoteChunkingWorkerBuilder<CustomersInfoDTO,CustomersInfoDTO>()
                .inputChannel(channel.requests())
                .outputChannel(channel.replies())
                .itemProcessor(workerItems.itemProcessor()) // 通用Processor
                .itemWriter(workerItems.infoWriter()) // 通用Writer
                .build();
    }
}

Worker通道配置类

@Configuration
@Profile("worker")
public class WorkerChannel {

    @Bean
    public DirectChannel requests() {
        return new DirectChannel();
    }

    @Bean
    public DirectChannel replies() {
        return new DirectChannel();
    }

    @Bean
    public IntegrationFlow inboundFlow(@Autowired ActiveMQConnectionFactory connectionFactory) {
        return IntegrationFlow.from(Jms.messageDrivenChannelAdapter(connectionFactory).destination("requests"))
                .channel(requests())
                .get();
    }

    @Bean
    public IntegrationFlow outboundFlow(@Autowired ActiveMQConnectionFactory connectionFactory) {
        return IntegrationFlow.from(replies())
                .handle(Jms.outboundAdapter(connectionFactory).destination("replies"))
                .get();
    }

}

application.properties配置

spring.profiles.active=manager
broker.url=tcp://localhost:61616
spring.activemq.broker-url=tcp://localhost:61616
spring.activemq.user=root
spring.activemq.password=root

spring.datasource.url=jdbc:postgresql://localhost:5432/remote-chunking-artemis?createDatabaseIfNotExist=true
spring.datasource.username=postgres
spring.datasource.password=12345
spring.batch.jdbc.initialize-schema=always
spring.jpa.hibernate.ddl-auto=create
spring.datasource.driver-class-name=org.postgresql.Driver
spring.batch.jdbc.schema=classpath:/org/springframework/batch/core/schema-postgresql.sql

file.input = data/customers-100.csv

启动方式

通过命令mvn spring-boot:run -D spring-boot.run.profiles="{profile}"在不同终端分别启动Worker和Manager,其中{profile}替换为worker或manager。

补充说明

正常运行任务时,日志会记录任务ID对应的完成或失败状态;但重启失败任务(配置为拾取首个失败任务)时,日志输出如下后管理器完全卡住,查看ActiveMQ控制台发现管理器未订阅replies队列,仅存在Worker的连接:

2024-03-08T00:13:47.622+05:30  INFO 19656 --- [           main] o.s.b.c.l.support.SimpleJobOperator      : Checking status of job execution with id=689
2024-03-08T00:13:47.678+05:30  INFO 19656 --- [           main] o.s.b.c.l.support.SimpleJobOperator      : Attempting to resume job with name=remoteChunkingJob and parameters={'run.id':'{value=5, type=class java.lang.Long, identifying=true}'}
2024-03-08T00:13:47.814+05:30  INFO 19656 --- [           main] o.s.b.c.l.support.SimpleJobLauncher      : Job: [SimpleJob: [name=remoteChunkingJob]] launched with the following parameters: [{'run.id':'{value=5, type=class java.lang.Long, identifying=true}'}]
2024-03-08T00:13:47.847+05:30  INFO 19656 --- [           main] o.s.batch.core.job.SimpleStepHandler     : Executing step: [managerStep]
2024-03-08T00:13:47.875+05:30  INFO 19656 --- [           main] o.s.b.i.c.ChunkMessageChannelItemWriter  : Waiting for 2 results

内容的提问来源于stack exchange,提问作者Aayush

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 09:27:03