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

Spring Batch任务仅首次运行正常,后续触发无预期执行效果的问题求助

Spring Batch任务仅首次运行正常,后续触发无预期执行效果的问题求助

我通过API触发批处理任务,每次服务器重启后第一次调用API都能正常运行——会命中所有断点,数据也能按预期保存到数据库中。但在这之后,日志显示任务已启动且步骤正在执行,却不会命中任何断点,数据也没有按预期保存。

以下是我的代码和配置:

// Controller
@RestController
@Component
public class MyController {
    @Autowired
    JobLauncher jobLauncher;

    @Autowired
    Job insertJob;
    
    @PostMapping("/start")
    public String startBatchJob(@RequestParam String validation) {
        try {
            String jobId = String.valueOf(System.currentTimeMillis());
            JobParameters jobParameters = new JobParametersBuilder()
                    .addString("jobId", jobId)
                    .addString("validation", validation)
                    .toJobParameters();

            JobExecution jobExecution = jobLauncher.run(insertJob, jobParameters);

            if (jobExecution.getStatus().isUnsuccessful() || jobExecution.getStatus() == BatchStatus.FAILED)
                return "Job FAILED with Job ID: " + jobId + ". Status: " + jobExecution.getStatus();

            // Get job status or other details
            return "Job finished successfully with Job ID: " + jobId + ". Status: " + jobExecution.getStatus();
        } catch (JobExecutionException e) {
            e.printStackTrace();
            return "Error starting job: " + e.getMessage();
        }
    }
}
// Batch Config
@Configuration
public class SpringBatchConfig {

    @Autowired
    ScheduledTasks scheduledTasks;

    @Autowired
    ValidationRepository validationRepository;

    @Autowired
    @Lazy
    PlatformTransactionManager transactionManager;

    @Autowired
    @Lazy
    JobRepository jobRepository;


    public List<Personality> getData() {
        return scheduledTasks.callRestApiForData();
    }

    @Bean(name = "insertJob")
    public Job insertJob(BaseWriter<Personality> writer,
                             BaseProcessor<Personality, Personality> processor,
                             BaseReader<Personality> reader) {
        return new JobBuilder("insertJob", jobRepository)
                .incrementer(new RunIdIncrementer())
                .listener(insertJobListener()).start(step_1(writer, processor, reader)).build();
    }

    @Bean
    public Step step_1(BaseWriter<Personality> writer,
                       BaseProcessor<Personality, Personality> processor,
                       BaseReader<Personality> itemReader) {
        return new StepBuilder("step_1", jobRepository)
                .<Personality, Personality> chunk(200, transactionManager)
                .reader(itemReader)
                .processor(processor)
                .writer(writer)
                .build();
    }

    @Bean
    public JobExecutionListener insertJobListener() {
        return new InsertJobCompletionListener();
    }

}
// listener
@Slf4j
public class InsertJobCompletionListener extends JobExecutionListenerSupport {


    @Override
    public void afterJob(JobExecution jobExecution) {
        if (jobExecution.getStatus() == BatchStatus.COMPLETED) {
            log.info("UPDATE BATCH COMPLETED");
        }
        else if(jobExecution.getStatus() == BatchStatus.FAILED){
            log.info("UPDATE BATCH JOB FAILED TO COMPLETE");
        }


    }
}
// reader

public interface BaseReader<T> extends ItemReader<T> {
}


@Configuration
public class ExchangeRateBaseReader<T> implements BaseReader<Personality> {

    @Autowired
    ScheduledTasks scheduledTasks;

    private boolean read = false;
    private List<Personality> data;
    private int index = 0;

    @Override
    public Personality read() throws Exception {
        if (data == null) {
            data = getData();
        }
        Personality item = null;
        if (index < data.size()) {
            item = data.get(index);
            index++;
        }
        return item;
    }

    public List<Personality> getData() {
        return scheduledTasks.callRestApiForData();
    }
}
// processor

public interface BaseProcessor<I,O> extends ItemProcessor<I, O> {
}

@Configuration
@Slf4j
public class CurrencyExchangeProcessor<I,O> implements BaseProcessor<Personality, Personality> {
    @Override
    public Personality process(Personality Personality) throws Exception {
        return Personality;
    }
}
// writer
public interface BaseWriter<T> extends ItemWriter<T> {
}


@Configuration
@Slf4j
public class CurrencyExchangeWriter<T> implements BaseWriter<Personality> {

    @Autowired
    private ValidationRepository validationRepository;

    @Override
    public void write(Chunk<? extends Personality> chunk) throws Exception {
        log.info("Total number of items to save - {}",chunk.getItems().size());
        validationRepository.saveAll(chunk.getItems());
    }
}

提前感谢各位的帮助!

备注:内容来源于stack exchange,提问作者Puneeth C

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:32:59