Spring Batch读取MongoDB时仅读取半数数据的问题排查
Spring Boot Batch MongoDB 读取数据每次仅处理半数的问题分析与解决
问题现象
使用Spring Boot Batch从MongoDB读取未处理文档(无processedAt字段),每次运行任务仅处理剩余数据的半数(如30k→15k→7.5k……),直到剩余文档不足1000条时才一次性处理完毕(因设置了cursorBatchSize(1000))。已排除内存问题(一次性读取所有文档时运行正常)和竞态条件(单线程运行问题仍存在)。
核心原因
问题出在MongoDB游标特性与查询逻辑的配合上:
- 无稳定排序条件:当前查询未指定排序规则,MongoDB默认按文档物理存储顺序返回结果。当文档被更新(添加
processedAt字段)时,若文档大小增加,MongoDB会将其移动到磁盘的新位置,导致后续查询无法按固定顺序遍历所有未处理文档。 - 非快照游标:默认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
相关产品推荐
相关产品推荐

