You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Batch读取MongoDB时仅读取半数数据的问题排查

Spring Boot Batch MongoDB 读取数据每次仅处理半数的问题分析与解决

问题现象

使用Spring Boot Batch从MongoDB读取未处理文档(无processedAt字段),每次运行任务仅处理剩余数据的半数(如30k→15k→7.5k……),直到剩余文档不足1000条时才一次性处理完毕(因设置了cursorBatchSize(1000))。已排除内存问题(一次性读取所有文档时运行正常)和竞态条件(单线程运行问题仍存在)。

核心原因

问题出在MongoDB游标特性与查询逻辑的配合上:

  1. 无稳定排序条件:当前查询未指定排序规则,MongoDB默认按文档物理存储顺序返回结果。当文档被更新(添加processedAt字段)时,若文档大小增加,MongoDB会将其移动到磁盘的新位置,导致后续查询无法按固定顺序遍历所有未处理文档。
  2. 非快照游标:默认MongoDB游标是动态的,会实时反映集合变化。当分批读取时,已处理的文档会被移出查询结果集,游标在获取下一批数据时,可能跳过部分未处理文档(因存储位置变化导致扫描范围不完整)。

解决方案

方案1:添加稳定排序条件(推荐)

为查询添加基于唯一且不变字段的排序(如_id),确保MongoDB按固定顺序遍历文档,不受物理存储位置变化影响。

修改ItemReader的查询逻辑:

@Bean(name = "myReader")
@StepScope
public MongoItemReader<TroublesomeModel> reader() {
    Query query = new Query();
    query.addCriteria(Criteria.where("processedAt").exists(false))
         .with(Sort.by(Sort.Direction.ASC, "_id")); // 添加按_id排序

    query.noCursorTimeout();
    query.cursorBatchSize(1000);

    MongoItemReader<TroublesomeModel> reader = new MongoItemReader<>();
    reader.setTemplate(mongoTemplate);
    reader.setTargetType(TroublesomeModel.class);
    reader.setSaveState(false);
    reader.setQuery(query);

    return reader;
}

方案2:使用快照游标

启用MongoDB的快照游标,让查询基于任务启动时的集合快照,不受后续文档更新影响。注意:快照游标不适用于分片集合,且会增加性能开销。

修改ItemReader的查询逻辑:

@Bean(name = "myReader")
@StepScope
public MongoItemReader<TroublesomeModel> reader() {
    Query query = new Query();
    query.addCriteria(Criteria.where("processedAt").exists(false))
         .snapshot(); // 启用快照

    query.noCursorTimeout();
    query.cursorBatchSize(1000);

    MongoItemReader<TroublesomeModel> reader = new MongoItemReader<>();
    reader.setTemplate(mongoTemplate);
    reader.setTargetType(TroublesomeModel.class);
    reader.setSaveState(false);
    reader.setQuery(query);

    return reader;
}

额外优化:修改Writer的更新方式

当前Writer使用updateMulti,但_id是唯一的,改用updateFirst更贴合业务逻辑(更新单个匹配文档):

@Override
public void write(List<? extends TroublesomeModel> items) {
    BulkOperations bulkOperations = mongoTemplate.bulkOps(BulkOperations.BulkMode.UNORDERED, TroublesomeModel.class);

    for (TroublesomeModel item: items) {
        item.setAsProcessed();
        bulkOperations.updateFirst( // 替换为updateFirst
                Query.query(Criteria.where("_id").is(item.getId())),
                new Update().set("processedAt", item.getProcessedAt())
        );
    }

    bulkOperations.execute();
}

验证

修改后重新运行任务,即可一次性读取并处理所有未处理文档,不会出现每次仅处理半数的情况。

内容的提问来源于stack exchange,提问作者AnonymousFox

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 17:15:16