Spring Batch 批处理

系统讲解 Spring Batch 批处理:Job/Step/Chunk 编程模型与 JobRepository 元数据表、ItemReader/Processor/Writer 组合、事务边界与提交间隔、重启续跑与失败重试跳过、本地与远程分区并行、Spring Batch 5 的 API 变化,以及生产监控与调优实践。

定时跑一次百万行数据迁移、每天凌晨对账、增量同步外部系统——这类任务用「循环 + 事务」硬写会迅速失控:中断了从哪继续?跑失败了重跑会不会重复写?单线程跑太慢怎么并行?Spring Batch 用 Job/Step/Chunk 三层模型 + 元数据仓库 + 可重启语义,把批处理的这些共性难题一次性解决。本文从编程模型讲到分区并行与生产调优。

一、批处理的核心问题

问题手工循环的痛点Spring Batch 的答案
中断续跑需自建断点记录JobRepository 记录执行状态,自动续跑
失败重试try/catch 手写RetryTemplate 声明式重试
容错跳过一行坏数据全批失败SkipPolicy 跳过并记录
并行提速手写线程池拆分Partitioning 本地/远程分区
事务边界全批一个事务,锁太久Chunk 提交,分批事务
可观测只有日志元数据表 + 监听器 + Micrometer
Spring Batch 的分层模型:
  Job       一次完整批处理(可含多个 Step)
    Step    一个独立阶段(读-处理-写 或 Tasklet)
      Chunk 事务性提交单元(N 条提交一次)
        ItemReader   →  ItemProcessor  →  ItemWriter

一句话总结: Spring Batch 的价值不是「帮你循环」,而是把中断续跑、失败重试、容错跳过、并行分区这些批处理共性难题标准化,让业务代码只关心「读什么、怎么转、写到哪」。

二、JobRepository 与元数据表

Spring Batch 的「可重启」靠元数据仓库(JobRepository)持久化执行状态,默认落库到一组 BATCH_* 表。

-- 核心元数据表
BATCH_JOB_INSTANCE      -- 作业实例(按 jobName + jobParameters 唯一)
BATCH_JOB_EXECUTION     -- 每次运行(含 START/END/STATUS/EXIT_CODE)
BATCH_JOB_EXECUTION_PARAMS -- 本次运行的参数
BATCH_STEP_EXECUTION    -- 每个 Step 的执行(含读写计数、提交次数)
BATCH_STEP_EXECUTION_CONTEXT -- Step 的上下文(断点位置)
BATCH_JOB_EXECUTION_CONTEXT  -- Job 的上下文
spring:
  batch:
    jdbc:
      initialize-schema: always   # 首次自动建表(生产建议 never + 手工 DDL)
    job:
      enabled: false              # 禁止启动时自动跑所有 Job(配合调度触发)
重启判定规则:
  JobInstance = jobName + 识别性 JobParameters
  同一 JobInstance 若已有 COMPLETED 的 Execution → 不能重跑(报 JobInstanceAlreadyCompleteException)
  若上次 FAILED → 可 restart,从上次失败的 Step 继续
  非识别参数(如时间戳)用 JobParametersBuilder.addLong("ts", ..., false) 排除

一句话总结: JobRepository 是 Spring Batch 的大脑,「重启从哪继续」完全由元数据表决定;设计 JobParameters 时想清楚哪些是识别性参数,直接决定能否重跑。

三、一个完整的 Chunk 型 Job

@Configuration
public class OrderSyncJobConfig {

    @Bean
    public Job orderSyncJob(JobRepository jobRepository, Step syncStep) {
        return new JobBuilder("orderSyncJob", jobRepository)
                .start(syncStep)
                .incrementer(new RunIdIncrementer())   // 每次新实例
                .build();
    }

    @Bean
    public Step syncStep(JobRepository jobRepository,
                         PlatformTransactionManager txManager,
                         ItemReader<OrderCsv> reader,
                         ItemProcessor<OrderCsv, Order> processor,
                         ItemWriter<Order> writer) {
        return new StepBuilder("syncStep", jobRepository)
                .<OrderCsv, Order>chunk(500, txManager)   // 每 500 条提交一次
                .reader(reader)
                .processor(processor)
                .writer(writer)
                .faultTolerant()
                    .skip(FlatFileParseException.class).skipLimit(100)
                    .retry(DeadlockLoserDataAccessException.class).retryLimit(3)
                .listener(new SyncStepListener())
                .build();
    }
}
@Bean
public FlatFileItemReader<OrderCsv> reader() {
    return new FlatFileItemReaderBuilder<OrderCsv>()
            .name("orderReader")
            .resource(new FileSystemResource("/data/orders.csv"))
            .linesToSkip(1)
            .delimited().delimiter(",")
            .names("orderNo", "userId", "amount", "createdAt")
            .targetType(OrderCsv.class)
            .build();
}

