这段代码我第一眼就不太信。
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 循环通常不是最值得怀疑的地方。

先看数据库往返次数,再看事务边界,最后才轮到线程数。顺序搞反了,线程开得越多,事故来得越快。

扫码领红包

微信赞赏支付宝扫码领红包

发表回复

后才能评论