每天 300 万条的对账
财务对账原来用 crontab + 一个 main 方法跑,失败了从头来,半夜崩了没人管。最惨的一次,任务跑到一半数据库连接断了,已处理的 150 万条没记录,重跑又重复处理了一遍,对账差了几十万。换成 Spring Batch 后,断点续跑和重试才真正可控,凌晨的 job 终于敢交给它,财务也敢信它的结果。
Chunk 处理模型
Spring Batch 的核心是 ItemReader → ItemProcessor → ItemWriter,按 chunk 提交:
@Bean
public Step settleStep() {
return stepBuilderFactory.get("settle")
.<RawBill, SettledBill>chunk(500)
.reader(billReader())
.processor(billProcessor())
.writer(billWriter())
.build();
}
chunk(500) 表示每攒 500 条才提交一次事务,既控内存又控事务粒度。我们读 300 万条,按 500 分批,内存稳定在 200MB 以内;如果 chunk 设成 50000,峰值内存会冲到 2G,GC 压力大。chunk 大小是性能旋钮:太小事务开销大,太大内存高,得测。
Reader 的坑
JdbcPagingItemReader 要配排序键,否则分页会漏读或重复读。我们一开始没配 sortKey,数据量大时出现了同一行被读两次,对账金额翻倍。补上唯一排序键(id)才正常:
new JdbcPagingItemReaderBuilder<RawBill>()
.sortKeys(Map.of("id", Order.ASCENDING))
.build();
并行步骤
对账分"收入"和"支出"两条独立线,用 Split 并行,缩短总时长:
return jobBuilderFactory.get("recon")
.split(new SimpleAsyncTaskExecutor())
.add(flow(incomeStep))
.add(flow(payoutStep))
.end()
.build();
原来串行跑要 40 分钟,并行后 22 分钟。注意 Split 用的线程池要够,SimpleAsyncTaskExecutor 默认无界,我们限制了并发数避免打爆数据库。
重启与跳过策略
这是它比裸 main 强的地方:
- 重启:Job 用同一个 JobParameters 重跑,已完成的 step 自动跳过,从断点继续;
- 跳过:某条脏数据解析失败,配 skip 策略跳过并记日志,不让整批挂掉:
.faultTolerant()
.skip(ParseException.class)
.skipLimit(100)
.listener(new SkipLoggingListener())
skipLimit(100) 表示最多跳过 100 条,超了才失败,避免无限跳过掩盖大问题。被跳过的行我们单独落表,白天人工复核,不丢数据。
幂等是关键
Spring Batch 的 restart 会从断点继续,但如果 Writer 不是幂等的,重跑会重复写。我们的 billWriter 用 "INSERT ... ON DUPLICATE KEY UPDATE",按业务主键去重,重跑安全。这点比框架本身更重要——框架给你续跑能力,幂等得自己保证。
监控与告警
Job 跑完我们接了监听器,把步级的读/处理/写计数和耗时打到 Prometheus,失败则发飞书。以前是"财务问起来了才知道挂了",现在是"挂了 1 分钟内告警就到"。还加了超时:单 step 超过 30 分钟没结束就告警,防慢 SQL 拖死。
文件类作业的坑
对账之外我们还用 Batch 处理银行回盘文件(每天几万个固定宽度的记录)。这种用 FlatFileItemReader,要配行 tokenizer 和字段映射。踩过的坑:文件最后一行没换行符,reader 会丢最后一条; Windows 和 Linux 换行符混用导致某字段带 \r 入库脏数据。我们统一在读取前用 sed 规整换行,并在 tokenizer 里 trim 每个字段,才干净。
跳过策略的边界
skipLimit 要慎设。我们一开始设 1000,结果有次源文件编码错了,几万条全解析失败,但因为没超 limit,任务"成功"跑完,产出了空结果,财务对账差一片。后来改成:解析类异常跳过上限设小(50),且跳过数超过阈值就告警人工介入,而不是默默吞掉。skip 是容错不是掩盖,超过预期就该亮红灯。
我们现在的批处理规范
沉淀出几条:chunk 大小按单条大小和内存测;Reader 必配排序键;Writer 必幂等;skip 必有上限且超限告警;Job 结束必发指标和通知。新接一个批任务,对着清单打勾就行,不用每次重新踩坑。Batch 这套框架把"可靠批处理"的最佳实践都内置了,关键是你得用对那几个旋钮。
并发处理提升吞吐
单 step 内还可以开多线程处理 chunk,用 task-executor 和 throttle-limit:
.taskExecutor(new SimpleAsyncTaskExecutor())
.throttleLimit(10)
这能并行处理多个 chunk,对 CPU 轻、IO 重的处理(比如每行调一次外部接口)提速明显。我们账单处理从单线程 25 分钟降到并行 8 分钟。但要注意:开并行后 chunk 之间顺序不保,writer 若依赖顺序(如按 id 递增写)会乱,得确保处理/写出是无序安全的,或者关掉并行。
一个生产真实案例
上个月银行回盘文件格式悄悄变了,多了一个字段,FlatFileItemReader 解析新格式时字段映射错位,50 万条里 3 万条脏数据。因为我们配了 skipLimit(50),超了直接失败告警,没默默产出错误对账,财务第一时间发现金额对不上。定位是文件格式变更,联系银行确认后更新 tokenizer 重跑,3 万条补平。如果当时没设 skip 上限,这 3 万条会被跳过、对账默默差一笔,月底才发现就晚了。这个上限配置,这次立了功。
留个问题
关于《Spring Batch 批处理作业设计与调优》里这个坑,你当时是怎么处理的?欢迎在评论区聊聊你踩过的类似情况。