@Bean
public JdbcBatchItemWriter<Order> writer(DataSource dataSource) {
    return new JdbcBatchItemWriterBuilder<Order>()
            .dataSource(dataSource)
            .sql("INSERT INTO orders(order_no,user_id,amount) VALUES (:orderNo,:userId,:amount)")
            .beanMapped()
            .build();
}

一句话总结: 一个 Chunk 型 Step = chunk(size, txManager) + Reader/Processor/Writer 三段;chunk(500) 的 500 就是事务提交间隔,它同时决定吞吐与失败回滚的粒度。

四、Chunk 的事务与提交语义

Chunk 执行循环(chunk-size = 500):
  1. Reader 读 1 条 → Processor 处理 → 累积到 buffer
  2. 重复直到 buffer 满 500 条
  3. 打开事务,Writer 批量写 500 条,提交
  4. 更新 StepExecution 的 read/write/commit 计数
  5. 循环直到 Reader 返回 null

失败时:
  若开启 faultTolerant → 逐条重试/跳过定位坏数据
  否则整个 chunk 回滚,Step 标记 FAILED
参数影响建议
chunk size提交频率、回滚粒度100 ~ 1000,看单条耗时
事务隔离并发写入冲突按需降级到 READ_COMMITTED
saveState是否记录断点需要重启续跑则 true
Processor 无状态否则重启结果不一致避免在 Processor 里存可变状态
// Reader 必须「可重启」:saveState=true 时会记录 read.count 断点
// 重启时跳过已读记录,需保证数据源顺序稳定(按主键排序)
@Bean
public JdbcPagingItemReader<Order> reader(DataSource ds) {
    JdbcPagingItemReader<Order> r = new JdbcPagingItemReader<>();
    r.setDataSource(ds);
    r.setPageSize(1000);
    r.setSortKeys(Map.of("id", Order.ASCENDING));   // 顺序稳定才能正确续跑
    // ...
    return r;
}

一句话总结: Chunk 的事务边界就是提交间隔,重启续跑要求 Reader 输出顺序稳定(按主键排序),否则断点续跑会漏读或重复读。

五、容错:重试、跳过与跳过策略

.faultTolerant()
    // 跳过:坏数据不中断整批
    .skip(FlatFileParseException.class).skipLimit(100)
    .skip(ValidationException.class).skipLimit(50)
    .noSkip(FileNotFoundException.class)      // 这类异常绝不跳过

    // 重试:瞬时故障自动重试
    .retry(DeadlockLoserDataAccessException.class).retryLimit(3)
    .retry(TransientDataAccessException.class).retryLimit(3)

    // 回滚策略:哪些异常触发回滚
    .rollback(DataIntegrityViolationException.class)

    // 跳过监听:记录被跳过的数据
    .skipListener(new SkipListener<OrderCsv, Order>() {
        public void onSkipInRead(Throwable t) { log.warn("读取跳过", t); }
        public void onSkipInWrite(Order o, Throwable t) { log.warn("写入跳过 {}", o, t); }
        public void onSkipInProcess(OrderCsv c, Throwable t) { log.warn("处理跳过", t); }
    })
跳过 vs 重试的判定:
  重试:瞬时故障(死锁、超时、连接抖动)—— 同一条数据再试可能成功
  跳过:脏数据(格式错、校验失败)—— 再试也不会成功,跳过并记账
  注意:skipLimit 用满后 Step 仍会失败(Fail 状态),需人工介入
一个重要细节:
  开启 faultTolerant 后,chunk 内每条记录会被「逐条」处理以定位坏数据,
  Processor 与 Writer 可能被重复调用(回滚重放),
  因此 Writer 的幂等性很重要 —— 用 upsert 而非纯 insert。

一句话总结: 容错的黄金准则是「瞬时故障重试、脏数据跳过」;一旦开启 faultTolerant,Writer 必须幂等,否则重放会写出重复数据。

六、分区并行:本地与远程

6.1 本地分区(Local Partitioning)

把一个大数据集按维度切成 N 份,用线程池并行处理:

@Bean
public Step partitionedStep(JobRepository jobRepository,
                            Step workerStep,
                            Partitioner partitioner) {
    return new StepBuilder("partitionedStep", jobRepository)
            .partitioner("workerStep", partitioner)
            .step(workerStep)
            .gridSize(8)                                   // 8 个分区
            .taskExecutor(new SimpleAsyncTaskExecutor("batch-"))  // 本地并行
            .build();
}

@Bean
public Partitioner columnRangePartitioner(DataSource ds) {
    return new ColumnRangePartitioner(ds, "orders", "id");
}
分区键选择原则:
  1. 分区之间数据量尽量均衡(按主键范围、按取模分片)
  2. 分区之间无共享写(避免并发写同一行)
  3. 分区数 = gridSize,配合线程池大小
  4. 每个分区独立 StepExecution,独立断点

