for (Long id : pendingIds) {
TradeRecord record = tradeDao.findById(id);
checkAmount(record);
tradeDao.updateResult(record);
}
查询一条,处理一条,更新一条。数据少的时候跑得挺欢,一到月末批量对账,任务进度就开始装死。
CPU不高,内存也没爆,数据库连接却一直被占着。再往下看,问题很直接:每条数据都要经历一次查询、一次业务处理、一次更新,事务还切得稀碎。
这种任务继续往线程池里塞,我一般不赞成。线程多了只会把数据库先打趴下。
后来把这段批处理改成 Spring Batch,吞吐提升到原来的数倍。标题里的 500% 不是框架会变魔术,而是原来那套单条处理,把数据库往返、事务提交和异常恢复全做错了。
旧代码最难受的地方,还不是慢。
任务跑到中间失败,已经处理到哪儿没人知道。重启之后,要么从头再来,要么临时写 SQL 排除已完成数据。碰上脏数据更麻烦,一条异常就能把整个定时任务顶住。
Spring Batch 比较对我胃口的地方,是它没有把批处理理解成“开几个线程跑循环”,而是拆成了三个动作:
读取一批数据,处理这一批数据,再批量写回。
配置里真正关键的是 chunk。
@Configuration
@RequiredArgsConstructor
public class ReconcileJobConfig {
private final JobRepository jobRepository;
private final PlatformTransactionManager transactionManager;
private final DataSource dataSource;
@Bean
public Step reconcileStep(
ItemReader<TradeRow> tradeReader,
ItemProcessor<TradeRow, ReconcileResult> tradeProcessor,
ItemWriter<ReconcileResult> tradeWriter) {
return new StepBuilder("reconcileTradeStep", jobRepository)
.<TradeRow, ReconcileResult>chunk(500, transactionManager)
.reader(tradeReader)
.processor(tradeProcessor)
.writer(tradeWriter)
.faultTolerant()
.retryLimit(3)
.retry(TransientDataAccessException.class)
.build();
}
}
这里的 500,表示每读取并处理 500 条数据,统一提交一次事务。
别小看这一下。
原来处理 500 条数据,可能伴随着 500 次提交。现在只提交一次,数据库日志刷盘、网络交互和连接占用都会少很多。
但 chunk 也不是越大越好。
我见过有人直接改成 10000,觉得批量越大越快。结果任务确实少提交了几次,锁持有时间却明显变长,回滚一次也更疼。这个值要结合单条数据大小、SQL耗时和数据库承受能力慢慢调,不要靠拍脑袋。
读取数据我更愿意用分页,而不是一次性全部塞进内存。
@Bean
public JdbcPagingItemReader<TradeRow> tradeReader() {
return new JdbcPagingItemReaderBuilder<TradeRow>()
.name("pendingTradeReader")
.dataSource(dataSource)
.pageSize(500)
.selectClause("""
select id, order_no, pay_amount, settle_amount
""")
.fromClause("from trade_reconcile")
.whereClause("where handle_status = 'WAITING'")
.sortKeys(Map.of("id", Order.ASCENDING))
.rowMapper((rs, rowNum) -> new TradeRow(
rs.getLong("id"),
rs.getString("order_no"),
rs.getBigDecimal("pay_amount"),
rs.getBigDecimal("settle_amount")
))
.build();
}
排序字段必须稳定,最好是唯一递增字段。
拿一个会重复的时间字段做分页排序,看起来没问题,数据一边更新一边翻页时,就可能出现重复读取或者漏数据。这种坑平时不响,一到补数据的时候特别难查。
处理逻辑反而不用写得太花。
@Bean
public ItemProcessor<TradeRow, ReconcileResult> tradeProcessor() {
return trade -> {
boolean matched = trade.payAmount()
.compareTo(trade.settleAmount()) == 0;
return new ReconcileResult(
trade.id(),
matched ? "MATCHED" : "DIFFERENT",
matched ? null : "支付金额与结算金额不一致"
);
};
}
写入阶段不要再退回单条 update。既然已经走到 Spring Batch,就把批量写入做完整。
@Bean
public JdbcBatchItemWriter<ReconcileResult> tradeWriter() {
return new JdbcBatchItemWriterBuilder<ReconcileResult>()
.dataSource(dataSource)
.sql("""
update trade_reconcile
set handle_status = :status,
fail_reason = :reason,
handled_at = current_timestamp
where id = :id
and handle_status = 'WAITING'
""")
.beanMapped()
.assertUpdates(true)
.build();
}
where 后面多加一层状态判断,不是多余。
批处理任务最怕重复执行。任务重启、人工补跑、调度重复触发,都可能让同一条记录再进来一次。更新语句本身具备幂等性,后面排查会省很多事。
Spring Batch 还有一个优势,是失败后能接着跑。
Job、Step、读取位置和执行状态都会记录到元数据表里。任务中途挂掉,不需要再写一套“上次执行到第几条”的土办法。换个相同业务参数重新启动,框架会根据执行状态恢复。
不过这里有个前提:JobParameters 要带业务标识。
JobParameters parameters = new JobParametersBuilder()
.addString("billDate", billDate.toString())
.addString("source", "PAY_CHANNEL_A")
.toJobParameters();
不要每次都塞一个当前时间强行创建新任务实例。那样看起来解决了“任务不能重复启动”,实际上把失败恢复也一起废了。
Spring Batch 真正提升效率的地方,不只是跑得快。
它把分页读取、分块提交、批量写入、重试和断点恢复放进了一套稳定流程里。以前一个批处理脚本越改越像临时工程,失败之后全靠人盯;改完以后,任务至少知道自己处理到哪儿、为什么失败、还能不能继续。
批量任务的数据量一上来,for 循环通常不是最值得怀疑的地方。
先看数据库往返次数,再看事务边界,最后才轮到线程数。顺序搞反了,线程开得越多,事故来得越快。
扫码领红包
微信赞赏
支付宝扫码领红包
