长任务状态、取消与可观测:执行跨过重启后怎样继续
每批提交 20 条,第 63 条转换失败,前三批的 60 条结果仍保存在数据库中。下一次启动要根据已提交位置找到剩余数据,并核对期间已经发生的外部操作。部署、连接中断和主动停止都会遇到类似问题。
Spring Batch 用作业实例保存批次身份,用执行记录和检查点支持续跑。调度平台负责按时启动,批处理程序负责分块提交和统计进度。平台再次调用时,传入的作业参数决定是继续原批次,还是创建新的批次。
长任务由哪些持久对象组成
Job、JobInstance 和 JobExecution
一个报表作业可以每天运行,也可以对某个历史窗口重新执行。Spring Batch 将定义、业务批次和运行尝试分开保存:
Job:copyWindow
├── JobInstance:window=window-A
│ ├── JobExecution 1:FAILED
│ │ └── StepExecution:copyRows,提交 60 条
│ └── JobExecution 2:COMPLETED
│ └── StepExecution:copyRows,继续提交 40 条
└── JobInstance:window=window-B
└── 独立的一批数据,拥有自己的执行记录Job 描述步骤及其顺序。JobInstance 由作业名称和参与身份计算的参数共同确定。JobExecution 记录一次运行,失败后以相同身份参数重新启动,会增加执行记录,并继续使用原实例的恢复信息。相关对象的完整定义见 Spring Batch 领域模型。
例如:
var parameters = new JobParametersBuilder()
.addString("window", "window-A", true)
.toJobParameters();第三个参数 true 表示 window 参与实例身份。故障注入开关、日志级别或执行速度通常不应改变这一身份。每次启动都追加随机数、当前时间或自增 run.id,会创建新实例;需要恢复旧批次时,这种写法会绕过旧检查点。
实例身份还要与读取条件一致。参数传入 window-A,Reader 却执行“不加窗口条件的全表扫描”,依然会读取全表。框架负责区分实例,数据范围由查询和输入文件决定。
步骤状态与业务结果
BatchStatus 表示框架的运行阶段;ExitStatus 保存退出码及附加说明,能够参与步骤之间的流程判断。业务结果还可能包括导出文件的位置、失败记录数量或外部结算批号,应放在专门的结果表中。
| 状态 | 含义 | 后续处理 |
|---|---|---|
STARTING / STARTED | 正在启动或执行 | 查询进度、运行进程与下游等待 |
STOPPING | 已请求停止 | 等待执行线程在检查点观察请求 |
STOPPED | 已停止,可按配置恢复 | 保留相同身份参数继续 |
FAILED | 运行失败 | 先定位异常,再决定重启或修复数据 |
COMPLETED | 本次实例已完成 | 相同实例再次启动通常被拒绝 |
ABANDONED | 明确放弃这次执行的后续恢复 | 需要受控管理操作,不能当作普通重试 |
UNKNOWN | 框架无法确认状态 | 核对事务与实际结果后再处理 |
业务允许“成功完成并跳过少量坏记录”时,COMPLETED 仍可能伴随跳过计数。报表页面应同时展示处理结果与异常记录,不能只显示一个绿色状态。
多步骤任务还要区分步骤级与作业级恢复。默认情况下,重启时已经完成的步骤会被跳过;allowStartIfComplete(true) 可要求该步骤再次执行,startLimit 则限制步骤的启动次数。启用这些选项前,需要确认重复运行前置校验、清理目录或结果发布步骤是否安全。步骤重启规则
JDBC 仓库保存什么
JobRepository 持久保存实例、运行状态、步骤计数和 ExecutionContext。它的数据库表与业务表承担不同职责:
BATCH_JOB_INSTANCE 作业名称与实例身份
BATCH_JOB_EXECUTION 每次运行的状态
BATCH_JOB_EXECUTION_PARAMS 参数与 identifying 标志
BATCH_STEP_EXECUTION 读写、提交、回滚等计数
BATCH_*_EXECUTION_CONTEXT 恢复需要的上下文
lab_source 输入数据
lab_output 已提交的业务结果这些表具有外键和版本控制字段。生产升级应使用对应版本的迁移脚本,并与旧数据兼容性一同验证;直接修改状态列或删除某张表,可能破坏实例之间的关系。元数据表结构
Spring Batch 6 默认的 @EnableBatchProcessing 和 DefaultBatchConfiguration 提供 ResourcelessJobRepository。它不保存跨进程恢复所需的元数据。持久批处理应显式启用 @EnableJdbcJobRepository,或配置 JdbcJobRepositoryFactoryBean 并指定数据源和事务管理器。JobRepository 配置
检查点怎样与业务提交保持一致
一个 chunk 的提交范围
chunk 表示一次事务处理的一组数据。以每批 20 条为例,运行过程是读取一批数据、逐条转换、批量写入,然后更新恢复位置并提交事务。某条转换失败后,当前事务回滚,先前已经提交的批次继续保留。
第 1 批: 1—20 写入 + 检查点 20 COMMIT
第 2 批: 21—40 写入 + 检查点 40 COMMIT
第 3 批: 41—60 写入 + 检查点 60 COMMIT
第 4 批: 61—80 第 63 条转换失败 ROLLBACK
重新启动:读取已保存位置,从第 61 条继续chunk 越小,提交事务和更新元数据越频繁;调大后,这些开销得到分摊,但内存占用、锁持有时间和失败重做范围也会增加。单条处理较慢时,大批次还会延长停止等待。批量数应根据实际处理成本与数据库负载调整。分块处理机制
数据库里的 read_count 可以大于 write_count。Reader 可能已把整批输入交给处理流程,Processor 随后失败;读到第 80 条并不意味着前 80 条全部提交。恢复位置以持久化的 ExecutionContext 为准。
同一事务里的两次写入
如果先保存检查点,再提交业务结果,业务事务失败后恢复可能跳过尚未落库的数据。反过来,业务结果已提交而检查点没有提交,重新启动会再次处理这一批。
示例让 JDBC 仓库与 JdbcTemplate 使用同一个 DataSource 和 DataSourceTransactionManager。Spring Batch 在 chunk 的事务中写业务数据并更新上下文,提交时一同生效:
var transactionManager = new DataSourceTransactionManager(dataSource);
var factory = new JdbcJobRepositoryFactoryBean();
factory.setDataSource(dataSource);
factory.setTransactionManager(transactionManager);
factory.afterPropertiesSet();
JobRepository repository = factory.getObject();
var step = new StepBuilder("copyRows", repository)
.<Integer, Integer>chunk(20)
.transactionManager(transactionManager)
.reader(reader)
.processor(processor)
.writer(writer)
.build();这段配置显式选择业务数据库的事务管理器;chunk builder 的默认值不会自动建立这笔业务事务。如果仓库和业务数据分属两个数据库,两次提交仍相互独立,需要通过幂等写入、可重复读取等方式处理提交间隙。
Reader 的恢复位置必须有稳定含义
数据库游标 Reader 可以保存已读取位置;分页 Reader 通常依据排序键恢复。无论使用哪一种方式,都需要稳定的输入范围和确定的排序。数据库 Reader
SELECT id
FROM lab_source
ORDER BY id;实验中的 lab_source 在运行期间保持不变,因此读到第 60 条后可继续读取第 61 条。实际任务的输入表仍在变化时,应先确定窗口或快照,例如固定 id <= upper_bound,再按唯一键向后读取。
OFFSET 100000 不是稳定的业务检查点。前面的行被删除或新增后,同一个偏移量可能指向不同数据;大偏移还会增加扫描代价。常见替代方案是保存 last_id,使用 WHERE id > :last_id AND id <= :upper_bound ORDER BY id。复合排序需要保存完整排序键,不能只保留其中一列。
文件处理还需要记录文件内容身份与位置。同名文件被覆盖后,应拒绝继续旧检查点,或创建新的业务批次。恢复前比对大小和校验值,比只检查文件名可靠。
外部副作用与结果发布
邮件、支付接口和对象存储通常不参加上述数据库事务。Processor 或 Writer 调用远程服务成功后,本地事务仍可能回滚。重启时再次调用,需要远端幂等键或业务结果查询。
文件导出可以把“生成分块文件”和“发布最终结果”分成两个步骤:分块文件使用确定的对象名和校验值,清单记录已完成块;全部块齐备后,再原子更新下载入口或发布清单。失败时删除本批次未引用的临时块,不应按宽泛目录清理其他运行的文件。
检查点负责定位剩余工作,重复副作用由任务幂等与补偿处理。两者应使用同一个业务窗口标识,便于恢复时查询已完成结果。
用 Spring Batch 验证跨进程续跑
环境、身份与源码入口
下载完整实验工程。环境为 Linux、Bash、Docker Engine、Compose v2、unzip,宿主账号需要有 Docker 使用权限。实验使用独立 PostgreSQL 数据卷,数据库端口仅绑定 127.0.0.1:18225。
| 组件 | 固定版本或设置 |
|---|---|
| Spring Batch | 6.0.5 |
| 构建环境 | maven:3.9.12-eclipse-temurin-17 |
| Java 运行环境 | eclipse-temurin:17.0.20_8-jdk |
| PostgreSQL | 18.6 |
| JDBC 驱动 | 42.7.8 |
| 应用容器身份 | 10001:10001 |
| 业务输入 | 100 条有序记录,每批 20 条 |
Maven 容器使用宿主 UID/GID,缓存目录由当前用户创建。数据库镜像按自身初始化流程切换数据库进程身份。示例中的固定口令仅供隔离实验,不能用于共享或生产环境。
在 ZIP 所在目录解压并构建:
unzip longjob-state-observability-lab.zip
cd longjob-state-observability-lab
mkdir -p .m2
docker run --rm --user "$(id -u):$(id -g)" \
-e MAVEN_CONFIG=/tmp/.m2 \
--mount "type=bind,source=$PWD,target=/work" \
--mount "type=bind,source=$PWD/.m2,target=/m2" \
-w /work maven:3.9.12-eclipse-temurin-17 \
mvn -B -ntp -Duser.home=/tmp -Dmaven.repo.local=/m2 clean verify
docker compose -p batch20 up -d --wait postgres
docker compose -p batch20 run --rm lab init初始化成功输出 schema=ready sourceRows=100。若提示 Refusing to initialize a nonempty public schema,说明该卷已有表。先检查是否使用了上次实验数据,不要直接清空数据库。clean verify 负责构建;以下独立 JVM 的运行结果才验证数据库行为。
在第 63 条制造失败
set +e
docker compose -p batch20 run --rm -e FAIL_AT=63 lab run
status=$?
set -e
test "$status" -eq 2
docker compose -p batch20 run --rm lab inspectFAIL_AT 只影响故障注入,不参与实例身份。进程结束时输出:
instance=1 execution=1 status=FAILED
windowRows=60第一次检查得到:
status=FAILED read_count=80 write_count=60 commit_count=3 rollback_count=1
outputRows=60
persistedContexts=1这里 read_count=80 来自第四批已经读取的输入,业务表只保存前三批共 60 条。rollback_count=1 对应失败批次的回滚。主键 (window_key,id) 保证同一窗口不能重复插入相同记录,后续如果错误地从头读取,会出现唯一键冲突。
由新 JVM 继续原实例
下面的命令创建新容器和新 JVM,没有继承上一个进程的 Java 对象:
docker compose -p batch20 run --rm lab run
docker compose -p batch20 run --rm lab inspect默认仍使用 window=window-A,输出:
instance=1 execution=2 status=COMPLETED
windowRows=100第二次执行的步骤计数为 read_count=40、write_count=40、commit_count=2、rollback_count=0。它处理剩余的 40 条,实例编号保持不变,执行编号增加。编号本身由数据库分配,在已有其他实验记录时可以不同。
再用同一身份启动:
set +e
docker compose -p batch20 run --rm lab run
status=$?
set -e
test "$status" -eq 3预期收到 JobInstanceAlreadyCompleteException,已完成实例被拒绝重复启动。需要计算新窗口时,显式更换 WINDOW_KEY;需要重做旧窗口时,应先确定结果覆盖、历史版本和审计策略,不要仅添加随机参数绕过限制。
从另一个进程请求停止
在终端 A 启动另一个实例,并把单条处理延迟设为 200 毫秒:
docker compose -p batch20 run --rm \
-e WINDOW_KEY=window-stop -e DELAY_MS=200 lab run终端 B 在终端 A 开始运行后执行:
docker compose -p batch20 run --rm lab stop
docker compose -p batch20 run --rm lab inspectstop 通过 JDBC 仓库找到运行中的执行记录,再调用真正的 JobOperator.stop。stopAccepted=true 表示请求被接受。查询时可能看到 STOPPING,也可能因为当前批次已经结束而直接看到 STOPPED。
当前版本的 chunk 处理在批次边界检查停止请求。实验中停止后曾保留 20 条或 40 条结果;其他时序也可能停在 60 或 80 条,首批提交前停止则可能保留 0 条。判断依据是执行已进入 STOPPED,已提交结果位于 0 到 100 之间且为 20 的倍数,不要求至少提交一批。如果已经完成 100 条,说明停止请求发得太晚,需要用新窗口再次观察。
以相同窗口恢复:
docker compose -p batch20 run --rm -e WINDOW_KEY=window-stop lab run该窗口最后达到 windowRows=100。停止流程只结束后续处理,没有撤销此前提交的结果。
实验结束保留数据:
docker compose -p batch20 down确认不再需要本实验元数据及输出后,才执行 docker compose -p batch20 down -v;它会删除这个 Compose 项目声明的数据卷,删除后不能依靠检查点恢复。
停止、故障与进度怎样处理
协作式停止与强制终止
Thread.interrupt、Future.cancel(true) 和框架停止操作的具体效果,取决于当前代码是否检查信号、阻塞调用是否支持中断。数据库网络调用、第三方 SDK 和本地计算可能有不同的退出速度。
停止流程通常包括停止领取新批次、完成或回滚当前事务、保存状态、关闭连接和文件。应设置最长批次时间、远程调用超时以及容器停止宽限期。宽限期过短,执行者还未结束事务就收到强制终止,恢复流程必须处理残留的运行状态。
SIGKILL 后,进程无法执行 Java 的清理代码。已提交的 chunk 仍然存在,未提交的数据库事务由连接终止触发回滚,但 Batch 元数据可能保留 STARTED。此时不能只改一列为 FAILED 后重启。先确认旧进程已经退出、数据库会话不再执行、外部副作用已核对,再使用该版本支持的管理恢复操作处理执行记录。
ABANDONED 应用于明确放弃恢复的执行,不能把它当作“强制解锁”。已经废弃的步骤或输出可能需要人工清理,再以新的业务版本启动。
进度应统计已提交工作
长任务页面至少要区分已扫描、已提交、已跳过、已失败和剩余工作量。若输入范围固定,committed / total 可以作为稳定进度;输入持续增加时,应显示累计处理量和积压量,而不是伪造一个不断回退的百分比。
日志可以记录 jobName、业务窗口、实例 ID、执行 ID、步骤名称和最后提交位置。输入内容、邮箱地址或完整报表数据无需写进日志。一次重启可能跨越多个进程,使用实例 ID 查询能保留连续历史。
指标适合回答吞吐和积压问题:
| 观察值 | 用途 |
|---|---|
| 最近一次成功提交的时间 | 判断任务是否持续推进 |
| 每批耗时与提交速率 | 区分处理慢与调度等待 |
| 已提交记录数、跳过数 | 对照业务结果和异常处理 |
| 重试次数、回滚次数 | 观察坏数据或下游波动 |
| 输入积压量与预计完成时间 | 判断能否赶上业务截止时间 |
| 运行中的执行数 | 检查并发上限与资源占用 |
处理线程被数据库锁卡住时,心跳线程可能照常发送消息。把心跳与最后提交时间放在一起观察,才能看出任务是否还在推进。心跳缺失也需要查原因:网络断开时,业务线程可能仍然运行。
Spring Batch 的运行计数和监控扩展应按所用版本配置;框架指标、业务指标和执行日志各自保留适合的粒度。执行 ID 等高基数值宜放日志或追踪属性,不应直接作为所有指标的标签。Spring Batch 监控入口
常见异常的排查顺序
| 现象 | 优先检查 | 下一步 |
|---|---|---|
| 重启从头读取 | 实例参数、Reader 名称、saveState、持久仓库 | 对比前后实例 ID 与上下文 |
| 重启出现业务主键冲突 | 检查点是否与业务提交分离、输入是否变化 | 核对已写记录,不直接忽略全部冲突 |
一直显示 STARTED | 实际进程、数据库会话、执行心跳 | 确认旧执行已停止再做管理恢复 |
停止请求长期处于 STOPPING | 单批耗时、阻塞调用和中断处理 | 获取线程栈,定位锁或网络等待 |
| 计数增长但结果不增加 | Processor 过滤、跳过配置、事务回滚 | 对照 write、filter、skip 与回滚计数 |
| 显示完成但文件不可下载 | 最终发布步骤、对象存储权限、清单更新 | 将结果发布作为可恢复步骤检查 |
批量大小、并发线程数和数据库连接池要一起考虑。提高并发后,数据库写入吞吐未必增加,锁冲突和 I/O 争用可能先上升。先测量单批提交耗时及等待原因,再调整分区或并行度。
权威资料与规范地址
任务模型与执行元数据
- Spring Batch 领域模型:https://docs.spring.io/spring-batch/reference/domain.html
- 元数据表结构:https://docs.spring.io/spring-batch/reference/schema-appendix.html
- JobRepository 配置:https://docs.spring.io/spring-batch/reference/job/configuring-repository.html
分块读写、重启与监控
- 分块处理机制:https://docs.spring.io/spring-batch/reference/step/chunk-oriented-processing.html
- 数据库 Reader:https://docs.spring.io/spring-batch/reference/readers-and-writers/database.html
- 步骤重启规则:https://docs.spring.io/spring-batch/reference/step/chunk-oriented-processing/restart.html
- Spring Batch 监控入口:https://docs.spring.io/spring-batch/reference/spring-batch-observability.html
