Java 基础体系 · 第 94/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Spring Batch:Job、Step、Chunk、Checkpoint、重启和幂等
Spring Batch 适合处理“有明确开始和结束、需要记录进度、能够失败后恢复”的批处理任务,例如:
- 从 CSV 文件导入数据库;
- 分页读取数据库后生成文件;
- 批量调用外部服务并记录结果;
- 夜间汇总、对账、数据清理;
- 大量数据迁移和重算。
它解决的不是简单的“循环读取并处理”,而是如何在长时间运行、部分成功、进程崩溃、数据库连接中断、重复启动和并发执行等情况下,保留可解释的状态,并尽量安全地继续执行。
本文使用以下术语:
- Job:一次完整的批处理作业定义及其运行实例;
- Step:Job 中的一个阶段;
- Chunk:Step 中一次事务处理的项目集合;
- Checkpoint:为了重启而持久化的处理进度;
- 重启:在失败的 JobExecution 上继续执行,而不是从头盲目重跑;
- 幂等:同一业务输入重复执行一次或多次,最终业务结果一致。
这些概念之间不是并列关系,而是一条完整链路:
flowchart LR
A[JobInstance] --> B[JobExecution]
B --> C[StepExecution]
C --> D[Chunk事务]
D --> E[读取 Item]
E --> F[处理 Item]
F --> G[写入 Item]
G --> H[提交事务]
H --> I[保存 Checkpoint]
I --> D
G --> J[异常]
J --> K[回滚当前 Chunk]
K --> L[JobExecution FAILED]
L --> M[重新启动]
M --> I
关键点是:Spring Batch 通常以 Chunk 事务提交点 作为安全进度边界。异常发生后,已经提交的 Chunk 不会自动回滚;当前未提交的 Chunk 通常会回滚,并在重启时根据检查点重新读取。
一、先建立执行模型:Job、JobInstance、JobExecution 和 StepExecution
1. Job 是作业定义,不等于一次运行
Job 描述一个完整批处理流程。例如:
导入订单 Job
├── 读取 CSV
├── 校验并转换订单
└── 写入数据库
Job 通常由一个或多个 Step 组成:
@Bean
Job importJob(JobRepository jobRepository, Step importStep) {
return new JobBuilder("importJob", jobRepository)
.start(importStep)
.build();
}
这里的 "importJob" 是 Job 的名称。它标识的是“作业类型”,不是某一次具体运行。
2. JobInstance 表示同一业务作业的一次逻辑运行
Spring Batch 用以下概念区分“同一个作业的同一次运行”和“重新执行”:
JobInstance = Job 名称 + 识别性 JobParameters
假设 Job 名为 importJob:
| Job 参数 | JobInstance |
|---|---|
input=orders-2025-01.csv |
实例 A |
input=orders-2025-02.csv |
实例 B |
input=orders-2025-01.csv |
仍然是实例 A |
因此,第二次使用完全相同的识别性参数启动时,Spring Batch 不会认为这是一个全新的业务批次。
如果实例 A 已经成功完成,再次以相同参数启动,通常会得到“该 Job 已经完成,不能再次执行”一类的启动失败。这不是异常行为,而是 Spring Batch 防止重复执行的默认保护。
3. JobExecution 表示一次尝试
一个 JobInstance 可以拥有多个 JobExecution:
JobInstance: importJob + input=orders-2025-01.csv
├── JobExecution 1: FAILED
└── JobExecution 2: COMPLETED
第一次执行失败,第二次针对同一个 JobInstance 启动,就是一次新的 JobExecution。第二次执行可以读取第一次执行留下的检查点。
常见状态包括:
STARTINGSTARTEDSTOPPINGSTOPPEDFAILEDCOMPLETEDABANDONEDUNKNOWN
其中:
FAILED通常表示可以考虑重启;COMPLETED表示该 JobInstance 已完成,不能再以相同识别性参数普通重启;STOPPED表示被请求停止,通常可以继续;ABANDONED表示明确放弃,不应再自动重启。
ExitStatus 是另一类状态信息,通常用于更细粒度地表达退出原因,例如 COMPLETED WITH SKIPS。不要把 BatchStatus 和 ExitStatus 当成同一个概念。
4. StepExecution 记录某个 Step 的一次执行
如果一个 Job 有三个 Step,每个 JobExecution 都会对应这些 Step 的执行记录:
JobExecution
├── StepExecution: readFile
├── StepExecution: transform
└── StepExecution: writeDatabase
StepExecution 中通常会记录:
- 读取数量;
- 写入数量;
- 过滤数量;
- 跳过数量;
- 提交次数;
- 回滚次数;
- 当前状态;
- Step 级执行上下文。
这些信息既用于重启,也用于监控和诊断。
二、JobRepository:状态为什么能在进程退出后保留
Spring Batch 不把执行状态只放在 Java 对象中,而是通过 JobRepository 持久化:
JobInstance
JobExecution
StepExecution
ExecutionContext
默认的 JDBC 实现会把这些数据保存到数据库,例如:
BATCH_JOB_INSTANCE
BATCH_JOB_EXECUTION
BATCH_JOB_EXECUTION_PARAMS
BATCH_STEP_EXECUTION
BATCH_STEP_EXECUTION_CONTEXT
BATCH_JOB_EXECUTION_CONTEXT
具体表结构由 Spring Batch 版本提供的数据库脚本决定。不同数据库的建表脚本不同,应使用与实际数据库匹配的官方 schema,而不是手工猜测字段。
在 Spring Boot 中,通常由自动配置创建 JobRepository、事务管理器以及必要的批处理基础设施。开发环境可以允许自动初始化 schema:
spring.batch.jdbc.initialize-schema=always
生产环境通常应使用数据库迁移工具管理这些表,并避免应用每次启动都初始化或删除 Batch 元数据。
为什么必须使用持久化 Repository
如果进度只保存在内存中:
long currentIndex = 100000;
JVM 崩溃后,程序无法知道已经处理到哪里,只能:
- 从头开始;
- 自己扫描目标系统推断进度;
- 冒险跳过一部分数据。
而 JobRepository 将“这次 Job 的身份、每个 Step 的统计和检查点”持久化下来,使重启具有明确依据。
三、Step:一个可独立观察和重启的处理阶段
Step 是 Job 的执行单元。一个 Step 通常完成一类相对独立的工作:
Job
├── Step 1:读取原始文件并导入临时表
├── Step 2:校验临时表数据
└── Step 3:合并到正式表
Step 的划分会影响:
- 事务边界;
- 检查点;
- 重启位置;
- 统计信息;
- 错误处理策略;
- 哪些阶段可以被单独跳过或重新执行。
如果把所有工作写进一个 Step,故障后可能需要重新执行很长的阶段。如果拆得过细,又会增加状态管理和阶段间数据传递的复杂度。
Step 有两种常见模型
1. Tasklet Step
Tasklet 适合一次性动作,例如:
- 创建目录;
- 删除临时文件;
- 执行一条存储过程;
- 调用一次归档操作。
@Bean
Step cleanupStep(JobRepository jobRepository,
PlatformTransactionManager transactionManager) {
return new StepBuilder("cleanupStep", jobRepository)
.tasklet((contribution, chunkContext) -> {
// 执行一次清理操作
return RepeatStatus.FINISHED;
}, transactionManager)
.build();
}
Tasklet 不代表天然可重启。它是否能正确重启,取决于任务本身是否记录状态、是否能够识别已经完成的工作。
2. Chunk-oriented Step
Chunk Step 适合大量项目的循环处理:
读取一个 Item
处理一个 Item
收集若干 Item
在事务中批量写入
提交事务
保存检查点
典型定义如下:
@Bean
Step importStep(JobRepository jobRepository,
PlatformTransactionManager transactionManager,
ItemReader<Order> reader,
ItemProcessor<Order, Order> processor,
ItemWriter<Order> writer) {
return new StepBuilder("importStep", jobRepository)
.<Order, Order>chunk(100, transactionManager)
.reader(reader)
.processor(processor)
.writer(writer)
.build();
}
<Order, Order> 表示:
- Reader 输出
Order; - Processor 输入
Order,输出Order; - Writer 接收一组
Order。
四、Chunk:事务处理的基本单位
1. Chunk 的准确含义
如果配置:
.chunk(3, transactionManager)
通常表示:
read item 1
process item 1
read item 2
process item 2
read item 3
process item 3
write item 1, item 2, item 3
commit
这里的 3 是提交间隔,也常称为 chunk size。它不是简单的“每 3 条调用一次 writer”规则,而是通常对应一个事务处理块。
默认的典型语义是:
一个 Chunk = 一次事务中的读取、处理和写入结果
实际事务边界由 Step 配置、事务管理器、Reader/Writer 行为以及错误处理策略共同决定,因此不能把所有 Reader 的内部行为都简化为“读取操作也一定由同一个数据库事务保护”。
2. Chunk 的数量和提交次数
假设总项目数为 N,chunk size 为 C,忽略过滤、跳过和提前结束:
提交次数 = ceil(N / C)
例如:
N = 10
C = 3
处理过程是:
Chunk 1: item 1 ~ item 3 -> commit
Chunk 2: item 4 ~ item 6 -> commit
Chunk 3: item 7 ~ item 9 -> commit
Chunk 4: item 10 -> commit
提交次数:
ceil(10 / 3) = 4
3. Chunk 失败时发生什么
假设 C = 3:
Chunk 1: item 1, 2, 3 -> 已提交
Chunk 2: item 4, 5, 6 -> 写入 item 5 时失败
如果异常没有被跳过或恢复:
- Chunk 2 的事务回滚;
- item 4、5、6 在事务内的写入结果被回滚;
- JobExecution 或 StepExecution 进入失败状态;
- 已提交的 item 1、2、3 不会被这个回滚撤销;
- 重启时通常从 Chunk 1 之后的检查点继续。
所以,Chunk 提供的是:
已提交 Chunk 的稳定性
+
未提交 Chunk 的整体回滚
它不是每条 Item 一个事务,也不是整个 Job 一个事务。
4. Chunk size 的实际取舍
Chunk 太小:
- 提交次数多;
- 数据库日志和事务管理开销增加;
- 吞吐量可能下降;
- 但单次失败回滚的数据较少。
Chunk 太大:
- 单事务锁持有时间变长;
- 内存中暂存的项目更多;
- 单次失败回滚更多工作;
- 某些数据库批量写入可能产生更大的日志压力。
Chunk size 的正确选择取决于:
- 单条记录大小;
- Reader 和 Writer 的实现;
- 数据库锁和日志行为;
- 外部服务调用耗时;
- 允许的重放范围;
- 失败恢复时间。
不能仅凭“批量越大越快”作出结论。
五、Reader、Processor、Writer 的数据流
Chunk Step 的基本数据流如下:
ItemReader<T>
↓
T
ItemProcessor<T, R>
↓
R 或 null
ItemWriter<R>
↓
提交事务
1. ItemReader
Reader 每次返回一个项目:
public interface ItemReader<T> {
T read() throws Exception;
}
当没有更多数据时返回 null。
常见 Reader 包括:
FlatFileItemReader:读取文本文件;JdbcCursorItemReader:使用 JDBC 游标读取;JdbcPagingItemReader:分页读取数据库;JpaPagingItemReader:分页读取 JPA 数据;- 自定义 Reader:读取 API、消息或其他来源。
Reader 的关键问题不是“能否读出数据”,而是“能否在失败后准确恢复位置”。
2. ItemProcessor
Processor 用于转换和校验:
public interface ItemProcessor<I, O> {
O process(I item) throws Exception;
}
返回 null 通常表示过滤该项目。被过滤的项目不会传给 Writer,但会影响处理统计和业务理解。
例如:
@Bean
ItemProcessor<Order, Order> processor() {
return order -> {
if (order.amount().signum() < 0) {
throw new IllegalArgumentException("金额不能为负数");
}
return order;
};
}
3. ItemWriter
Writer 一次接收当前 Chunk 中的多个项目:
public interface ItemWriter<T> {
void write(Chunk<? extends T> chunk) throws Exception;
}
Writer 的写入结果通常与当前事务一起提交或回滚。对于数据库 Writer,这个事务边界通常比较清晰;对于 HTTP、文件系统、消息系统等外部副作用,则需要单独分析。
六、一个端到端的 CSV 导入示例
下面的示例使用:
- Spring Boot 自动配置;
- Spring Batch Chunk Step;
- H2 数据库;
- CSV 文件作为输入;
- JDBC 批量写入;
inputJob 参数指定输入文件。
示例 API 使用 Spring Batch 5 风格的 JobBuilder、StepBuilder 和 chunk(size, transactionManager) 方法。具体 Spring Boot 版本应使用其依赖管理提供的兼容 Spring Batch 版本,不应手动混用不兼容的 Spring Framework、Spring Batch 和 Boot 版本。
1. 依赖
Maven 依赖的核心部分如下:
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-batch</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>com.h2database</groupId>
<artifactId>h2</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
Java 25 LTS 可以作为运行环境,但是否能够使用某个具体 Spring Boot 小版本,应以该版本的官方兼容性说明为准。Java 版本不会改变 Job、Step 和 Chunk 的事务语义。
2. 配置文件
spring.datasource.url=jdbc:h2:file:./data/batch-db
spring.datasource.username=sa
spring.datasource.password=
spring.datasource.driver-class-name=org.h2.Driver
spring.batch.jdbc.initialize-schema=always
spring.batch.job.enabled=true
jdbc:h2:file:./data/batch-db 使数据库状态保存在文件中。若改成内存数据库,应用进程退出后 Batch 元数据也会消失,无法演示跨进程重启。
生产环境不应简单照搬 initialize-schema=always,而应使用正式数据库和迁移脚本管理 Batch 元数据。
3. 输入文件
创建 orders.csv:
id,amount
1,10.50
2,20.00
3,7.80
4,100.00
5,12.30
4. 数据表
可以在应用初始化 SQL 中创建业务表:
create table if not exists imported_order (
id bigint primary key,
amount decimal(19, 2) not null
);
id 设置为主键有两个作用:
- 保证业务上不会存在两个相同订单;
- 在重复处理时让数据库暴露冲突,而不是静默产生重复数据。
但主键冲突本身不是完整的幂等方案,后文会说明如何根据业务选择忽略、更新或拒绝。
5. Java 代码
package example.batch;
import java.math.BigDecimal;
import javax.sql.DataSource;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.configuration.annotation.StepScope;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.database.JdbcBatchItemWriter;
import org.springframework.batch.item.database.builder.JdbcBatchItemWriterBuilder;
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.LineMapper;
import org.springframework.batch.item.file.mapping.BeanWrapperFieldSetMapper;
import org.springframework.batch.item.file.mapping.DefaultLineMapper;
import org.springframework.batch.item.file.transform.DelimitedLineTokenizer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.Resource;
import org.springframework.transaction.PlatformTransactionManager;
@Configuration
public class BatchConfig {
@Bean
@StepScope
FlatFileItemReader<Order> orderReader(
@Value("#{jobParameters['input']}") Resource input) {
var tokenizer = new DelimitedLineTokenizer();
tokenizer.setNames("id", "amount");
var fieldSetMapper = new BeanWrapperFieldSetMapper<Order>();
fieldSetMapper.setTargetType(Order.class);
LineMapper<Order> lineMapper = new DefaultLineMapper<>() {{
setLineTokenizer(tokenizer);
setFieldSetMapper(fieldSetMapper);
}};
return new FlatFileItemReader<Order>() {{
setName("orderReader");
setResource(input);
setLinesToSkip(1);
setLineMapper(lineMapper);
}};
}
@Bean
ItemProcessor<Order, Order> orderProcessor() {
return order -> {
if (order.getId() <= 0) {
throw new IllegalArgumentException("订单 ID 必须为正数");
}
if (order.getAmount() == null
|| order.getAmount().signum() < 0) {
throw new IllegalArgumentException("订单金额不能为负数");
}
return order;
};
}
@Bean
JdbcBatchItemWriter<Order> orderWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<Order>()
.dataSource(dataSource)
.sql("""
insert into imported_order(id, amount)
values (:id, :amount)
""")
.beanMapped()
.build();
}
@Bean
Step importStep(
JobRepository jobRepository,
PlatformTransactionManager transactionManager,
ItemReader<Order> orderReader,
ItemProcessor<Order, Order> orderProcessor,
ItemWriter<Order> orderWriter) {
return new StepBuilder("importStep", jobRepository)
.<Order, Order>chunk(3, transactionManager)
.reader(orderReader)
.processor(orderProcessor)
.writer(orderWriter)
.build();
}
@Bean
Job importJob(JobRepository jobRepository, Step importStep) {
return new JobBuilder("importJob", jobRepository)
.start(importStep)
.build();
}
public static class Order {
private long id;
private BigDecimal amount;
public long getId() {
return id;
}
public void setId(long id) {
this.id = id;
}
public BigDecimal getAmount() {
return amount;
}
public void setAmount(BigDecimal amount) {
this.amount = amount;
}
}
}
启动类可以是:
package example.batch;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class BatchApplication {
public static void main(String[] args) {
SpringApplication.run(BatchApplication.class, args);
}
}
使用命令行启动:
java -jar batch-app.jar input=file:/absolute/path/orders.csv
这里的 input 会成为 Job 参数,并被 @StepScope Reader 延迟获取。@StepScope 很重要:Reader Bean 不在应用上下文创建时立即解析 Job 参数,而是在 Step 执行时解析当前 JobExecution 的参数。
预期结果是:
Chunk 1: 订单 1、2、3 -> 写入并提交
Chunk 2: 订单 4、5 -> 写入并提交
数据库中应有五行数据:
select id, amount
from imported_order
order by id;
结果:
1 | 10.50
2 | 20.00
3 | 7.80
4 | 100.00
5 | 12.30
由于 chunk(3, ...),前 3 条形成第一个事务,后 2 条形成第二个事务。
七、Checkpoint:重启依赖的进度快照
1. Checkpoint 是什么
Checkpoint 是“已经安全处理到哪里”的持久化记录。
对一个文件 Reader 来说,检查点可能包含:
已经读取的条数
当前行号
资源位置
Reader 自身恢复所需的其他状态
对分页数据库 Reader 来说,可能包含:
当前页
页内位置
排序键
但具体字段由 Reader 实现决定,应用不应依赖某个具体实现的内部字段名称,除非明确接受版本耦合。
Checkpoint 通常保存在 ExecutionContext 中:
StepExecution
└── ExecutionContext
├── reader 的恢复状态
├── writer 的恢复状态
└── 自定义状态
ExecutionContext 是键值形式的持久化上下文。它不是普通的临时 Map,也不是用来保存任意大对象的缓存。
2. Checkpoint 何时更新
在典型 Chunk 流程中:
开始事务
读取和处理一批 Item
Writer 写入
更新 Step 统计
提交事务
持久化执行上下文
具体内部时序由 Spring Batch 的实现和事务配置决定,但可以使用一个可靠的抽象:
一个检查点只有在对应处理进度被成功提交并持久化后,才应被视为可恢复进度。
如果在当前 Chunk 处理中 JVM 崩溃:
上一个检查点:item 3
当前处理到:item 5
item 4、5 尚未提交
重启时通常从 item 4 附近重新开始,而不是从 item 6 开始。
3. 完整算例
设:
总数据:item 1 到 item 10
chunk size:3
处理过程:
| 阶段 | 当前 Chunk | 结果 | 检查点 |
|---|---|---|---|
| 1 | 1、2、3 | 提交成功 | 3 |
| 2 | 4、5、6 | item 5 写入失败,回滚 | 仍为 3 |
| 3 | 重启 | 重新读取 4、5、6 | 处理成功后更新为 6 |
| 4 | 7、8、9 | 提交成功 | 9 |
| 5 | 10 | 提交成功 | 10 |
这里的“重新读取 4、5、6”不是浪费,而是为了保证检查点和目标写入状态一致。
4. Checkpoint 不等于业务进度
这是一个常见误解。
假设 Writer 调用外部支付系统:
1. 外部支付已成功
2. 本地程序在保存检查点前崩溃
3. 重启后重新读取该订单
4. 再次调用外部支付
Spring Batch 只能知道自己的 Step 进度,不能自动知道外部支付是否已经成功。因此:
Checkpoint 解决的是批处理框架的恢复位置;
它不自动解决外部系统的业务状态一致性。
八、重启:不是“再次运行”,而是恢复失败的 JobInstance
1. 一个正确的重启流程
假设第一次运行:
java -jar batch-app.jar input=file:/data/orders.csv
在第二个 Chunk 失败,Job 状态为 FAILED。
修复故障后,再次使用相同识别性参数启动:
java -jar batch-app.jar input=file:/data/orders.csv
Spring Batch 会根据:
Job 名称 + input 参数
找到原来的 JobInstance,并创建新的 JobExecution。然后:
- 检查上一次 JobExecution 的状态;
- 判断 Job 是否允许重启;
- 找到失败 Step 的 StepExecution;
- 读取其 ExecutionContext;
- 让支持状态恢复的 Reader 从检查点附近继续;
- 完成剩余 Chunk;
- 将 Job 标记为
COMPLETED。
2. 重启与新运行的区别
如果想处理一个新的文件:
java -jar batch-app.jar input=file:/data/orders-2025-02.csv
这是新的 JobInstance。
如果想强制把同一个文件当成新的业务批次,可以增加一个识别性参数:
java -jar batch-app.jar \
input=file:/data/orders.csv \
batchId=20250201
但这会创建新的 JobInstance,意味着它通常会从头读取文件。它不是“从失败点继续”的重启。
因此:
相同识别性参数 + FAILED/STOPPED -> 尝试重启
不同识别性参数 -> 新建 JobInstance,通常从头开始
不要为了绕过“Job 已完成”错误,随意加 run.id,然后把它称为重启。那通常是一次新的运行,可能重新执行所有业务副作用。
3. 哪些情况不能直接重启
以下情况需要谨慎处理:
Job 已经 COMPLETED
同一 JobInstance 已完成时,Spring Batch 通常拒绝再次运行。若确实要重跑,应明确创建新的业务运行标识,并保证 Writer 幂等。
Step 被标记为不可重启
某些 Step 可以配置为不允许重启。此时失败后再次启动会被拒绝。是否允许重启必须与业务语义一致。
使用了不支持状态恢复的 Reader
如果 Reader 没有保存位置,重启可能只能从头读。此时必须依赖幂等 Writer、业务去重或人为切分输入。
外部副作用已经发生
即使 Batch 元数据显示失败,外部系统可能已经执行成功。直接重启可能造成重复发送、重复扣款或重复发货。
4. saveState 的边界
许多有状态 Reader 和 Writer 默认会保存执行状态,但也存在 saveState 一类配置可以关闭状态保存。
关闭状态保存的后果是:
当前 Step 失败后,框架无法依赖该组件的检查点恢复位置
这有时适用于无状态、可重算或每次都从头读取的场景,但不能把它当作性能优化后无代价地启用。
九、事务、回滚和故障路径
1. 数据库 Writer 的典型故障路径
以 JDBC Writer 为例:
读取 item 1、2、3
处理 item 1、2、3
开始数据库事务
批量 INSERT 1、2、3
数据库提交
保存进度
如果 INSERT item 2 违反唯一键:
INSERT 1、2、3
↓
item 2 失败
↓
事务回滚
↓
item 1、3 的本次写入也撤销
↓
Step 失败或进入 skip/retry 流程
这要求目标数据库写入确实参与了同一个 Spring 事务。若 Writer 内部另开连接、异步提交或调用不受该事务管理的数据源,就不能假定会自动回滚。
2. 一个 Chunk 中的读写并非都具有同样的回滚语义
数据库写入通常可以回滚,但 Reader 从文件中读取的文件指针不能通过数据库事务“倒退”。框架需要依赖 Reader 保存的状态,在重启或回滚路径中重新定位,或者重新读取一部分数据。
因此要区分:
数据库事务回滚:撤销已执行的数据库变更
Reader 状态恢复:重新确定下一次应读取的位置
两者共同构成 Chunk 的恢复机制,但不是同一件事。
3. Retry 和 Skip 会改变简单的回滚模型
如果配置了重试:
当前 Chunk 中某 Item 失败
↓
对该 Item 再次尝试
↓
成功后继续当前流程
如果配置了跳过:
某 Item 持续失败
↓
记录 skip
↓
跳过该 Item,继续处理其他 Item
这时“失败即整个 Chunk 永久失败”的模型不再完整。必须同时记录:
- 哪些异常可重试;
- 最大重试次数;
- 哪些异常可跳过;
- 跳过的数据如何保存;
- 最终报告是否包含失败记录。
跳过本质上是业务决策,不是简单的异常吞掉。订单金额无法解析可以进入错误表;但数据库连接中断通常不能被当成坏数据跳过。
十、幂等:重启安全的核心业务条件
1. 幂等的形式化定义
设一次业务操作为函数:
f(state, input) -> state'
若对相同输入重复执行满足:
f(f(state, input), input) = f(state, input)
则称该操作具有幂等性。
直觉是:
执行一次和执行多次,最终状态相同。
注意,幂等不等于“每次调用都返回相同响应”,而是最终业务状态不会因为重复执行而继续错误累积。
2. 非幂等示例
update account
set balance = balance + 100
where account_id = 1;
执行一次余额增加 100,执行两次增加 200:
f(f(balance, +100), +100) != f(balance, +100)
这是非幂等操作。
3. 幂等示例
update account
set status = 'PAID'
where order_id = 1001;
重复执行多次,状态仍然是 PAID。在没有其他条件副作用的前提下,这是幂等的。
但以下操作即使状态赋值相同,也可能不是幂等的:
设置订单为 PAID
同时发送一次通知
订单状态可以幂等,通知发送可能重复。幂等性必须针对完整业务副作用分析,而不是只看一条 SQL。
十一、Spring Batch 重启为什么仍然需要幂等
很多人会推导出一个错误结论:
Spring Batch 有 Checkpoint
所以重启不会重复处理
正确结论是:
Checkpoint 尽量减少重复读取范围;
幂等性保证重复读取或重复写入不会破坏业务结果。
典型故障时间线
T1: 读取订单 1001
T2: 调用外部物流系统,创建运单成功
T3: 本地程序崩溃
T4: Checkpoint 尚未提交
T5: 重启后再次读取订单 1001
T6: 再次创建运单
如果外部物流系统没有幂等键,就可能创建两张运单。
即使 Writer 是数据库 Writer,也存在类似窗口:
T1: 数据库写入成功
T2: 数据库事务已提交
T3: 应用在 Batch 元数据更新前崩溃
T4: 重启后重新处理该 Item
是否重复取决于业务表约束和写入语句:
- 普通
INSERT:可能主键冲突; INSERT ... ON CONFLICT DO UPDATE:可以转为幂等 Upsert;- 先查询再插入:在并发下仍可能竞态;
- 唯一键加冲突处理:通常比“先查后写”更可靠。
用业务唯一键实现去重
假设订单导入的业务唯一键是 order_id,可以建立唯一约束:
create unique index uk_imported_order_id
on imported_order(id);
然后根据需求选择:
方案一:重复数据视为异常
适用于重复表示数据质量问题:
insert into imported_order(id, amount)
values (:id, :amount);
重复时让 Step 失败或进入专门的 skip 处理。
方案二:重复数据更新
适用于导入是“最终状态覆盖”:
merge into imported_order t
using (values (:id, :amount)) s(id, amount)
on t.id = s.id
when matched then
update set amount = s.amount
when not matched then
insert (id, amount) values (s.id, s.amount);
具体 MERGE 语法依赖数据库,应使用目标数据库的正确方言。
方案三:重复数据忽略
适用于“只要曾经成功写入即可”:
如果业务键已经存在,则认为该 Item 已完成
这种做法必须保证已存在的记录确实代表本次业务处理成功,而不是上一次执行只写入了一半数据。
十二、批量导入中的幂等设计模式
1. 业务键加唯一约束
适用于单表导入:
业务键:order_id
数据库唯一约束:order_id
重复写入:冲突、更新或忽略
数据库约束是并发安全的重要组成部分,因为它在数据库层强制执行,而不是依赖多个应用实例的检查逻辑。
2. Staging 表加最终合并
把文件先写入临时表:
CSV
↓
staging_order
↓
校验、去重
↓
正式表
临时表可以携带:
batch_id
source_file
source_line
business_key
payload_hash
processed_at
正式合并时根据:
batch_id + business_key
判断是否已经应用。这样可以把“文件读取”和“正式业务变更”分成两个更容易恢复的阶段。
3. 幂等请求键
对于外部 HTTP 服务,为每次业务操作生成稳定的幂等键:
idempotency_key = order_id + ":" + operation_type
重启后再次调用时,携带同一个键:
POST /shipments
Idempotency-Key: order-1001:create-shipment
前提是外部服务真正支持该语义,并且会保存和复用这个键对应的结果。仅仅在客户端生成 Header,而服务端不识别,不能产生幂等效果。
4. Outbox 模式
如果一个本地数据库事务既要修改业务状态,又要发送消息,可以把消息写入 Outbox 表:
业务表更新 + Outbox 记录
↓ 同一数据库事务
事务提交
↓
独立发布器读取 Outbox
↓
发送消息
↓
标记已发送
发布器仍可能重复发送,因此消费者也应使用消息 ID 去重。Outbox 解决的是本地状态和待发送事件之间的原子记录问题,不等于整个消息链路天然 exactly-once。
十三、Checkpoint 与幂等的关系
可以用下表区分它们:
| 问题 | 主要解决者 |
|---|---|
| 进程崩溃后从哪里继续读 | Checkpoint |
| 当前 Chunk 的数据库写入是否回滚 | 事务 |
| 已经提交的数据是否重复插入 | 幂等 Writer / 唯一约束 |
| 外部 API 是否重复执行 | 外部幂等键 / 业务去重 |
| 同一批业务是否被重复启动 | JobInstance 参数设计 |
| 失败数据是否继续处理 | Retry / Skip / 错误表 |
它们之间存在依赖关系:
没有 Checkpoint:
可能从头重读,需要更强的幂等性
没有事务:
一个 Chunk 可能部分写入,恢复更复杂
没有幂等性:
即使 Checkpoint 正确,故障窗口仍可能造成重复副作用
因此,重启安全通常不是单个开关,而是以下条件的组合:
持久化 JobRepository
+ 可恢复的 Reader/Writer
+ 清晰的事务边界
+ 幂等的业务写入
+ 可识别的 Job 参数
+ 明确的异常处理策略
十四、并发执行时,Chunk 和幂等问题会扩大
1. 默认单线程执行
普通 Chunk Step 默认按顺序处理:
线程 1:读取 1、2、3 -> 写入 -> 提交
线程 1:读取 4、5、6 -> 写入 -> 提交
这时顺序和检查点相对容易理解。
2. 多线程 Step
多线程 Step 可能变成:
线程 1:读取 1、2、3
线程 2:读取 4、5、6
线程 3:读取 7、8、9
此时:
- Chunk 的提交顺序可能不同;
- 输出顺序可能不再稳定;
- Reader 必须支持并发访问,或使用合适的同步策略;
- 数据分片不能重叠;
- Writer 必须支持并发调用;
- 检查点必须能正确表示并发处理状态。
一个只能维护单一游标的 Reader,不能因为外层加了线程池就自动变成线程安全 Reader。
3. Partitioning
Partitioning 将输入空间切成多个独立分区:
master Step
├── partition 0: id 1 ~ 10000
├── partition 1: id 10001 ~ 20000
└── partition 2: id 20001 ~ 30000
分区条件必须满足:
各分区没有重叠
所有输入都被覆盖
每个分区可以独立重启或重新执行
如果两个分区都可能写入同一个业务键,那么即使每个分区内部都有 Checkpoint,整体仍可能产生重复或竞争。
4. 不要并发启动同一个 JobInstance
如果两个进程使用完全相同的 Job 名称和识别性参数同时启动:
importJob + input=orders.csv
它们竞争的是同一个 JobInstance。最终结果取决于数据库锁、Repository 隔离级别和启动时序,可能出现实例已运行、重复执行或状态竞争。
常见做法是:
- 同一业务批次只允许一个调度器启动;
- 用数据库锁或调度平台保证互斥;
- 为可并行处理的批次分配不同业务参数;
- 在目标表使用唯一约束作为最后一道保护。
十五、常见误解与实际失败表现
误解一:chunk(100) 表示最多只会重复 100 条
通常它表示一个事务块大小,但实际重放范围还取决于:
- 检查点持久化时机;
- Reader 是否保存状态;
- Writer 是否已在外部系统产生副作用;
- retry/skip 配置;
- 失败发生在提交前还是提交后。
因此不能把 100 当作所有情况下的绝对重复上限。
误解二:数据库事务可以回滚 HTTP 调用
下面的代码不会因为数据库事务回滚而撤销远端调用:
writer.write(chunk);
httpClient.post("/payment");
除非远端系统实现了兼容的分布式事务协议,并且整个系统正确配置,否则本地事务和 HTTP 调用不是一个原子操作。
误解三:加 run.id 就是重启
增加参数会改变 JobInstance 身份:
旧实例:input=a.csv
新实例:input=a.csv, run.id=2
新实例通常从头执行。它可以用于“有意重跑”,但不等于从旧实例的检查点恢复。
误解四:处理成功但 Job 失败,数据一定没有写入
可能发生以下情况:
业务表事务已提交
Batch 元数据还没来得及更新
JVM 崩溃
这时业务表中可能已有数据,但 Job 状态仍是 FAILED。重启必须依赖幂等写入,不能只看 Batch 状态判断业务表是否为空。
误解五:把所有异常都 skip
如果把数据库宕机、网络中断、程序缺陷都配置为 skip,可能得到:
Job 状态 COMPLETED
但大量业务数据没有处理
Skip 适用于明确可丢弃或可隔离的业务数据异常,不适用于基础设施故障。
十六、诊断重启问题的方法
遇到“重启后重复写入”“从头读取”“Job 无法再次启动”时,应按以下顺序核对。
1. 检查 JobInstance 参数
确认:
Job 名称是否相同
识别性参数是否相同
是否意外增加了 run.id 或时间戳
参数变化会直接改变 JobInstance。
2. 检查 JobExecution 和 StepExecution
重点查看:
status
exit_code
start_time
end_time
read_count
write_count
commit_count
rollback_count
skip_count
failure_exceptions
如果 Step 显示已提交 300 条,但业务表只有 200 条,需要检查 Writer 是否使用了同一事务和同一数据源。
3. 检查 ExecutionContext
确认 Reader 是否保存了状态,以及失败前最后一个成功提交点是什么。不要只看日志中的“读取到第几条”,因为读取计数、写入计数和提交点可能不同。
4. 检查目标表的唯一约束
如果没有业务唯一约束,重复执行可能静默产生重复行。诊断时应查询:
select business_key, count(*)
from target_table
group by business_key
having count(*) > 1;
5. 检查失败发生在哪个时间窗口
需要区分:
Writer 执行前失败
Writer 执行中失败
数据库事务提交前失败
数据库事务提交后、Batch 状态更新前失败
外部调用成功后本地崩溃
不同时间点对应不同恢复策略,不能只依据异常消息判断。
十七、如何验证一个 Step 真正支持重启
一个完整的重启测试不应只测试“正常执行成功”,而应人为制造失败:
- 准备 10 条输入;
- 配置
chunk(3, ...); - 在处理到第 5 条时抛出一次性异常;
- 运行 Job,确认状态为
FAILED; - 检查数据库只有第一个 Chunk 的结果;
- 移除故障;
- 使用相同 Job 参数重新启动;
- 确认第 4 条附近被重新处理;
- 确认最终数据共 10 条且没有重复;
- 检查 Step 的最终
read_count、write_count和提交次数。
还应测试以下情况:
- 失败发生在 Writer 中;
- 数据库连接在提交附近中断;
- 同一文件被重复作为新 JobInstance 执行;
- 两个进程同时启动同一 Job 参数;
- 输入中存在坏数据;
- 外部调用成功但本地进程随后崩溃。
只有这些故障窗口都能解释,才能认为设计真正具备可恢复性。
十八、设计时的核心判断
对于一个批处理业务,可以按下面的因果关系设计:
如果输入数据可定位
-> 让 Reader 保存检查点
如果目标写入在数据库事务中
-> 让数据库事务覆盖当前 Chunk
如果可能发生提交后崩溃
-> 让 Writer 依据业务键幂等
如果调用外部系统
-> 使用外部幂等键、状态查询或 Outbox
如果坏数据和基础设施故障不同
-> 只对明确的业务异常 skip
如果需要重跑
-> 使用新的业务批次参数,而不是伪装成重启
如果需要并发
-> 先证明分区不重叠、组件可并发、结果可幂等
Spring Batch 的 Job、Step 和 Chunk 解决了批处理的结构化执行;JobRepository 和 Checkpoint 解决了进度持久化;事务解决了当前处理块的原子性;重启机制解决了失败后的继续执行。但最终业务是否安全,还取决于目标系统是否能够承受重复执行。
最可靠的理解方式是:
Checkpoint 让系统知道“应当从哪里继续”;
事务让系统能够撤销“尚未提交的这一块”;
幂等让系统能够安全承受“不可避免的重复”。
三者缺一时,长时间运行的批处理就很难在真实故障中保持可恢复和可解释。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Spring 调度任务:线程池、Cron、时区、防重入和分布式锁
- 下一篇:Spring Integration:Channel、Adapter、Gateway、流程和错误处理
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论