
做数据迁移和批处理的时候,踩过坑的团队基本都经历过从“crontab + 脚本一把梭”到“引入专业框架”的阵痛。数据量小的时候怎么折腾都行,一旦订单量破千万、日志文件上 gb,内存溢出、中断后不知道跑哪儿了、一条脏数据卡死整个任务这些问题就会集中爆发。spring batch 在 spring 生态里算是个“重武器”,它的分层抽象、chunk 事务边界和内置的状态持久化,确实能把这些脏活累活规范化。下面直接进正题,聊聊在生产环境里怎么用 spring batch 搭一套能扛事、好排查的批处理管线。
1. 为什么传统脚本扛不住海量数据
早年团队喜欢用 linux crontab 调个 java 或 shell 脚本,数据量少的时候跑得挺欢,但规模一上来,底层缺陷根本藏不住:
- oom 是常态:很多人图省事,直接
select *或者file.readallbytes()把数据全塞进内存。jvm 堆一打满,应用直接 crash,连日志都来不及吐。 - 状态黑盒:脚本中途挂了,根本不知道处理到第几条。哪些入库了?哪些还在半事务?只能靠人工翻日志或查业务表,恢复成本极高。
- 容错全靠硬编码:遇到格式错误或唯一键冲突,整个任务直接中断。重试逻辑写死在循环里,没有退避策略,瞬间打满数据库连接池。
- 扩容困难:单节点串行跑,遇到几十 gb 的文件或亿级表,耗时直接拉胯。想水平扩展得自己写分片逻辑,容易引入数据倾斜或重复消费。
批处理的核心诉求其实就四点:流式读写保内存、分片并行提吞吐、chunk 切事务边界、元数据落库保断点续传。spring batch 的设计正好踩在这些点上。
2. 核心模型:job / step / chunk 到底在干什么
spring batch 把批处理拆得很干净,job → step → chunk 是它的骨架,reader → processor → writer 是数据流转的血脉。
2.1 分层职责
- job:顶层容器,代表一次完整的业务流。里面可以串多个 step,定义执行顺序和重启策略。
- step:实际干活的单元。事务控制、容错策略、监听器都是挂在 step 上的。
- chunk:这才是 spring batch 的灵魂。它不是指数据块大小,而是事务提交边界。比如配了
chunk(1000),框架会读 1000 条 → 走 processor → 交给 writer → 统一提交事务 → 记录进度。这种“微批提交”直接砍掉了长事务带来的锁竞争和 undo 日志膨胀问题。
2.2 数据流转管道
- itemreader:只负责读。每次
read()吐一条记录,底层必须是无状态或游标驱动的。千万别自己搞全量查询塞 list,框架自带的 reader 都是按需拉取的,内存占用恒定。 - itemprocessor:做转换、过滤、校验。返回
null就自动丢弃该条,不往下走。这里尽量保持轻量,别塞重型计算或阻塞 i/o,否则会把整个 step 拖慢。 - itemwriter:负责写。一次接收一个 list(长度就是 chunk size),批量刷到目标端。writer 自己不碰事务,全由 step 代理提交。
这套模型之所以稳,是因为内存里永远只扛着固定数量的对象。处理完一批就清掉,o 级别的内存开销,彻底跟 oom 绝缘。
3. 生产级配置实战
现在基本都用 spring boot starter 集成。下面直接上分片、大文件流读和 jdbc 批量写入的配置,附带线上调过参数。
3.1 分片并行(partitioner)
单表破亿或者文件超 50gb,单线程肯定不够看。spring batch 的 partitionstep 可以把数据切块,分给多个 worker 跑。
@bean
public step partitionedstep(partitionhandler partitionhandler) {
return stepbuilderfactory.get("partitionedstep")
.partitioner("workerstep", rangepartitioner())
.partitionhandler(partitionhandler)
.build();
}
@bean
public partitioner rangepartitioner() {
return gridsize -> {
// 实际项目里建议从 db 查 min/max id,这里简化示意
long minid = 1l;
long maxid = 10000000l;
long range = (maxid - minid + 1) / gridsize;
map<string, executioncontext> map = new hashmap<>();
for (int i = 0; i < gridsize; i++) {
executioncontext ctx = new executioncontext();
ctx.putlong("startid", minid + i * range);
ctx.putlong("endid", (i == gridsize - 1) ? maxid : minid + (i + 1) * range - 1);
map.put("partition_" + i, ctx);
}
return map;
};
}
配合 taskexecutorpartitionhandler 和线程池,就能在单机内多线程跑分片。如果节点多,结合 spring cloud data flow 或 k8s job,把 executioncontext 序列化下发,就能实现跨节点并行。
3.2 大文件流式读取
csv 或文本文件解析,必须走流。flatfileitemreader 底层就是 bufferedreader,每次只读一行,内存几乎不涨。
@bean
public flatfileitemreader<userrecord> flatfileitemreader() {
flatfileitemreader<userrecord> reader = new flatfileitemreader<>();
reader.setresource(new filesystemresource("/data/export_2023.csv"));
reader.setlinestoskip(1); // 跳过表头
defaultlinemapper<userrecord> linemapper = new defaultlinemapper<>();
linemapper.setlinetokenizer(new delimitedlinetokenizer(","));
beanwrapperfieldsetmapper<userrecord> fieldmapper = new beanwrapperfieldsetmapper<>();
fieldmapper.settargettype(userrecord.class);
linemapper.setfieldsetmapper(fieldmapper);
reader.setlinemapper(linemapper);
reader.setstrict(false); // 文件不存在时不直接抛异常,方便幂等重试
return reader;
}
遇到编码乱码或者自定义分隔符,自己实现个 linetokenizer 就行。注意大文件务必保证 reader 和 writer 的线程安全,spring batch 的 reader 默认是单线程安全的,分片模式下每个分片会拿到独立的 reader 实例。
3.3 jdbc 批量写入调优
jdbcbatchitemwriter 的性能卡点通常在网络往返和驱动层。配置本身不复杂,但参数不对吞吐差好几倍。
@bean
public jdbcbatchitemwriter<userrecord> jdbcbatchwriter(datasource datasource) {
jdbcbatchitemwriter<userrecord> writer = new jdbcbatchitemwriter<>();
writer.setdatasource(datasource);
writer.setsql("insert into users (id, name, age, created_at) values (:id, :name, :age, :created_at)");
writer.setitemsqlparametersourceprovider(new beanpropertysqlparametersourceprovider());
writer.setassertupdates(true); // 校验影响行数,防止静默失败
return writer;
}
线上实测有效的几个点:
- batch size 对齐:writer 内部调用
addbatch()的次数默认跟 chunk size 一致。建议 chunk 设在500~2000之间。设太大 jdbc 驱动层的内部缓冲区会撑爆,gc 频繁。 - 关闭主键回填:除非业务强依赖自增 id 回写,否则别用
generatedkeyholder。关掉return_generated_keys,mysql 场景下能稳定提升 30% 以上的吞吐。 - 连接池与 url 参数:hikaricp 的
maximumpoolsize得跟分片线程数对齐,别省连接。mysql jdbc url 必须带rewritebatchedstatements=true,不然驱动会把批量 sql 拆成单条发,白白浪费带宽。
4. 容错与断点续传
生产环境的数据从来都不干净,框架得能优雅兜底。
4.1 skip 与 retry 的配置逻辑
- retry:对付瞬时故障,比如网络抖动、db 锁等待、下游服务偶尔超时。配指数退避,别死循环重试。
- skip:对付脏数据,比如格式解析失败、违反唯一约束。记个日志或扔死信队列,任务继续往前跑。
@bean
public step faulttolerantstep(itemreader<user> reader, itemwriter<user> writer, platformtransactionmanager tm) {
return stepbuilderfactory.get("faulttolerantstep")
.<user, user>chunk(1000, tm)
.reader(reader)
.writer(writer)
.faulttolerant()
.retrylimit(3)
.retry(deadlockloserdataaccessexception.class, resourceaccessexception.class)
.skiplimit(50)
.skip(dataintegrityviolationexception.class, illegalargumentexception.class)
.listener(new customskiplistener()) // 把跳过的记录落地到死信表,事后人工核对
.build();
}
注意:skiplimit 和 retrylimit 是并存的。框架会先重试,达到上限后再走 skip。别把 skip 策略写得太宽泛,否则脏数据全漏过去,下游对账能哭死。
4.2 事务边界与避坑
spring batch 的事务由 step 全权代理。chunk 提交前,reader 和 writer 的操作在同一个事务里。中途报错,整个 chunk 回滚,元数据表不会更新,完美保证原子性。
血泪经验:绝对不要在 processor 里手动 @transactional 或者调 datasource 开新事务。这会打断框架的代理链,导致 jobrepository 里的进度更新失败,断点续传直接失效。processor 只负责纯计算或轻量校验。
4.3 断点续传的正确姿势
spring batch 把执行状态存在 6 张 batch_* 表里。重启时,框架去表里捞上次成功的 chunk 偏移量,接着读。
这里有个很多人搞错的细节:断点续传依赖相同的 jobparameters。jobinstance 的标识是 jobname + jobparameters。如果你每次启动都加个时间戳:
jobparameters params = new jobparametersbuilder()
.addlong("triggertime", system.currenttimemillis()) // ❌ 每次参数不同,框架会当成新任务,无法续跑
.addstring("source", "order_center")
.tojobparameters();
这样永远跑的是新实例。想续传,要么传一模一样的参数,要么通过 jobexplorer 查出上次挂掉的 jobexecution,调用 joblauncher.restart(jobexecution)。如果确实需要每次生成新实例但又想支持手动回滚,建议在业务表层面做幂等(insert ignore 或 on duplicate key update),别死磕框架状态。
5. 线上调优与可观测性
架构搭起来只是第一步,线上稳不稳,全看调优和监控。
5.1 内存与游标选择
分页读取(jpapagingitemreader)在深分页时 limit offset, size 会越拖越慢,而且如果排序键不唯一,容易漏数据或重复。很多人想换游标(jdbccursoritemreader),但游标会一直占用数据库连接,step 跑半小时,连接就占半小时,连接池很容易被打满。
我的建议是:数据量在千万级以内,用分页但务必加唯一排序键(比如主键 id);超大规模迁移优先用游标,但配合 fetchsize 控制(mysql jdbc 默认是 integer.min_value 流式拉取,别改),并且单独给 reader 配个连接池,别跟 writer 抢连接。
5.2 长事务与批大小动态化
别把 chunk 大小写死。业务高峰期数据库压力大,chunk 设 500 能减少锁持有时间;夜间低峰期拉到 2000,减少事务提交开销。可以通过 steplistenersupport 结合 prometheus 指标动态调整,或者干脆用两个不同参数的 step 错峰执行。
写入前如果表里有一堆非唯一索引,归档场景下可以临时 alter table ... disable keys,写完再 enable keys,重建索引比逐行维护快得多。当然,线上核心交易表别这么玩,只适用于日志/流水类归档表。
5.3 监控大盘怎么接
spring boot actuator 原生暴露了 batch 的 micrometer 指标,直接接 grafana 就行。
management:
endpoints:
web:
exposure:
include: health,metrics
metrics:
tags:
application: data-sync-service重点盯这几个指标:
spring.batch.job.execution.count:任务触发频次spring.batch.step.read.count/write.count:读写量对齐,排查数据丢失spring.batch.step.skip.count:跳过率突增通常是上游数据源出了问题- 耗时分布:通过自定义 timer 记录 processor 阶段耗时,80% 的瓶颈都在数据转换逻辑太重。
配个告警规则,skip 率超过 5% 或者 job 状态变 failed,直接推钉钉/企微。别等对账发现少了数据才去查日志。
6. 选型建议与调度协同
spring batch 不是万能药,得看场景用。
如果是夜间 etl、跨库同步、大文件清洗这类离线任务,spring batch 的事务控制和状态机机制非常顺手。但如果是秒级高频微批,或者需要滑动窗口聚合的实时流,硬塞 spring batch 只会适得其反。它的元数据落库和 chunk 状态机本身就有开销,延迟下不去。这种场景直接上 kafka streams 或 flink 更合适,轻量级任务用定时任务+消息队列就能搞定。
生产环境的标准打法是“调度层 + 执行层”解耦:
调度层交给 xxl-job、airflow 或 dolphinscheduler。它们负责任务依赖编排、分片参数下发、全局失败重试和告警。执行层用 spring batch 接收调度器传来的 shardindex 和 shardtotal,通过 partitioner 切分数据范围,独立跑流式处理,自己管 chunk 事务和 jobrepository 状态。调度器看宏观进度,batch 管微观数据流转,两边通过 webhook 或 mq 对齐状态。这套组合拳打下来,数据管道基本能做到“调度可管、执行可控、断点可续、异常可溯”。
框架只是工具,真正决定稳定性的是对事务边界、数据流向和失败场景的预判。把 chunk 粒度调对,把 skip 策略写准,把调度与执行拆干净,批处理就不再是定时炸弹。
到此这篇关于spring batch海量数据迁移与批处理架构完整教学的文章就介绍到这了,更多相关spring batch数据迁移与批处理内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论