《Spring Boot 实战》13.2 Spring Batch 实战

以每月结算逾期罚金为例讲透 Spring Batch:Job、Step、chunk 模型与事务边界、JobRepository 元数据表、JobParameters 与失败重启、跳过与重试策略,并诚实说明几万行以内的批处理用一句 SQL 往往比它更划算。

本节目标:掌握 Spring Batch 的 Job/Step/chunk 模型,理解 chunk 与事务边界的关系、JobRepository 的作用、JobParameters 如何支撑重启,并判断一个批处理任务到底值不值得上 Spring Batch。
适用版本:Spring Boot 4.1.x(Java 21)

13.2 Spring Batch 实战

13.1 解决的是「什么时候触发」,本节解决「触发之后怎么处理一大批数据」。借阅服务里有一类任务天然是批量的:每月把上个月的借阅流水汇总成报表、把逾期未还的记录逐条算出罚金并入库。这类任务有几个共同点——数据量大、要能重跑、要能记录「跑到哪了」。在 @Scheduled 方法里写个 for 循环能跑通,但一旦中途失败,你不知道处理到第几条、也不知道怎么接着跑。Spring Batch 就是为这些诉求设计的。

先给结论:Spring Batch 的能力很强,但引入成本也高。本节的顺序是「先讲它怎么用,再讲什么时候不该用」——最后那节不是客套,是很多团队真实的教训。

13.2.1 Spring Batch 的模型:Job / Step / chunk

三个核心概念:

  • Job:一次批处理作业的整体,比如「2026-09 罚金结算」。
  • Step:Job 里的一个阶段。一个 Job 可以有多个 Step,串行或并行。
  • chunk:Step 内部的处理单位。面向 chunk 的 Step 由「读—处理—写」三件套组成,按固定大小成批提交。

「读—处理—写」对应三个接口,它们是 Spring Batch 的核心抽象:

接口方法职责
ItemReader<T>T read()逐条读;返回 null 表示读完
ItemProcessor<I,O>O process(I item)转换/过滤;返回 null 表示丢弃这条
ItemWriter<T>void write(Chunk<? extends T>)成批写;一次收到一个 chunk 的数据

注意 ItemWriter.write 的参数是 Chunk 而不是 List——这是较新版本的口径(早期是 write(List)),写自定义 writer 时不要照抄老教程。Chunk 通过 getItems() 拿到本批数据。

13.2.2 一个完整的月报 Job

下面用「逾期罚金结算」把三件套串起来。假设 loan 表里有未归还且已逾期的借阅,罚金按逾期天数 × 每日费率计算,结果写入 fine_record 表。

先写一个分页读取的 ItemReader。真实项目可以用 JdbcPagingItemReader,但这里手写一个,是为了把「分页读取」这件事讲透:

import org.springframework.batch.infrastructure.item.ItemReader;
import org.springframework.jdbc.core.JdbcTemplate;
import java.time.LocalDate;
import java.util.Collections;
import java.util.Iterator;
import java.util.List;

class OverdueLoanReader implements ItemReader<Loan> {

    private final JdbcTemplate jdbc;
    private final int pageSize;
    private int page = 0;
    private Iterator<Loan> buffer = Collections.emptyIterator();

    OverdueLoanReader(JdbcTemplate jdbc, int pageSize) {
        this.jdbc = jdbc;
        this.pageSize = pageSize;
    }

    @Override
    public Loan read() {
        if (!buffer.hasNext() && !loadNextPage()) {
            return null;          // 返回 null:告诉 Batch 这个 Step 读完了
        }
        return buffer.next();
    }

    private boolean loadNextPage() {
        List<Loan> rows = jdbc.query("""
                select id, book_id, member_id, due_at
                from loan
                where status = 'BORROWED' and due_at < current_date
                order by id
                limit ? offset ?
                """,
                (rs, n) -> new Loan(rs.getLong("id"), rs.getLong("book_id"),
                        rs.getLong("member_id"), rs.getObject("due_at", LocalDate.class)),
                pageSize, page * pageSize);
        page++;
        buffer = rows.iterator();
        return !rows.isEmpty();
    }
}

再写 ItemProcessor 和 ItemWriter:

import org.springframework.batch.infrastructure.item.Chunk;
import org.springframework.batch.infrastructure.item.ItemProcessor;
import org.springframework.batch.infrastructure.item.ItemWriter;
import org.springframework.jdbc.core.JdbcTemplate;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.time.temporal.ChronoUnit;

class FineCalculationProcessor implements ItemProcessor<Loan, FineRecord> {

    private final BigDecimal dailyRate;

    FineCalculationProcessor(BigDecimal dailyRate) {
        this.dailyRate = dailyRate;
    }

