【问题标题】:Running job stops if start new running the same job - Spring Batch如果开始新运行相同的作业,则运行作业停止 - Spring Batch
【发布时间】:2021-06-25 08:08:36
【问题描述】:

问题

我有点困惑,因为当通过 HTTP 请求开始执行 Spring Batch 作业时,如果我收到另一个 HTTP 请求来启动相同的作业,但在作业执行时使用不同的参数,正在执行的作业已执行停止未完成并开始处理新作业。

上下文

我开发了一个 API REST 来加载和处理 Excel 文件的内容。 Web 服务公开了两个端点,一个用于加载、验证和存储数据库中 Excel 文件的内容,另一个用于开始处理存储在数据库中的记录。

它是如何工作的

  • POST /api/excel/上传 此端点接收 Excel 文件。当收到请求时,每个文件都被分配一个唯一的标识符并验证其内容。如果内容正确,则将其插入到等待处理的临时表中。

  • GET /api/Excel/process?id=x 该端点接收要处理的文件的标识符。当收到请求时,会启动一个 Spring Batch 作业来处理临时表中的记录。

一些代码

  • 控制器
@PostMapping(produces = {APPLICATION_JSON_VALUE})
public ResponseEntity<Page<ExcelLoad>> post(@RequestParam("file") MultipartFile multipartFile)
{
    return super.getResponse().returnPage(service.upload(multipartFile));
}

@GetMapping(value = "/process", produces = APPLICATION_JSON_VALUE)
public DeferredResult<ResponseEntity<Void>> get(@RequestParam("id") Integer idCarga)
{
    DeferredResult<ResponseEntity<Void>> response = new DeferredResult<>(1000L);
    response.onTimeout(() -> response.setResult(super.getResponse().returnVoid()));

    ForkJoinPool.commonPool().submit(() -> service.startJob(idCarga));

    return response;
}

我使用 DeferredResult 在收到请求后向客户端发送响应,而不等待作业完成

  • 服务
public void startJob(int idCarga)
{
    JobParameters params = new JobParametersBuilder()
            .addString("mainJob", String.valueOf(System.currentTimeMillis()))
            .addString("idCarga", String.valueOf(idCarga))
            .toJobParameters();

    try
    {
        jobLauncher.run(job, params);
    }
    catch (JobExecutionException e)
    {
        log.error("---ERROR: {}", e.getMessage());
    }
}
  • 批次
@Bean
public Step mainStep(ReaderImpl reader, ProcessorImpl processor, WriterImpl writer)
{
    return stepBuilderFactory.get("step")
            .<List<ExcelLoad>, Invoice>chunk(10)
            .reader(reader)
            .processor(processor)
            .writer(writer)
            .faultTolerant().skipPolicy(new ExceptionSkipPolicy())
            .listener(stepSkipListener)
            .build();
}

@Bean
public Job mainJob(Step mainStep)
{
    return jobBuilderFactory.get("mainJob")
                            .listener(mainJobExecutionListener)
                            .incrementer(new RunIdIncrementer())
                            .start(mainStep)
                            .build();
}

执行一些测试,我观察到以下行为:

  1. 如果我向端点/进程请求在不同时间处理每个文件:在这种情况下,存储在临时表中的所有记录都被处理:

    • 记录已处理的文件 1:3606(预期为 3606)。
    • 记录已处理的文件 2:1776(预期为 1776)。
  2. 如果我向端点 /process 发出请求以首先处理 file1,并且在它完成之前我再次请求处理 file2:在这种情况下,不会处理存储在临时表中的所有记录:

    • 记录已处理的文件 1:1080(预期 3606)
    • 记录已处理的文件 2:1774(预期为 1776)

