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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 20:19:58