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

MongoDB聚合Java优化:高效实现多维度去重计数查询

MongoDB 百万级数据聚合性能优化方案

问题背景

现有MongoDB events集合含百万级文档,结构如下:

{
    "id": "fg3df23",
    "config_id": "49gjd3s",
    "timestamp": 1649158819315,
    "status": "active",
    "user_data": {
        "ip": "<some IP>",
        "country": "US"
    }
}

需求为:筛选指定config_id列表、timestamp在指定范围内的事件,按国家统计:

  1. 各状态的事件分布
  2. 每个国家的去重IP数量
  3. 全局总去重IP数量

期望返回结构对应Java类:

public class CountriesResponse  {
    private int countryTotalDistinctIps;
    private Set<CountryResponse> data;

    public static class CountryResponse {
        private String country;
        private int distinctIPCount;
        private Map<String, Integer> statusBreakdown;
    }
}

现有Spring Data MongoDB聚合方案处理130万条匹配事件耗时约15秒,性能不佳;曾尝试按user_data.country、user_data.ip、status分组,引发$group阶段内存不足,开启allowDiskUsage:true后性能更差。

优化方案

一、重构聚合Pipeline,并行处理统计任务

核心思路是用$facet并行执行两个统计分支,拆分原有的单次分组逻辑,避免大集合addToSet的内存开销:

聚合步骤说明

  1. Match阶段:保留原过滤逻辑,利用索引快速筛选数据
  2. Project阶段:仅保留所需字段(user_data.country、user_data.ip、status),减少Pipeline内数据传输量
  3. Facet阶段:并行处理两个统计任务:
    • 分支1:按国家+状态分组,统计各状态事件数
    • 分支2:先按国家+IP分组去重,再按国家统计去重IP数量
  4. 手动合并两个分支的结果,组装成期望的响应结构

Spring Data MongoDB 代码实现

// 1. Match阶段:过滤符合条件的数据
MatchOperation matchOp = match(Criteria.where("config_id").in(configIdsList)
        .and("timestamp").gte(start).lte(end));

// 2. Project阶段:仅保留需要的字段,减少数据传输
ProjectionOperation keepNeededFields = project("user_data.country", "user_data.ip", "status");

// 3. Facet阶段:并行执行两个统计分支
FacetOperation facetOp = facet(
        // 分支1:统计每个国家的状态分布
        Arrays.asList(
                group("user_data.country", "status")
                        .count().as("count"),
                project("count")
                        .and("country").previousOperation().andInclude("user_data.country")
                        .and("status").previousOperation().andInclude("status")
        ), "statusBreakdowns",
        // 分支2:统计每个国家的去重IP数量
        Arrays.asList(
                group("user_data.country", "user_data.ip"), // 按国家+IP去重
                group("user_data.country")
                        .count().as("distinctIPCount"), // 统计国家去重IP数
                project("distinctIPCount")
                        .and("country").previousOperation().andInclude("user_data.country")
        ), "countryIpCounts"
);

// 构建聚合查询
Aggregation aggregation = Aggregation.newAggregation(matchOp, keepNeededFields, facetOp);

// 执行聚合并处理结果
AggregationResults<Document> results = mongoTemplate.aggregate(aggregation, "events", Document.class);
Document facetResult = results.getUniqueMappedResult();
List<Document> statusBreakdowns = (List<Document>) facetResult.get("statusBreakdowns");
List<Document> countryIpCounts = (List<Document>) facetResult.get("countryIpCounts");

// 组装状态分布映射
Map<String, Map<String, Integer>> countryStatusMap = new HashMap<>();
for (Document doc : statusBreakdowns) {
    String country = doc.getString("country");
    String status = doc.getString("status");
    int count = doc.getInteger("count");
    countryStatusMap.computeIfAbsent(country, k -> new HashMap<>()).put(status, count);
}

// 组装最终响应
Set<CountryResponse> data = new HashSet<>();
int countryTotalDistinctIps = 0;
for (Document doc : countryIpCounts) {
    String country = doc.getString("country");
    int distinctIPCount = doc.getInteger("distinctIPCount");
    countryTotalDistinctIps += distinctIPCount;

    CountryResponse cr = new CountryResponse();
    cr.setCountry(country);
    cr.setDistinctIPCount(distinctIPCount);
    cr.setStatusBreakdown(countryStatusMap.getOrDefault(country, Collections.emptyMap()));
    data.add(cr);
}

CountriesResponse response = new CountriesResponse();
response.setCountryTotalDistinctIps(countryTotalDistinctIps);
response.setData(data);

二、索引优化

  1. 保留核心过滤索引:config_id_1_timestamp_-1是Match阶段最有效的索引,前缀config_id适配in查询,后续timestamp适配范围查询,可快速过滤数据。
  2. 创建覆盖索引:新增config_id_1_timestamp_-1_user_data.country_1_user_data.ip_1_status_1索引,覆盖Match和Project阶段所需的所有字段,MongoDB无需回表读取原文档,直接从索引获取数据,大幅降低磁盘IO开销。

三、其他优化建议

  • 避免使用addToSet处理大集合:该操作需在内存中维护去重集合,数据量大时内存开销极高,改用两次分组去重的方式更高效。
  • 开启查询分析:用explain()检查聚合Pipeline是否正确使用索引,定位内存或IO瓶颈。
  • 分片集群扩容:若数据量持续增长,可将events集合分片,分散聚合计算压力到多个节点。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 03:54:57