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
相关产品推荐
相关产品推荐

