本节目标:掌握 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 执行流程是:
- 从 reader 读 200 条;
- 每条过一遍 processor;
- 200 条交给 writer 一次性写出;
- 提交事务(
transactionManager在这里起作用); - 回到第 1 步,直到 reader 返回
null。
也就是说,一个 chunk = 一个事务。这带来两个必须权衡的后果:
- chunk 太大:单次事务时间长、持有的锁多、回滚代价大,数据库日志与锁竞争压力上升;失败时一个 chunk 要整体重来。
- chunk 太小:提交次数多,事务开销与网络往返占比上升,整体变慢。
所以 chunk 大小的选择标准是「让单次事务的耗时与锁范围可控」,而不是「越大越快」。经验区间通常在百到千量级,具体值要结合单条处理耗时、数据库对长事务的容忍度来压测确定,没有通用数字。
还要注意一个语义细节:writer 是批量的,reader 是逐条的。 如果 writer 里对每条记录都单独发一条 SQL,那 chunk 的批量优势就没了——正确做法是像上面那样用 batchUpdate 一次提交一批。
13.2.4 JobRepository 与元数据表
JobRepository 是 Spring Batch 的「账本」:它把每次运行的 Job 实例、执行、步骤、参数、上下文都记下来。它落在哪里,是 4.x 里一个必须搞清楚的变化。
| Starter | JobRepository 存哪 | 适用 |
|---|---|---|
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 分布式调度与幂等 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。