6.2 远程分区(Remote Chunking / Partitioning)

远程分区:
  Manager(主节点)切分数据 → 通过中间件分发 ExecutionContext
  Worker(多机)执行分片 Step
  中间件:Kafka / JMS / RabbitMQ

远程 Chunking(另一种模式):
  Manager 负责读,Worker 负责处理与写
  适合处理逻辑重、IO 轻的场景
  注意:网络往返会成为瓶颈
模式适合缺点
本地分区单机多核,数据源可并行读受单机资源限制
远程分区海量数据,多机横向扩展部署复杂,需中间件
远程 Chunking处理重、读轻网络往返瓶颈

一句话总结: 分区并行的关键在分区键——切得均衡、互不干扰才有加速比;单机用本地分区,海量数据用远程分区配中间件。

七、Spring Batch 5 的变化

变化说明
需要 Java 17+基线提升
@EnableBatchProcessing 不再必需自动配置已提供 JobRepository
Builder API 强制JobBuilder/StepBuilder 取代链式 JobBuilderFactory
JobOperator 接口调整获取运行信息的方式变化
元数据表结构调整部分表 schema 版本升级
默认 JobRepository 由 @EnableBatchProcessing 改为自动配置自定义需显式声明
// Spring Batch 4(旧)
@Bean
public Job job(JobBuilderFactory jobs, Step s) {
    return jobs.get("job").start(s).build();
}

// Spring Batch 5(新)
@Bean
public Job job(JobRepository repo, Step s) {
    return new JobBuilder("job", repo).start(s).build();
}
# 触发 Job 的常见方式
# 1. 定时任务(推荐)
spring.batch.job.enabled=false    # 关闭启动自动运行
# 再用 @Scheduled 或调度平台触发 JobLauncher.run(...)

# 2. 命令行一次性触发
java -jar app.jar --spring.batch.job.name=orderSyncJob run.id=$(date +%s)

一句话总结: 迁移到 Batch 5 的主要工作是把 Factory 换成 Builder 并去掉 @EnableBatchProcessing;生产环境记得 spring.batch.job.enabled=false,用调度器显式触发。

八、监控与调优实践

// 用监听器暴露进度与耗时
@Component
public class SyncStepListener implements StepExecutionListener {
    @Override
    public ExitStatus afterStep(StepExecution stepExecution) {
        log.info("step={} read={} write={} skip={} commit={} 耗时={}ms",
                stepExecution.getStepName(),
                stepExecution.getReadCount(),
                stepExecution.getWriteCount(),
                stepExecution.getSkipCount(),
                stepExecution.getCommitCount(),
                stepExecution.getEndTime().toEpochMilli()
                        - stepExecution.getStartTime().toEpochMilli());
        return stepExecution.getExitStatus();
    }
}
# 结合 Micrometer 暴露批处理指标
management:
  metrics:
    tags:
      application: order-batch
  endpoints:
    web:
      exposure:
        include: health,metrics,prometheus
调优清单:
  1. chunk size:单条慢 → 调大;回滚代价高 → 调小
  2. 批写:JdbcBatchItemWriter 的 batch size 与 chunk 对齐
  3. 读优化:JdbcPagingItemReader 用主键游标 + 覆盖索引
  4. 并行:本地分区 gridSize ≈ CPU 核数
  5. 数据库:关闭自动提交,避免每行一次往返
  6. 内存:流式读,避免一次性 load 全表
  7. 幂等:Writer 用 upsert,容忍重放
常见生产坑:
  1. 元数据表无限增长 → 定期清理历史 Execution
  2. 大事务锁表 → 调小 chunk,缩短事务
  3. 重启后重复写 → Writer 不幂等
  4. 分区数据倾斜 → 分区键选错
  5. 忘记 run.id → 第二次运行报 JobInstanceAlreadyCompleteException

一句话总结: 调优的核心是提交间隔、并行度与幂等写三件事;监控上抓 read/write/skip/commit 四个计数与耗时,就能快速定位瓶颈。

小结

维度要点
模型Job → Step → Chunk(Reader/Processor/Writer)
状态JobRepository 元数据表驱动重启续跑
事务chunk size = 提交间隔 = 回滚粒度
容错瞬时故障重试、脏数据跳过、Writer 幂等
并行本地分区(单机)/ 远程分区(多机)
版本Batch 5 用 Builder,去掉 @EnableBatchProcessing

Spring Batch 把批处理从「一堆 for 循环」升级为「有状态、可重启、可观测、可并行」的工程化作业。掌握 JobRepository 的重启语义与 Chunk 的事务边界,是写出可靠批处理的两块基石。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java」更多文章

  1. Testcontainers 集成测试
  2. Micrometer 可观测性
  3. JPMS 模块系统