    @Override
    public FineRecord process(Loan loan) {
        long overdueDays = ChronoUnit.DAYS.between(loan.dueAt(), LocalDate.now());
        if (overdueDays <= 0) {
            return null;          // 返回 null:过滤掉这条,不进入 writer
        }
        BigDecimal amount = dailyRate.multiply(BigDecimal.valueOf(overdueDays));
        return new FineRecord(loan.id(), loan.memberId(), amount);
    }
}

class FineWriter implements ItemWriter<FineRecord> {

    private final JdbcTemplate jdbc;

    FineWriter(JdbcTemplate jdbc) {
        this.jdbc = jdbc;
    }

    @Override
    public void write(Chunk<? extends FineRecord> chunk) {
        jdbc.batchUpdate("""
                insert into fine_record(loan_id, member_id, amount)
                values (?, ?, ?)
                """,
                chunk.getItems().stream()
                        .map(f -> new Object[]{f.loanId(), f.memberId(), f.amount()})
                        .toList());
    }
}

最后用 Java 配置把它们组装成 Job 与 Step:

import org.springframework.batch.core.job.Job;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.Step;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.transaction.PlatformTransactionManager;
import javax.sql.DataSource;
import java.math.BigDecimal;

@Configuration
class FineBatchConfig {

    @Bean
    Job settleFineJob(JobRepository jobRepository, Step settleFineStep) {
        return new JobBuilder("settleFineJob", jobRepository)
                .start(settleFineStep)
                .build();
    }

    @Bean
    Step settleFineStep(JobRepository jobRepository,
                        PlatformTransactionManager transactionManager,
                        DataSource dataSource) {
        return new StepBuilder("settleFineStep", jobRepository)
                .<Loan, FineRecord>chunk(200, transactionManager)
                .reader(new OverdueLoanReader(new JdbcTemplate(dataSource), 200))
                .processor(new FineCalculationProcessor(new BigDecimal("0.50")))
                .writer(new FineWriter(new JdbcTemplate(dataSource)))
                .build();
    }
}

两点关于写法的说明:其一,JobBuilder/StepBuilder 用构造器接收名字与 JobRepository,这是当前推荐风格;旧的 JobBuilderFactory/StepBuilderFactory 已被移除,不要再用。其二,Spring Boot 的 starter 通过自动配置提供 JobRepository、JobLauncher 等基础设施,通常不需要再手动加 @EnableBatchProcessing。还要留意 6.0 的包结构做过重组:ItemReader / ItemProcessor / ItemWriter / Chunk 现在在 org.springframework.batch.infrastructure.item 下,而 Job / Step 及其 Builder、JobParameters、JobRepository 仍在 org.springframework.batch.core.*。照抄 5.x 教程里的 org.springframework.batch.item.* 会直接编译失败——上面几段示例的 import 都已按 6.0 的实际包路径给出。

13.2.3 chunk 大小与事务边界

chunk(200, transactionManager) 里那个 200 不是随便填的——它直接决定事务边界。面向 chunk 的 Step 执行流程是:

  1. 从 reader 读 200 条;
  2. 每条过一遍 processor;
  3. 200 条交给 writer 一次性写出;
  4. 提交事务(transactionManager 在这里起作用);
  5. 回到第 1 步,直到 reader 返回 null。

也就是说,一个 chunk = 一个事务。这带来两个必须权衡的后果:

  • chunk 太大:单次事务时间长、持有的锁多、回滚代价大,数据库日志与锁竞争压力上升;失败时一个 chunk 要整体重来。
  • chunk 太小:提交次数多,事务开销与网络往返占比上升,整体变慢。

所以 chunk 大小的选择标准是「让单次事务的耗时与锁范围可控」,而不是「越大越快」。经验区间通常在百到千量级,具体值要结合单条处理耗时、数据库对长事务的容忍度来压测确定,没有通用数字。

还要注意一个语义细节:writer 是批量的,reader 是逐条的。 如果 writer 里对每条记录都单独发一条 SQL,那 chunk 的批量优势就没了——正确做法是像上面那样用 batchUpdate 一次提交一批。

13.2.4 JobRepository 与元数据表

JobRepository 是 Spring Batch 的「账本」:它把每次运行的 Job 实例、执行、步骤、参数、上下文都记下来。它落在哪里,是 4.x 里一个必须搞清楚的变化。

StarterJobRepository 存哪适用
spring-boot-starter-batch内存(不落库)简单、一次性、不需要重启续跑
spring-boot-starter-batch-jdbc数据库(JDBC)需要失败重启、需要查看历史执行