【问题讨论】:

    标签: java excel spring-boot spring-batch


    【解决方案1】:

    JobLauncher 不会停止作业执行,它只会启动它们。 Spring Batch 提供的默认作业启动器是SimpleJobLauncher,它将作业启动委托给TaskExecutor。现在,根据您使用的任务执行器实现以及如何配置它来启动并发任务,您可以看到不同的行为。例如,当您启动一个新的作业执行并将一个新任务提交给任务执行器时,任务执行器可以决定如果所有工作人员都忙,则拒绝此提交,或者将其放入等待队列,或者停止另一个任务并提交新的那一个。这些策略取决于几个参数(TaskExecutor 实现、后台使用的队列类型、RejectedExecutionHandler 实现等)。

    在您的情况下,您似乎正在使用以下内容:

    ForkJoinPool.commonPool().submit(() -> service.startJob(idCarga));
    

    因此,您需要检查此池的行为,以了解它如何处理新任务提交(我想这就是停止您的工作的原因,但您需要确认这一点)。也就是说,我不明白你为什么需要这个。如果您的要求如下:

    我使用 DeferredResult 在收到请求后向客户端发送响应,而无需等待作业完成

    然后您可以在作业启动器中使用异步任务执行器实现(如ThreadPoolTaskExecutor),请参阅Running Jobs from within a Web Container

    【讨论】:

    • 我没有使用自定义任务执行器,也没有配置为启动并发任务。我需要用户 ForkJoinPool 向客户端发送响应,否则端点在作业完成之前不会响应
    • 我试过这个,行为是一样的:\@Bean public JobLauncher jobLauncher (JobRepository jobRepository, TaskExecutor taskExecutor) { SimpleJobLauncher jobLauncher = new SimpleJobLauncher (); jobLauncher.setJobRepository (jobRepository); jobLauncher.setTaskExecutor (taskExecutor);返回作业启动器; } \@Bean public TaskExecutor taskExecutor () { SimpleAsyncTaskExecutor taskExecutor = new SimpleAsyncTaskExecutor (); taskExecutor.setDaemon (true); taskExecutor.setThreadPriority (Thread.MIN_PRIORITY);返回任务执行器; }
    • otherwise the endpoint not response until job finish:这不正确。这是我链接的文档的摘录:The controller launches a Job using a JobLauncher that has been configured to launch asynchronously, which immediately returns a JobExecution. The Job will likely still be running, however, this nonblocking behaviour allows the controller to return immediately, which is required when handling an HttpRequest。启动器将立即返回。
    • 请编辑问题并添加代码,而不是在 cmets 中添加代码。 I have tried with this and the behavior is the same:声明JobLauncher bean 不是提供自定义启动器的方式,您需要提供自定义BatchConfigurer 并覆盖getJobLauncher。这就是为什么不考虑您的自定义启动器的原因。这在此处的文档中进行了解释:docs.spring.io/spring-batch/docs/4.3.x/reference/html/….
    • 正如您所建议的,我已经能够取消使用 DeferredResult 类创建和配置自定义任务执行器以启动并发任务。现在看来这两个文件都在正确处理,但作者没有在数据库中插入任何东西
    【解决方案2】:

    感谢@Mahmoud Ben Hassine 对answer 的帮助,我能够解决这个问题。为了帮助实施,如果有人提出这个问题,我分享在我的情况下已经解决问题的代码:

    • 控制器
    @Autowired
    private JobLauncher jobLauncher;
    
    @Autowired
    private Job job;
    
    @GetMapping(value = "/process", produces = APPLICATION_JSON_VALUE)
    public void get(@RequestParam("id") Integer idCarga) throws JobExecutionException
    {
        JobParameters params = new JobParametersBuilder()
                .addString("mainJob", String.valueOf(System.currentTimeMillis()))
                .addString("idCarga", String.valueOf(idCarga))
                .toJobParameters();
    
        jobLauncher.run(job, params);
    }
    
    • 批处理配置、作业和步骤
    @Configuration
    @EnableBatchProcessing
    public class BatchConfig extends DefaultBatchConfigurer
    {
        @Autowired
        private JobBuilderFactory jobBuilderFactory;
    
        @Autowired
        private StepBuilderFactory stepBuilderFactory;
    
        @Autowired
        private StepSkipListener stepSkipListener;
    
        @Autowired
        private MainJobExecutionListener mainJobExecutionListener;
    
        @Bean
        public TaskExecutor taskExecutor()
        {
            ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor();
            taskExecutor.setMaxPoolSize(10);
            taskExecutor.setThreadNamePrefix("batch-thread-");
    
            return taskExecutor;
        }
    
        @Bean
        public JobLauncher jobLauncher() throws Exception
        {
            SimpleJobLauncher jobLauncher = new SimpleJobLauncher();
            jobLauncher.setJobRepository(getJobRepository());
            jobLauncher.setTaskExecutor(taskExecutor());
            jobLauncher.afterPropertiesSet();
    
            return jobLauncher;
        }
    
        @Bean
        public Step mainStep(ReaderImpl reader, ProcessorImpl processor, WriterImpl writer)
        {
            return stepBuilderFactory.get("step")
                    .<List<ExcelLoad>, Invoice>chunk(10)
                    .reader(reader)
                    .processor(processor)
                    .writer(writer)
                    .faultTolerant().skipPolicy(new ExceptionSkipPolicy())
                    .listener(stepSkipListener)
                    .build();
        }
    
        @Bean
        public Job mainJob(Step mainStep)
        {
            return jobBuilderFactory.get("mainJob")
                                    .listener(mainJobExecutionListener)
                                    .incrementer(new RunIdIncrementer())
                                    .start(mainStep)
                                    .build();
        }
    }
    

    如果在应用此代码后,就像我遇到的那样,您在将记录插入数据库中时也遇到了问题,您可以通过 question 我也将适用于我的代码放入其中。

    【讨论】:

      猜你喜欢
      • 2017-11-09
      • 1970-01-01
      • 1970-01-01
      • 2015-03-05
      • 1970-01-01
      • 2019-11-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多