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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 12:53:12