批处理“失联”又“重跑”: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 中记录时间戳,在定时检查线程或通过 StepExecutiongetStartTime 计算持续时间,当超过阈值时告警。可与外部调度平台集成。


五、方案三:精准重启——让作业从断点继续,而非重头再来

Spring Batch 的“从失败处重启”是建立在执行上下文持久化幂等读写之上的。

5.1 利用 ExecutionContext 保存断点位置

ItemReader 中,将当前读取位置(如分页偏移、文件行号)保存在 StepExecutionExecutionContext 中,当步骤失败后,下次重启时读取器会从这个位置继续。

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)。
  • 对于外部服务,在执行前检查 StepExecutionwriteCountcommitCount,或者利用 ItemWriterbeforeStep 记录已处理条数,但更可靠的是使用“已处理日志表”。
  • 使用 JobParameters 的唯一标识(如 run.id)配合业务唯一键避免重复。

5.3 跳过已处理记录与重试

Spring Batch 提供了 skip-limitretry-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 长时间未更新的记录,可能是僵尸作业(进程被杀但状态未更新)。检测并报警,必要时手动标记为 FAILEDABANDONED

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 同时运行 JobInstanceJobExecution ID 区分,Prometheus 标签使用 job_namestatus
重启时丢失作业参数 未将识别参数与非识别参数分开 使用 JobParametersBuilder 时,非识别参数不参与实例标识

八、最佳实践:打造可观测、可恢复的批处理

  1. 必须启用数据库元数据持久化spring.batch.jdbc.initialize-schema=always 或手动建表,确保重启能力。
  2. 每个 ItemReader 实现状态保存:使用 ItemStreamSupportExecutionContext,让作业可恢复。
  3. 幂等 Writer 是黄金法则:如果不能,就用“已处理记录表”进行去重。
  4. 全局监听器统一告警:失败、完成、超时都要通知到人。
  5. 利用 Micrometer 暴露进度指标:读计数、写计数、作业耗时,让 Grafana 大屏可见。
  6. 合理分割 JobParameters:业务日期等识别参数保持不变,run.id 等非识别参数每次递增。
  7. 定期归档元数据:保留适当历史,避免性能退化。
  8. 编写集成测试模拟失败重启:验证状态恢复和幂等性,防止上线后翻车。
  9. 使用 JobOperator 提供管理 API:允许手动停止、重启,而不是直接操作数据库。
  10. 对长时间运行的任务设置超时:在 Step 中配置 throttle-limit 或通过 TaskExecutor 超时控制。

九、结语:让批处理从“黑箱”变“白盒”

Spring Batch 为批处理提供了强大的持久化和重启框架,但监控和恢复策略需要你亲手设计。当你把进度暴露到监控大屏,把断点保存到执行上下文,把幂等写入每条记录,批处理就不再是凌晨的噩梦。现在,检查你的作业:ItemReader 有没有保存状态?元数据表还能查吗?失败后重启会重复写入吗?补上这些空缺,你就能安心入睡,因为你的批处理知道在哪里跌倒,就在哪里爬起。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