Java批处理实战指南:从笨重脚本到高性能流式处理的架构跃迁
📚 目录导读
- 为什么还需要批处理?—— 不仅仅是历史包袱
- Java批处理的三大核心模式:从OOP到函数式
- 案例实战:金融对账系统的百万级数据批处理
- 1 传统JDBC批处理优化(PreparedStatement + addBatch)
- 2 并行流式处理(Parallel Stream)与ForkJoinPool陷阱
- 3 内存高效策略:分页拉取 vs 游标(Cursor)模式
- 批处理中的事务边界与失败恢复(Checkpoint机制)
- Spring Batch框架精讲:重试、跳过、监听器
- 高频问答:解决你写批处理时的5个痛点
- 从“能跑”到“跑得稳”的工程化思考
为什么还需要批处理?—— 不仅仅是历史包袱
在微服务与实时流处理(如Kafka Streams)盛行的今天,很多人误以为批处理是“老古董”,但实际上,财务结算、ETL数据仓库、报表生成、日终对账这些场景,依然依赖高吞吐量的批处理任务,原因在于:数据一致性要求极高(不允许丢数据)、成本可控(相比实时计算,批量吞吐能显著降低硬件开销)。

Java在批处理领域的地位无法撼动,得益于其JVM内存管理、强类型安全以及极其丰富的生态(Spring Batch、Apache Flink的批模式),本文将通过一个完整的电商订单对账案例,演示如何实现每秒处理5万条记录且不发生OOM的批处理程序。
Java批处理的三大核心模式
在写任何代码前,我们必须先区分三种批处理模式,这会直接影响架构设计:
- 模式A:读-处理-写(Chunk-Oriented):这是最主流的,每读取N条数据(如1000条)作为一个Chunk,处理完统一提交事务,这能极大减少数据库连接开销。
- 模式B:流式逐条处理(ItemProcessor):适合数据源本身不要求强一致性的场景,但要注意逐条提交事务会慢100倍。
- 模式C:并行分区处理(Partitioner):将大表按ID范围或时间分片,用多线程/分布式节点并行读取,这是性能提升的最关键手段。
关键认知:“批处理”不等于“一次性把所有数据加载到内存”,凡是写List<Data> allData = dao.findAll()的代码,在千万级数据面前都是“自杀式”写法。
案例实战:金融对账系统的百万级数据批处理
假设我们有交易流水表(trade_flow),共800万条记录,需计算每日手续费并写入汇总表(trade_summary)。
❌ 错误示范(新手常见):
// 一次性全查,内存直接爆掉
List<TradeFlow> list = jdbcTemplate.query("SELECT * FROM trade_flow", mapper);
for (TradeFlow t : list) {
// 逐条insert,效率极低
jdbcTemplate.update("INSERT INTO trade_summary...", t.getAmount());
}
✅ 正确方案:分页游标 + 批量提交
步骤1:使用JDBC游标模式(流式读取)
设定fetchSize为1000,但不真正加载1000条到内存,而是由数据库驱动逐条从网络流中读取。
// 连接必须设置为只向前、只读模式,否则无效
try (Connection conn = dataSource.getConnection();
PreparedStatement ps = conn.prepareStatement("SELECT * FROM trade_flow WHERE settle_date = ?",
ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY)) {
ps.setFetchSize(1000); // 关键:让MySQL/PostgreSQL流式返回
ps.setDate(1, today);
try (ResultSet rs = ps.executeQuery()) {
// 累积到1000条后批量处理
List<TradeFlow> chunk = new ArrayList<>(1000);
while (rs.next()) {
TradeFlow flow = mapper.mapRow(rs, 0);
chunk.add(flow);
if (chunk.size() == 1000) {
processChunkAndFlush(chunk); // 批量写入
chunk.clear();
}
}
// 处理最后不足1000条的余量
if (!chunk.isEmpty()) processChunkAndFlush(chunk);
}
}
步骤2:批量写入使用addBatch()
不要用update()每次提交,而是用PreparedStatement的addBatch()累积1000条后统一执行。
private void processChunkAndFlush(List<TradeFlow> chunk) {
String sql = "INSERT INTO trade_summary (acct_id, total_fee, trade_count, calc_date) VALUES (?,?,?,?)";
try (PreparedStatement ps = conn.prepareStatement(sql)) {
for (TradeFlow t : chunk) {
ps.setString(1, t.getAcctId());
ps.setBigDecimal(2, t.getFee());
ps.setInt(3, 1);
ps.setDate(4, t.getDate());
ps.addBatch();
}
ps.executeBatch(); // 一次网络往返,执行1000条
conn.commit(); // 手动提交,控制事务边界
}
}
批处理中的事务边界与失败恢复(Checkpoint机制)
批处理最怕的是跑到第500万条时突然断电,如果事务跨度太大(比如一条事务包含100万条),回滚成本极高。业界标准解法是设置Checkpoint(检查点):
- 策略:每处理完1000条Chunk,记录一个
batch_id和last_processed_id到独立的batch_status表。 - 恢复:重启任务时,先查询
last_processed_id,然后用WHERE id > last_processed_id重新拉取,保证幂等性。
⚠️ 重要提示:不要在批处理中开启分布式事务(如XA),因为长时间占用数据库锁会导致死锁,正确的做法是“最终一致性”:允许短时数据不一致,通过补单Job修复。
Spring Batch框架精讲:重试、跳过、监听器
如果不想手写这么多底层控制,Spring Batch提供了完善的分步式批处理(Step) 架构,其核心优势在于三个机制:
- 重试(Retry):对于偶发的数据库连接超时,配置
retry-limit="3",且只对TransientDataAccessException异常生效。 - 跳过(Skip):如果是脏数据(如无法解析日期),我们希望能跳过这一条而不影响整个Batch,写入
skip.log文件中,这比手动写try-catch优雅得多。 - 监听器(Listener):通过
ItemReadListener和ChunkListener在批次前/后执行系统通知或清理内存。
核心配置片段示例:
<batch:job id="settlementJob">
<batch:step id="step1">
<batch:tasklet>
<batch:chunk reader="jdbcReader" writer="jdbcWriter" commit-interval="1000" retry-limit="2" skip-limit="10">
<batch:skippable-exception-classes>
<batch:include class="java.lang.IllegalArgumentException"/>
</batch:skippable-exception-classes>
<batch:retryable-exception-classes>
<batch:include class="org.springframework.dao.DeadlockLoserDataAccessException"/>
</batch:retryable-exception-classes>
</batch:chunk>
</batch:tasklet>
</batch:step>
</batch:job>
高频问答:解决你写批处理时的5个痛点
Q1:我的批处理程序跑着跑着突然OOM(内存溢出),最常见原因是什么?
A:90%是因为使用了
List收集了所有查询结果,或者Stream收集器.collect(Collectors.toList())。解决:改用流式游标模式,或者用Stream的limit配合skip实现手动分页,但始终要控制单片大小(如每片1000)。
Q2:批处理时,数据库连接池连接不够用了怎么办?
A:首先检查是否忘记关闭
ResultSet或Statement,在并行批处理场景下,不要将Connection设为多线程共享,而是使用ThreadLocal<Connection>管理,且连接池最大连接数要设置为CPU核数 * 2 + 数据库磁盘数。
Q3:如何保证批处理任务不会重复执行(比如两台服务器同时跑)?
A:使用分布式锁(如Redis的
SETNX),或者数据库唯一约束(如batch_date字段设置唯一索引),最稳妥的是在启动前先向batch_control表插入一条RUNNING状态记录,获取成功才继续。
Q4:批处理性能一直提不上去,CPU利用率只有10%?
A:很可能是单线程读取数据库的瓶颈,尝试对主键ID进行取模分片,比如
WHERE MOD(id, 8) = 0,然后用8个线程并行查询,配合ForkJoinPool(可用parallelStream(),但注意共享变量线程安全)。
Q5:批处理中需要调用外部API(如短信通知),如何防止接口超时导致整个Batch失败?
A:不要在Chunk中间调用外部API,这会让事务锁挂起,正确的做法是将需要通知的数据先写入
notification_queue表,批处理只负责快速写队列,由独立的小线程池异步消费队列发送通知,这样解耦了速度差异。
从“能跑”到“跑得稳”的工程化思考
编写高效的Java批处理案例,本质上是对IO模型和内存模型的理解,作为开发者,我们需要始终记住三个铁律:
- 永远不要相信一次性加载是安全的——除非数据量确定小于1000条。
- 事务提交频率过高(每条提交)和过低(几百万条提交)都是灾难——推荐以1000-5000为一个Chunk。
- 监控与恢复机制比处理逻辑本身更重要——在线上环境中,
失败重启的幂等性是批处理的第一准则。
希望本文的综合解析能够帮你跳出“写死循环跑批”的困境,用更工程化、更高性能的Java技术去应对数据量持续增长的业务挑战。