Java应用执行MongoDB聚合时遇Prematurely reached end of stream错误求助
问题描述
在Java应用中执行长时间运行的聚合操作时,持续遇到以下错误:
连接到mymongoserver:27017的[connectionId{localValue:5, serverValue:32617}]发生套接字异常,所有到该地址的连接将被关闭。
org.springframework.data.mongodb.UncategorizedMongoDbException: Prematurely reached end of stream; 嵌套异常为com.mongodb.MongoSocketReadException: Prematurely reached end of stream。
当前配置
MongoDB URI:
mongodb://user:mymongoserver:27017/dbname?authSource=dbname&replicaSet=rs-mdb&readPreference=secondaryPreferred&socketTimeoutMS=10000000&maxIdleTimeMS=10000000&connectTimeoutMS=10000000
聚合操作代码
public void cleanDuplicateActiveRecords() { // Step 1: Find duplicate active records grouped by externalId & // serviceSpecification.name // Step 1: Re-map dot-notated fields to simple field names try { int batchSize = 1000; long offset = 0; boolean hasMoreData = true; AggregationOptions options = AggregationOptions.builder().maxTime(Duration.ofHours(2)).allowDiskUse(true) .build(); while (hasMoreData) { Aggregation aggregation = Aggregation.newAggregation( Aggregation.match(Criteria.where("order.state").in(Arrays.asList("ACTIVE", "SUSPENDED"))), Aggregation.project("_id", "lastUpdated").and("order.externalId").as("externalId") .and("order.serviceSpecification.name").as("orderName"), Aggregation.sort(Sort.by(Sort.Direction.DESC, "lastUpdated")), Aggregation.group("externalId", "orderName").push("_id").as("ids").push("lastUpdated") .as("timestamps").count().as("count"), Aggregation.match(Criteria.where("count").gt(1)), Aggregation.project("ids", "timestamps"), Aggregation.skip(offset), Aggregation.limit(batchSize) ).withOptions(options); List<DuplicateOrder> duplicates = mongoTemplate .aggregate(aggregation, "orders", DuplicateOrder.class) .getMappedResults(); if (duplicates.size() < batchSize) { hasMoreData = false; } else { offset += batchSize; } log.info("The number of duplicated records {}", duplicates.size()); // Step 2: Delete all but the latest record for each group for (DuplicateOrder duplicate : duplicates) { List<String> ids = duplicate.getIds(); if (ids.size() > 1) { // Keep only the most recent record, delete the rest List<String> idsToDelete = ids.subList(1, ids.size()); log.info("deleting records with ids: {}", idsToDelete); Query deleteQuery = new Query(Criteria.where("_id").in(idsToDelete)); mongoTemplate.remove(deleteQuery, "orders"); } } } } catch (Exception e) { log.error("Exception happened: {}", e); e.printStackTrace(); } }
解决方案
1. 优化聚合分页逻辑,避免重复扫描数据
当前用skip+limit的分页方式,数据量大时每次聚合都要扫描前面所有数据,拖长查询时间触发超时。改用游标分批处理:
- 移除
skip和limit,通过游标批次获取结果,避免重复扫描 - 修改代码示例:
// 调整聚合选项,添加游标批次配置 AggregationOptions options = AggregationOptions.builder() .maxTime(Duration.ofHours(2)) .allowDiskUse(true) .cursor(new Document("batchSize", 1000)) // 设置游标每次返回的批次大小 .build(); // 获取聚合结果的迭代器 AggregationResults<DuplicateOrder> results = mongoTemplate.aggregate(aggregation, "orders", DuplicateOrder.class); Iterator<DuplicateOrder> iterator = results.iterator(); int processedGroups = 0; while (iterator.hasNext()) { DuplicateOrder duplicate = iterator.next(); // 处理重复记录删除逻辑 List<String> ids = duplicate.getIds(); if (ids.size() > 1) { List<String> idsToDelete = ids.subList(1, ids.size()); log.info("deleting records with ids: {}", idsToDelete); Query deleteQuery = new Query(Criteria.where("_id").in(idsToDelete)); mongoTemplate.remove(deleteQuery, "orders"); } processedGroups++; if (processedGroups % 1000 == 0) { log.info("Processed {} duplicate groups", processedGroups); } }
2. 调整MongoDB连接参数,适配长时操作
- 确认服务器端
net.socketTimeoutMS配置不小于客户端socketTimeoutMS(服务器默认5分钟,你的客户端设置约2.7小时,需同步修改服务器配置或调整客户端参数匹配) - 添加
heartbeatSocketTimeoutMS参数,避免副本集心跳超时误判连接失效,建议设置为30000ms(30秒) - 修改后的URI示例:
mongodb://user:mymongoserver:27017/dbname?authSource=dbname&replicaSet=rs-mdb&readPreference=secondaryPreferred&socketTimeoutMS=10000000&maxIdleTimeMS=10000000&connectTimeoutMS=10000000&heartbeatSocketTimeoutMS=30000
3. 添加复合索引,加速聚合查询
针对聚合的match、sort、group阶段创建复合索引,减少磁盘IO和查询时间:
// 在MongoDB Shell中执行索引创建 db.orders.createIndex({ "order.state": 1, "order.externalId": 1, "order.serviceSpecification.name": 1, "lastUpdated": -1 })
4. 批量执行删除操作,减少写请求频率
当前每处理一个分组就执行一次删除,频繁写操作会增加MongoDB负载。改为积累一定量的待删除ID后批量删除:
List<String> allIdsToDelete = new ArrayList<>(); while (iterator.hasNext()) { DuplicateOrder duplicate = iterator.next(); List<String> ids = duplicate.getIds(); if (ids.size() > 1) { allIdsToDelete.addAll(ids.subList(1, ids.size())); } // 每积累1000个ID执行一次批量删除 if (allIdsToDelete.size() >= 1000) { Query deleteQuery = new Query(Criteria.where("_id").in(allIdsToDelete)); mongoTemplate.remove(deleteQuery, "orders"); log.info("Deleted {} records", allIdsToDelete.size()); allIdsToDelete.clear(); } } // 处理剩余未删除的ID if (!allIdsToDelete.isEmpty()) { Query deleteQuery = new Query(Criteria.where("_id").in(allIdsToDelete)); mongoTemplate.remove(deleteQuery, "orders"); log.info("Deleted {} remaining records", allIdsToDelete.size()); }
5. 排查MongoDB服务器及网络状态
- 查看MongoDB服务器日志,确认是否存在负载过高、网络波动、防火墙超时等导致连接中断的原因
- 副本集环境下,可临时将
readPreference改为primary,测试是否因节点切换导致连接异常
内容的提问来源于stack exchange,提问作者Nijaoui Khalil
相关产品推荐
相关产品推荐

