批处理“失联”又“重跑”:Spring Batch 监控缺失与重启暴雷,你的作业还安全吗?
批处理“失联”又“重跑”:Spring Batch 监控缺失与重启暴雷,你的作业还安全吗?
凌晨两点,ETL 作业跑了 6 个小时,你不知道它卡在哪一步,也不敢重启,怕数据重复;唯一知道的是数据库连接池快满了,CPU 却在摸鱼。终于,作业因为一条“脏数据”崩了,你想从失败点重启,却发现它根本不是从断点继续,而是把所有数据又处理了一遍,主键冲突炸穿日志。更绝望的是,Spring Boot Actuator 只告诉你“作业正在运行”,其他一概不知。这不是批处理的宿命,而是你缺少一套透明的监控体系和精准的重启机制。
本文深入 Spring Boot + Spring Batch 的批处理监控与重启疑难杂症,从元数据表解读、进度暴露、失败恢复策略到优雅重启,给你一套让批处理从“盲盒”变“透明”的完整方案,确保每一次重启都精准落在断裂的齿扣上。
一、血泪现场:监控失明与重启陷阱
1.1 执行进度成谜,排查全靠猜
一个数据迁移作业跑了 5 小时,日志一直在滚动,但没人知道到底迁移了多少条,还有多少没迁。运维只能在数据库里 SELECT COUNT(*) 对比,但数据在变化,结果不准。作业最终失败,大量时间被浪费。
1.2 失败后重启,数据重复写入
你点击重启,Spring Batch 确实“从失败步骤继续”,但你的 ItemWriter 没做幂等,之前已成功写入的 10 万条数据被再次插入,主键冲突导致大量异常,或者重复记录污染业务表。
1.3 只想跳过坏记录,结果整个作业都被跳过了
你配置了 skip-limit=10,以为只跳过 10 条坏数据。但作业重启后,框架重新处理了所有记录,导致之前被跳过的错误又被触发,加上新的错误,跳过计数累计,直接超过上限,作业再次失败。
1.4 没人知道作业失败了
作业异常终止,没有告警。直到第二天业务方询问报表为何没出,你才发现数据库里的 BATCH_JOB_EXECUTION 状态是 FAILED,但无人知晓。
这些痛苦源于两个缺失:运行时的可观测性 和 失败后的精确恢复策略。Spring Batch 其实内置了强大的元数据持久化和重启能力,但你需要正确配置和使用。
二、Spring Batch 的监控根基:四大元数据表
Spring Batch 将作业和步骤的执行状态持久化到数据库中(默认表前缀 BATCH_),这是监控和重启的基础。
BATCH_JOB_INSTANCE:作业实例,由jobName和唯一jobKey(识别参数)决定。BATCH_JOB_EXECUTION:每次作业执行记录,包含START_TIME,END_TIME,STATUS,EXIT_CODE。BATCH_STEP_EXECUTION:步骤执行记录,包含READ_COUNT,WRITE_COUNT,FILTER_COUNT,COMMIT_COUNT,STATUS等。BATCH_JOB_EXECUTION_PARAMS:执行参数。
监控的起点就是查询这些表。Spring Batch 提供了 JobExplorer API 来访问它们,并可通过 Actuator 暴露端点。
三、方案一:利用 Spring Boot Actuator 快速暴露作业状态
Spring Batch 集成自动配置了 JobLauncherCommandLineRunner,并且 Actuator 可以提供基本的作业信息。
开启批处理端点(Spring Boot 2.x 和 3.x 略有不同):
management:
endpoints:
web:
exposure:
include: health,info,metrics,scheduledtasks,batch
endpoint:
batch:
enabled: true
访问 /actuator/batch 会返回最近执行的作业概要。但信息有限,还需要自定义增强。
3.1 自定义 Actuator 端点暴露详细进度
通过 JobExplorer 获取正在运行的步骤的读写计数,计算进度。
@Component
@Endpoint(id = "batch-progress")
public class BatchProgressEndpoint {
@Autowired
private JobExplorer jobExplorer;
@ReadOperation
public Map<String, Object> getProgress(@Selector String jobName) {
// 查找最近一个正在运行的作业实例
JobInstance jobInstance = jobExplorer.getLastJobInstance(jobName);
if (jobInstance == null) return Map.of("status", "NOT_FOUND");
List<JobExecution> executions = jobExplorer.getJobExecutions(jobInstance);
JobExecution latest = executions.get(executions.size() - 1);
// 收集步骤执行信息
List<Map<String, Object>> steps = new ArrayList<>();
for (StepExecution se : latest.getStepExecutions()) {
steps.add(Map.of(
"stepName", se.getStepName(),
"status", se.getStatus().toString(),
"readCount", se.getReadCount(),
"writeCount", se.getWriteCount(),
"commitCount", se.getCommitCount()
));
}
return Map.of(
"jobId", latest.getJobId(),
"status", latest.getStatus().toString(),
"startTime", latest.getStartTime(),
"steps", steps
);
}
}
现在通过 /actuator/batch-progress/jobName 就可以实时看到进度。
3.2 Micrometer 指标推送 Prometheus
Spring Batch 自动注册 Micrometer 指标:spring_batch_job_duration_seconds, spring_batch_job_status, spring_batch_step_read_count 等。只需添加 micrometer-registry-prometheus,配置 Actuator 暴露 Prometheus 端点,即可在 Grafana 中可视化。
自定义 Gauge 追踪进度百分比:
@Bean
public MeterBinder batchProgressMetrics(JobExplorer jobExplorer) {
return registry -> {
Gauge.builder("batch.job.progress", () -> {
// 根据预估总量和已读数量计算百分比,需要业务传入总量
// 这里简化为读取计数
return getCurrentReadCount();
}).register(registry);
};
}
四、方案二:JobExecutionListener 与 StepExecutionListener 记录详细日志与告警
4.1 全局监听器记录生命周期
@Component
public class BatchJobNotificationListener implements JobExecutionListener, StepExecutionListener {
@Autowired
private NotificationService notificationService;
@Override
public void beforeJob(JobExecution jobExecution) {
notificationService.send("Job " + jobExecution.getJobInstance().getJobName() + " started");
}
@Override
public void afterJob(JobExecution jobExecution) {
if (jobExecution.getStatus() == BatchStatus.FAILED) {
String msg = "Job FAILED: " + jobExecution.getExitStatus().getExitDescription();
notificationService.sendAlert(msg);
}
}
@Override
public ExitStatus afterStep(StepExecution stepExecution) {
log.info("Step {}: read={}, write={}, commit={}, status={}",
stepExecution.getStepName(), stepExecution.getReadCount(),
stepExecution.getWriteCount(), stepExecution.getCommitCount(),
stepExecution.getStatus());
return stepExecution.getExitStatus();
}
}
4.2 长时间运行任务告警
在 beforeStep 中记录时间戳,在定时检查线程或通过 StepExecution 的 getStartTime 计算持续时间,当超过阈值时告警。可与外部调度平台集成。
五、方案三:精准重启——让作业从断点继续,而非重头再来
Spring Batch 的“从失败处重启”是建立在执行上下文持久化和幂等读写之上的。
5.1 利用 ExecutionContext 保存断点位置
在 ItemReader 中,将当前读取位置(如分页偏移、文件行号)保存在 StepExecution 的 ExecutionContext 中,当步骤失败后,下次重启时读取器会从这个位置继续。
public class RetryableItemReader implements ItemReader<String> {
private int current = 0;
private int total = 1000;
@BeforeStep
public void restoreState(StepExecution stepExecution) {
if (stepExecution.getExecutionContext().containsKey("current")) {
this.current = stepExecution.getExecutionContext().getInt("current");
}
}
@Override
public String read() {
if (current >= total) return null;
String data = fetchData(current);
// 在合适的时机更新上下文(如每 chunk 提交时)
ExecutionContext ctx = StepSynchronizationManager.getContext().getStepExecution().getExecutionContext();
ctx.putInt("current", current + 1);
current++;
return data;
}
}
当步骤在某个 chunk 失败时,Spring Batch 会回滚事务,并且 ExecutionContext 中的 current 会保留在失败前的值(因为事务回滚了未提交的上下文更新?实际上 ExecutionContext 的更新是在 chunk 提交后持久化,所以失败 chunk 的更新不会保存,下次重启会从上一个成功提交的位置继续)。这个机制天然保证了重启只处理未完成的部分。
5.2 非幂等 Writer 的重启保护
如果写操作不是幂等的(如发送邮件、调用非幂等 API),重启会导致重复操作。解决方式:
- 将写操作设计为事务性且幂等(如 SQL INSERT IGNORE 或 Merge)。
- 对于外部服务,在执行前检查
StepExecution的writeCount和commitCount,或者利用ItemWriter的beforeStep记录已处理条数,但更可靠的是使用“已处理日志表”。 - 使用
JobParameters的唯一标识(如run.id)配合业务唯一键避免重复。
5.3 跳过已处理记录与重试
Spring Batch 提供了 skip-limit 和 retry-limit。但重启时,skip 计数会重置(除非使用 ExecutionContext 记录)。如果希望重启后继续保留跳过配额,需要在 SkipPolicy 中自定义,结合数据库记录已跳过的记录标识。
更好的做法是:对于可预见的脏数据,用 ItemProcessor 返回 null 来过滤,而非依赖 skip。
5.4 重启失败的作业
通过 JobOperator 可以命令行或 API 重启:
@Autowired
private JobOperator jobOperator;
public void restartLastFailedJob(String jobName) {
JobInstance jobInstance = jobExplorer.getLastJobInstance(jobName);
List<JobExecution> executions = jobExplorer.getJobExecutions(jobInstance);
JobExecution lastExecution = executions.get(executions.size() - 1);
if (lastExecution.getStatus() == BatchStatus.FAILED) {
jobOperator.restart(lastExecution.getId());
}
}
或者使用 Spring Batch Admin 等工具。
注意事项:重启相同的 JobParameters 会直接返回原来的 JobInstance,并继续未完成的步骤。如果参数有变化,会创建新实例,无法从旧实例重启。因此,需要保持作业识别参数(如日期)不变,而通过 run.id 等非识别参数区分每次执行。
六、方案四:健康检查与元数据表维护
6.1 元数据表膨胀问题
长时间运行后,元数据表会积累大量历史数据,影响查询性能。Spring Batch 提供 JobRepository.deleteJobExecution() 等方法,可以编写定时任务清理旧记录。
@Scheduled(cron = "0 0 4 * * ?")
public void cleanUpHistory() {
for (String jobName : jobNames) {
List<JobInstance> instances = jobExplorer.findJobInstancesByJobName(jobName, 0, 100);
// 保留最近10个实例,删除更旧的
}
}
注意:删除作业实例会级联删除执行记录和上下文。谨慎操作,确保不再需要这些审计数据。
6.2 监控元数据表本身
查询 BATCH_JOB_EXECUTION 中的 STATUS=STARTED 长时间未更新的记录,可能是僵尸作业(进程被杀但状态未更新)。检测并报警,必要时手动标记为 FAILED 或 ABANDONED。
6.3 与外部调度平台集成
Spring Batch 的作业可以由 Spring Cloud Data Flow、Kubernetes CronJob 或 XXL-JOB 触发,这些平台可以提供额外的监控面板和重启接口,但作业自身的元数据仍是最底层保障。
七、常见疑难杂症速查表
| 现象 | 根因 | 解决方案 |
|---|---|---|
| 重启后从头开始,而非从断点 | ItemReader 未持久化状态到 ExecutionContext | 实现 ItemStream 或使用 @BeforeStep/@AfterStep 管理状态 |
| 重启导致数据重复 | Writer 非幂等,或者状态保存与实际写入不一致 | 设计幂等写,或使用唯一键去重;确保状态更新在事务内 |
| 跳过记录在重启后再次导致失败 | skip count 重置,跳过策略未持久化 | 记录已跳过的记录 ID,重启时通过 Reader 过滤 |
| 作业执行状态为 UNKNOWN | 进程被杀,元数据未更新 | 启动时检测并标记 FAILED,配合 JobExplorer 修复 |
| Actuator 端点暴露的信息过少 | 默认仅返回最后几次执行概要 | 自定义 Endpoint,或使用 Spring Batch Admin UI |
| 元数据表查询慢 | 数据过多未清理 | 定期清理历史,添加索引 |
| 并行作业的监控混乱 | 多个 JobExecution 同时运行 | 按 JobInstance 和 JobExecution ID 区分,Prometheus 标签使用 job_name 和 status |
| 重启时丢失作业参数 | 未将识别参数与非识别参数分开 | 使用 JobParametersBuilder 时,非识别参数不参与实例标识 |
八、最佳实践:打造可观测、可恢复的批处理
- 必须启用数据库元数据持久化:
spring.batch.jdbc.initialize-schema=always或手动建表,确保重启能力。 - 每个 ItemReader 实现状态保存:使用
ItemStreamSupport或ExecutionContext,让作业可恢复。 - 幂等 Writer 是黄金法则:如果不能,就用“已处理记录表”进行去重。
- 全局监听器统一告警:失败、完成、超时都要通知到人。
- 利用 Micrometer 暴露进度指标:读计数、写计数、作业耗时,让 Grafana 大屏可见。
- 合理分割 JobParameters:业务日期等识别参数保持不变,
run.id等非识别参数每次递增。 - 定期归档元数据:保留适当历史,避免性能退化。
- 编写集成测试模拟失败重启:验证状态恢复和幂等性,防止上线后翻车。
- 使用
JobOperator提供管理 API:允许手动停止、重启,而不是直接操作数据库。 - 对长时间运行的任务设置超时:在
Step中配置throttle-limit或通过TaskExecutor超时控制。
九、结语:让批处理从“黑箱”变“白盒”
Spring Batch 为批处理提供了强大的持久化和重启框架,但监控和恢复策略需要你亲手设计。当你把进度暴露到监控大屏,把断点保存到执行上下文,把幂等写入每条记录,批处理就不再是凌晨的噩梦。现在,检查你的作业:ItemReader 有没有保存状态?元数据表还能查吗?失败后重启会重复写入吗?补上这些空缺,你就能安心入睡,因为你的批处理知道在哪里跌倒,就在哪里爬起。
更多推荐

所有评论(0)