这是 4.0 的破坏性变更:默认的 spring-boot-starter-batch 改成了无数据库的内存模式,升级后元数据不再写进原来的库;要恢复「写库」的旧行为,必须换成 spring-boot-starter-batch-jdbc。如果你从 3.x 迁移过来,发现 BATCH_* 表不再更新,原因就在这里。

JDBC 模式下会用到一批以 BATCH_ 开头的元数据表,常见的包括 BATCH_JOB_INSTANCE(一次作业的唯一实例)、BATCH_JOB_EXECUTION(实例的每次执行,含状态)、BATCH_JOB_EXECUTION_PARAMS(本次参数)、BATCH_STEP_EXECUTION(每个 Step 的执行与读写计数)、以及两张执行上下文表 BATCH_JOB_EXECUTION_CONTEXT / BATCH_STEP_EXECUTION_CONTEXT(保存可续跑的进度)。

这些表由 starter 的 schema 初始化脚本创建。生产环境通常关掉自动初始化,改用 Flyway/Liquibase 这类迁移工具统一管理表结构——避免「应用启动顺手建表」带来的权限与版本一致性问题。spring.batch.jdbc 命名空间下有一组与 schema 初始化相关的属性(例如初始化出错是否继续由 spring.batch.jdbc.continue-on-error 控制),具体以官方文档为准。

13.2.5 JobParameters 与失败重启

JobInstance 的唯一性由「Job 名 + 标识性参数」决定。 这句话是理解重启的钥匙。用 JobParametersBuilder 传参:

import org.springframework.batch.core.job.parameters.JobParameters;
import org.springframework.batch.core.job.parameters.JobParametersBuilder;

JobParameters params = new JobParametersBuilder()
        .addString("month", "2026-09")               // 标识性参数:区分不同月份
        .addLong("run.id", System.currentTimeMillis()) // 每次运行都不同
        .toJobParameters();
jobLauncher.run(settleFineJob, params);

两个容易踩的点:

  • 同一组标识性参数重复运行会报错。 如果 2026-09 这个月已经成功跑完,再用完全相同的参数跑,Batch 会拒绝(认为这个 JobInstance 已完成)。要允许重跑,就加一个每次都变的 run.id;反过来,如果你希望「同一个月只能成功一次」,就不要加 run.id,让参数保持稳定。
  • run.id 会破坏「重启」语义。 失败重启时,参数必须与失败那次完全一致,Batch 才会识别出「这是同一个 JobInstance 的续跑」而不是「一个全新的实例」。所以 run.id 这类参数要谨慎使用:它换来的是「可重复运行」,代价是「不能用同一组参数续跑」。

重启的依据来自 JobRepository 记录的 JobExecution 状态。失败后,用相同的标识性参数再次 run,Batch 会找到那条 FAILED 的执行记录并从失败的 Step 续跑(已成功的 Step 默认不重跑)。续跑能做到「从哪断的接着跑」,靠的是 reader 的可重入性:面向 chunk 的 Step 默认按「已处理条数」定位,所以 reader 必须保证按稳定顺序读取(这就是上面 order by id 的意义)——顺序不稳定,续跑就会漏读或重读。

13.2.6 跳过与重试

一批几十万条数据里,总有几条是「脏的」——比如某条 loan 关联的 member 已不存在。让整批因为一条脏数据回滚显然不合理。Spring Batch 的容错能力通过 faultTolerant() 打开:

@Bean
Step settleFineStep(JobRepository jobRepository,
                    PlatformTransactionManager transactionManager,
                    DataSource dataSource) {
    return new StepBuilder("settleFineStep", jobRepository)
            .<Loan, FineRecord>chunk(200, transactionManager)
            .reader(new OverdueLoanReader(new JdbcTemplate(dataSource), 200))
            .processor(new FineCalculationProcessor(new BigDecimal("0.50")))
            .writer(new FineWriter(new JdbcTemplate(dataSource)))
            .faultTolerant()
            .skip(IllegalStateException.class).skipLimit(50)      // 脏数据跳过,最多 50 条
            .retry(TransientDataAccessException.class).retryLimit(3) // 瞬时故障重试 3 次
            .build();
}

跳过(skip) 和 重试(retry) 要分清:

  • skip 针对确定性错误——这条数据本身有问题,重试多少次都一样,只能记下来跳过。跳过的条数达到 skipLimit 后,整个 Step 判为失败。
  • retry 针对瞬时错误——比如数据库连接超时、死锁,重试可能成功。重试是「原地重试当前这条」还是「重试整个 chunk」,取决于配置的 retry 策略。

