使用JobOperator重启失败批执行ID未恢复任务,反而生成新完成实例
Spring Batch 重启失败任务无执行逻辑问题
环境信息
- Java: 17
- Spring Boot: 3.0.5
- Spring Batch: 5.0
问题描述
为测试任务重启功能,故意让主批处理任务执行失败,此时BATCH_JOB_EXECUTION表中STATUS和EXIT_CODE字段均标记为FAILED。任务初始通过JobLauncher.run(Jobname, Jobparams)启动。
注意事项
- 任务未声明为Spring Bean
- 未使用自动增量器,任务构建代码如下:
Job testBatchResumeJob = new JobBuilder("TEST_BATCH_RESUME_JOB", jobRepository).start(step).build();
重启现象
使用jobOperator.restart(failedBatchExecutionId)尝试重启失败任务时,出现以下异常现象:
- 系统创建新的任务实例和步骤记录,
BATCH_JOB_EXECUTION和BATCH_STEP_EXECUTION表中任务及步骤状态均标记为COMPLETED - Reader、Writer和Processor组件完全未执行任何业务逻辑
已尝试但无效的方案
- 移除主批处理中的自动增量器
- 使用后置处理器Bean(全局配置和任务级配置均尝试)
@Bean public JobRegistryBeanPostProcessor jobRegistryBeanPostProcessor(@Qualifier("myBatchJobRegistry") JobRegistry jobRegistry) { JobRegistryBeanPostProcessor postProcessor = new JobRegistryBeanPostProcessor(); postProcessor.setJobRegistry(jobRegistry); return postProcessor; }
- 显式注册任务并移除后置处理器Bean
if (!jobRegistry.getJobNames().contains(testBatchResumeJob .getName())) { jobRegistry.register(new ReferenceJobFactory(testBatchResumeJob )); }
- 将任务改为Spring Bean并调整相关基础设施配置
示例代码
批处理通用自定义配置
@Configuration public class MyBatchConfig { @Primary @Bean(name = "myBatchDataSource") public DataSource batchDataSource() { DataSource myBatchDbSrc = DataSourceBuilder.create().username(getUsername()).password(getPassword()).url(getUrl()).build(); if (myBatchDbSrc != null && myBatchDbSrc instanceof HikariDataSource) { @SuppressWarnings("resource") HikariDataSource hikariDatsource = (HikariDataSource) myBatchDbSrc; hikariDatsource.setSchema(getSchema()); } return myBatchDbSrc; } @Bean(name = "transactionManager") public JdbcTransactionManager batchTransactionManager(@Qualifier("myBatchDataSource") DataSource dataSource) { return new JdbcTransactionManager(dataSource); } @Bean(name = "myBatchJobRepository") public JobRepository jobRepository(@Qualifier("myBatchDataSource") DataSource batchDataSource, @Qualifier("transactionManager") JdbcTransactionManager batchTransactionManager) throws Exception { JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean(); factory.setDataSource(batchDataSource); factory.setTransactionManager(batchTransactionManager); factory.afterPropertiesSet(); return factory.getObject(); } @Bean(name = "myBatchJobLauncher") public JobLauncher jobLauncher(@Qualifier("myBatchJobRepository") JobRepository jobRepository) throws Exception { TaskExecutorJobLauncher jobLauncher = new TaskExecutorJobLauncher(); jobLauncher.setJobRepository(jobRepository); jobLauncher.afterPropertiesSet(); return jobLauncher; } @Bean (name = "myBatchJobExplorer") public JobExplorer jobExplorer(@Qualifier("myBatchDataSource") DataSource dataSource,@Qualifier("transactionManager")JdbcTransactionManager batchTransactionManager) throws Exception { final JobExplorerFactoryBean bean = new JobExplorerFactoryBean(); bean.setDataSource(dataSource); bean.setTransactionManager(batchTransactionManager); bean.setTablePrefix("BATCH_"); bean.setJdbcOperations(new JdbcTemplate(dataSource)); bean.afterPropertiesSet(); return bean.getObject(); } @Bean (name ="myBatchJobRegistry") public JobRegistry jobRegistry() throws Exception { return new MapJobRegistry(); } @Bean (name = "myBatchJobOperator") public JobOperator jobOperator(@Qualifier("myBatchJobLauncher") JobLauncher jobLauncher, @Qualifier("myBatchJobRepository") JobRepository jobRepository, @Qualifier("myBatchJobRegistry") JobRegistry jobRegistry, @Qualifier("myBatchJobExplorer") JobExplorer jobExplorer) { final SimpleJobOperator jobOperator = new SimpleJobOperator(); jobOperator.setJobLauncher(jobLauncher); jobOperator.setJobRepository(jobRepository); jobOperator.setJobRegistry(jobRegistry); jobOperator.setJobExplorer(jobExplorer); return jobOperator; } }
批处理重启任务配置
@Component("TESTRESUMEBATCHJOB") @Configuration public class TestMyBatchResumeConfig implements IBatchJobFactory{ @Autowired private ApplicationContext appContext; @Autowired @Qualifier("transactionManager") private JdbcTransactionManager transactionManager; /** 任务仓库 */ @Autowired @Qualifier("myBatchJobRepository") private JobRepository jobRepository; @Autowired @Qualifier("myBatchJobRegistry") private JobRegistry jobRegistry; @Override public Job getBatchJob() throws Exception{ ItemProcessor<TestBatchModel, TestBatchModel> itemProcessor = appContext .getBean(TestBatchProcessor.class); ItemWriter<TestBatchModel> itemWriter = appContext.getBean(TestBatchWriter.class); TestBatchReader reader = appContext.getBean(TestBatchReader.class); Step step = new StepBuilder("TEST_BATCH_RESUME_STEP", jobRepository) .<TestBatchModel, TestBatchModel>chunk(5, transactionManager) .reader(reader.getPagingItemReader()).processor(itemProcessor).writer(itemWriter) .taskExecutor(testBatchTaskExecutor()) .throttleLimit(2).build(); Job testBatchResumeJob = new JobBuilder("TEST_BATCH_RESUME_JOB", jobRepository).start(step).build(); return testBatchResumeJob ; } public SimpleAsyncTaskExecutor testBatchTaskExecutor() { SimpleAsyncTaskExecutor acctStmtTaskExecuter = new SimpleAsyncTaskExecutor(); acctStmtTaskExecuter.setConcurrencyLimit(100); acctStmtTaskExecuter.setThreadPriority(1); acctStmtTaskExecuter.setThreadNamePrefix("TEST_BATCH_RESUME"); return acctStmtTaskExecuter; } }
ItemReader实现
@Component public class SubhayuTestBatchModelReader { @Autowired @Qualifier("myBatchDataSource") private DataSource myBatchDataSource; /** 常量SELECT_CLAUSE_PARTY */ private static final String SELECT_CLAUSE_PARTY = "SELECT SOURCE_REFERENCE_ID, PARTY_ID, PARTY_NAME"; /** 常量FROM_CLAUSE_PARTY */ private static final String FROM_CLAUSE_PARTY = "FROM MY_PARTY "; /** 常量WHERE_CLAUSE_PARTY */ private static final String WHERE_CLAUSE_PARTY = "CONSOLIDATE_STATEMENT_FLAG='Y'"; public JdbcPagingItemReader<SubhayuTestBatchModel> getPagingItemReader() throws Exception { JdbcPagingItemReader<SubhayuTestBatchModel> reader = new JdbcPagingItemReader<>(); reader.setDataSource(myBatchDataSource); reader.setFetchSize(5); reader.setPageSize(5); reader.setRowMapper(new BeanPropertyRowMapper<>(SubhayuTestBatchModel.class)); Map<String, Order> sortKeys = new HashMap<>(); sortKeys.put("PARTY_ID", Order.ASCENDING); SqlPagingQueryProviderFactoryBean factory = new SqlPagingQueryProviderFactoryBean(); factory.setDataSource(myBatchDataSource); factory.setSelectClause(SELECT_CLAUSE_PARTY); factory.setFromClause(FROM_CLAUSE_PARTY); factory.setWhereClause(WHERE_CLAUSE_PARTY); factory.setSortKeys(sortKeys); reader.setQueryProvider(factory.getObject()); reader.afterPropertiesSet(); return reader; } }
ItemWriter实现
@Component @StepScope public class TestBatchModelWriter implements ItemWriter<TestBatchModel> { /** 常量GREP_KEY */ private static final String GREP_KEY = "TEST_RESUME_BATCH_PROCESS|"; @Override public void write(Chunk<? extends TestBatchModel> arg0) throws Exception { for(TestBatchModel testModelItem : arg0 ) { LOGGER.info(GREP_KEY + "Inside Writer, Writing Item for Party ID: "+ testModelItem.getParty_id()); if(testModelItem.getParty_id().equals("TEST0027")) { throw new NullPointerException("For Party ID: "+ testModelItem.getParty_id()); } } } }
ItemProcessor实现
@Component @StepScope public class TestBatchModelProcessor implements ItemProcessor<TestBatchModel, TestBatchModel> { /** 常量GREP_KEY */ private static final String GREP_KEY = "TEST_RESUME_BATCH_PROCESS|"; @Override public TestBatchModel process(TestBatchModel arg0) throws Exception { LOGGER.info(GREP_KEY+ "Inside Processor for line item: "+ arg0.getParty_id() ); return arg0; } }
内容的提问来源于stack exchange,提问作者Subhayu
相关产品推荐
相关产品推荐

