如何实现Spring Batch作业的每日分批定时执行?
解决方案:Spring Batch 每日分批读取CSV文件
不需要更换现有调度方式,Spring自带的@Scheduled结合cron表达式完全能满足每日执行的需求。以下是具体修改步骤:
1. 持久化读取偏移量(解决重启丢失进度问题)
静态变量count在应用重启后会丢失进度,因此需要将偏移量持久化到数据库。
第一步:创建偏移量跟踪实体与Repository
import jakarta.persistence.*; @Entity public class BatchOffset { @Id @GeneratedValue(strategy = GenerationType.IDENTITY) private Long id; private String jobName; // 关联对应的Batch作业名称 private int currentBatchIndex; // 记录当前已处理的批次序号(第0批读前100条,第1批读下100条...) // Getter & Setter }
import org.springframework.data.jpa.repository.JpaRepository; import java.util.Optional; public interface BatchOffsetRepository extends JpaRepository<BatchOffset, Long> { Optional<BatchOffset> findByJobName(String jobName); }
2. 修改Reader实现动态分批次读取
将FlatFileItemReader设置为StepScope,使其能获取每次Job执行时传入的偏移量参数,动态调整跳过行数和读取上限。
@Configuration @AllArgsConstructor public class SpringBatchConfig { private final GrpDataWriter writer; private static final int BATCH_SIZE = 100; // 标记为StepScope,允许注入Job参数 @Bean @StepScope public FlatFileItemReader<GrpData> reader(@Value("#{jobParameters['batchIndex']}") Integer batchIndex) { FlatFileItemReader<GrpData> itemReader = new FlatFileItemReader<>(); itemReader.setResource(new FileSystemResource("src/main/resources/grpData.csv")); itemReader.setName("dataReader"); // 跳过表头 + 已处理的所有记录数(batchIndex * 100) int linesToSkip = 1 + (batchIndex != null ? batchIndex * BATCH_SIZE : 0); itemReader.setLinesToSkip(linesToSkip); itemReader.setMaxItemCount(BATCH_SIZE); // 限制本次最多读100条 itemReader.setLineMapper(lineMapper()); return itemReader; } // 原lineMapper、processor方法保持不变... @Bean public Step step(JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("step", jobRepository) .<GrpData, GrpData>chunk(10, transactionManager) .reader(reader(null)) // StepScope会自动注入实际参数 .processor(processor()) .writer(writer) .build(); } // 原job方法保持不变... }
3. 修改调度器实现每日执行与进度更新
调整调度周期为每日执行,每次执行前读取当前偏移量,执行后更新进度,处理文件读取完成的边界情况。
@Component @AllArgsConstructor public class JobScheduler { private final JobLauncher jobLauncher; private final Job job; private final BatchOffsetRepository offsetRepository; private static final String JOB_NAME = "importGrpData"; private static final int BATCH_SIZE = 100; // cron表达式:每天凌晨0点执行 @Scheduled(cron = "0 0 0 * * ?") public void scheduleJob() { // 获取或初始化偏移量记录 BatchOffset offset = offsetRepository.findByJobName(JOB_NAME) .orElseGet(() -> { BatchOffset newOffset = new BatchOffset(); newOffset.setJobName(JOB_NAME); newOffset.setCurrentBatchIndex(0); return offsetRepository.save(newOffset); }); int currentBatchIndex = offset.getCurrentBatchIndex(); // 构造唯一的Job参数,确保Spring Batch允许重复执行 JobParameters jobParameters = new JobParametersBuilder() .addInt("batchIndex", currentBatchIndex) .addLong("timestamp", System.currentTimeMillis()) .toJobParameters(); try { JobExecution execution = jobLauncher.run(job, jobParameters); // 检查本次读取的记录数,判断是否已读完文件 StepExecution stepExecution = execution.getStepExecutions().stream().findFirst().orElse(null); if (stepExecution != null) { int readCount = stepExecution.getReadCount(); if (readCount < BATCH_SIZE) { // 文件已读完,可选择重置偏移量或停止调度 offset.setCurrentBatchIndex(0); } else { // 偏移量+1,下次读取下一批 offset.setCurrentBatchIndex(currentBatchIndex + 1); } offsetRepository.save(offset); } } catch (Exception e) { e.printStackTrace(); } } }
关键说明
- StepScope的作用:让Reader在每次Step执行时动态实例化,从而获取当前Job的参数(批次序号),实现动态读取逻辑。
- 偏移量持久化:确保应用重启后不会丢失读取进度,继续从上次的位置开始处理。
- 调度方式选择:当前需求用
@Scheduled的cron表达式足够,如果后续需要复杂调度(如节假日跳过),再考虑替换为Quartz等专业调度框架。
内容的提问来源于stack exchange,提问作者Arjun Koroth
相关产品推荐
相关产品推荐