有一条隐蔽的代价必须知道:一旦对 chunk 开启 skip,Batch 会退化为「逐条提交」来精确定位是哪一条出错。 因为一个 chunk 是一条事务,出错时整个 chunk 已回滚,要找出坏数据只能把 chunk 拆成单条一条条试。所以容错不是免费的——它在「出错少」时几乎无感,在「出错多」时会显著变慢。别把 skipLimit 设得过大来「掩盖」系统性问题。

13.2.7 分区与并行 Step 的适用条件

单个 Step 是单线程的——reader 一条条读、writer 一批批写,都在一个线程里。要提速,Spring Batch 提供两条路:

  • 并行 Step:一个 Job 里多个互不依赖的 Step 用 TaskExecutor 并行跑。适合「不同 Step 之间无数据依赖」的场景。
  • 分区(Partitioning):把一个 Step 按数据范围切成多个分区,每个分区是一个独立的 Step 执行,由 Partitioner 划分、PartitionHandler 分发到线程池或多个执行器。适合「同一个逻辑步骤、数据可按 key 切分」的场景。

两者都有前提:数据库连接池要够大,否则并行度一上去就互相等连接;writer 要能接受并发写,否则会争锁。分区还要求分区键能均匀切分数据,切歪了会「一个分区忙死、其余空转」。判断是否值得:先确认瓶颈是 CPU 还是数据库——如果瓶颈是数据库(大多数批处理都是),把线程数开大只会让数据库更慢,这时该优化的是 SQL 和索引,不是并行度。

13.2.8 诚实说明:几万行以内,别急着上 Batch

Spring Batch 的能力很强,但它的复杂度是实打实的:要维护 BATCH_* 元数据表、要理解 JobInstance/JobExecution 的语义、要处理重启与容错、要配置 chunk 与并行。这些成本只有在数据量和可靠性要求真正达到时才能被摊平。

判断标准可以给得很直接:

数据规模 / 诉求建议
几千行、一次跑完、失败重跑整批可接受一句 SQL(insert ... select ...)或流式处理
几万行、需要分批、失败要能续跑流式分页处理(reader 式循环)足够
几十万行以上、多阶段、要重启与容错、要审计Spring Batch
需要并行/分区、跨多个数据源Spring Batch

换句话说:如果「把一条 SQL 跑两遍」是可接受的,那你根本不需要 Batch。 罚金结算这种「按逾期天数算一个值、批量插库」的活,很多情况下一条 insert into fine_record ... select ... from loan where ... 就完成了,而且它天然是原子的、失败重跑即幂等。只有当计算逻辑复杂到 SQL 写不动、或者数据量大到一条 SQL 撑不住、或者需要「记录进度、失败续跑」时,Spring Batch 的复杂度才物有所值。

选型时先问三个问题:数据量到几十万了吗?失败后能整批重跑吗?处理逻辑能在 SQL 里表达吗? 三个都是「是」/「能」,就先用 SQL;有一个是「否」,再考虑 Batch。

小结

  • Spring Batch 的模型是 Job → Step → chunk,面向 chunk 的 Step 由 ItemReader/ItemProcessor/ItemWriter 组成;ItemWriter.write 收的是 Chunk 而不是 List。
  • chunk(n, txManager) 的 n 就是事务边界:一个 chunk 一个事务;太大锁范围大、太小开销高,要按实测确定而非照抄。
  • 4.0 起 spring-boot-starter-batch 是内存 JobRepository;要写库必须用 spring-boot-starter-batch-jdbc——这是从 3.x 升级时 BATCH_* 表不再更新的根因。
  • JobInstance 的唯一性 = Job 名 + 标识性参数;失败重启要求参数完全一致;加 run.id 换来「可重复运行」但代价是「不能续跑」。
  • 续跑依赖 reader 按稳定顺序读取,所以分页查询必须有确定的 order by。
  • skip 处理确定性脏数据、retry 处理瞬时故障;对 chunk 开启 skip 会退化为逐条提交来定位坏数据,代价不小。
  • 并行与分区的前提是连接池够大、writer 能并发、分区键能切匀;瓶颈在数据库时,加并行度只会更慢。
  • 最诚实的一条:几万行以内、失败能整批重跑、逻辑能用 SQL 表达时,一句 SQL 往往比 Spring Batch 更划算。 Batch 是为「几十万行以上 + 需要进度、重启、容错」准备的。

13.3 回到调度本身,解决 @Scheduled 在多实例部署下的最后一个大坑:任务被重复执行,以及由此引出的一整套幂等与告警设计。

阅读导航:上一节:13.1 调度方案选型 · 下一节:13.3 分布式调度与幂等 。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「java」更多文章

  1. 《Spring Boot 入门》18.3 打包与运行
  2. 《Spring Boot 入门》18.2 实现
  3. 《Spring Boot 入门》18.1 需求与设计