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

使用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)尝试重启失败任务时,出现以下异常现象:

  1. 系统创建新的任务实例和步骤记录,BATCH_JOB_EXECUTION和BATCH_STEP_EXECUTION表中任务及步骤状态均标记为COMPLETED
  2. Reader、Writer和Processor组件完全未执行任何业务逻辑

已尝试但无效的方案

  1. 移除主批处理中的自动增量器
  2. 使用后置处理器Bean(全局配置和任务级配置均尝试)
@Bean
public JobRegistryBeanPostProcessor jobRegistryBeanPostProcessor(@Qualifier("myBatchJobRegistry") JobRegistry jobRegistry) {
JobRegistryBeanPostProcessor postProcessor = new JobRegistryBeanPostProcessor();
postProcessor.setJobRegistry(jobRegistry);
return postProcessor;
}
  1. 显式注册任务并移除后置处理器Bean
if (!jobRegistry.getJobNames().contains(testBatchResumeJob .getName())) {
 jobRegistry.register(new ReferenceJobFactory(testBatchResumeJob ));
}
  1. 将任务改为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 15:08:08