• mongodb.aggregate 索引查询+分组group+排序sort 优化查询效率


    一、主要问题

    系统中有一张温控终端状态的表tcState,记录了所有温控终端的温控状态,大约有1600万条数据。需求就是通过列表的形式展示出所有温控终端最新的温控终端状态,查询条件有公司id、终端分组id、温控终端id、状态读取时间

    基本的查询逻辑就是根据查询条件、索引筛选数据,对数据根据温控终端进行分组、按照最新时间排序。

    但是第一版做出来,发现查询速度很慢,一次查询用了7/8秒钟,完全无法接受,于是开始对整个过程进行拆解分析,找出问题和瓶颈。

    二、索引查询

    复合索引

    首先索引问题,需要检查索引是否命中,可以使用explain()函数查看MongoDB在执行过程的准备过程。

    这是我修改后建立的复合索引:

    companyId ASC, tcGroupId ASC, tempControllerId ASC, readingTime DESC

    索引失效:

    复合索引生效的原则是前缀生效

    以下的场景索引是生效的,因为它们都包含了索引的前缀

    1. db.tcState.find({companyId :18,tcGroupId :1,readingTime :'2022-08-31'})
    2. db.tcState.find({companyId :18,tcGroupId :1})
    3. db.tcState.find({companyId :18})
    4. db.tcState.find({tcGroupId :1,companyId :1})
    5. db.tcState.find({companyId :1,readingTime :'2022-08-31'})

    以下场景都是无效,因为它们没有包含索引前缀

    1. db.student.find({tcGroupId :1})
    2. db.student.find({tcGroupId :1,readingTime :'2022-08-31'})

    重复建立索引:

     如图所示,index1索引和groupQuery索引

    以下的场景使用的索引都是 groupQuery,第一个索引index1永远不会发挥作用,反而浪费了存储和降低了插入性能

    1. db.tcState.find({companyId :1})
    2. db.tcState.find({companyId :1,readingTime :'2022-08-31'})

     但是,为什么我还是要建立index1这个索引呢?

    这是因为我发现,如果我用readingTime 排序的时候,groupQuery索引就无法命中了,具体原因,暂时不太清楚

    db.tcState.find({companyId:1}).sort({readingTime:-1})

    三、排序sort

     这里排序出现的主要问题就是到底在什么时候排序?

    1. db.tcState.aggregate([
    2. { $match : { companyId : 1}},
    3. { $sort : { readingTime : -1}},
    4. { $group : { _id : "$tempControllerId",tempControllerId:{$first:"$tempControllerId"}, readingTime:{$first:"$readingTime"}}},
    5. { $skip : 0},
    6. { $limit : 20}
    7. ]);

    第一次,我选择了,先查询数据,再根据时间进行排序找到最新数据,最后进行分组 ,但这样其实相当于要对全部数据进行排序,对资源的消耗是巨大的,虽然现在有了索引,但是结合分组时,默认只取第一条数据的逻辑,其实我们可以不用按时间排序,默认分组后,就是获取到的每一个终端最新的一条数据。这样我们就可以把时间排序放在分组后,分组之后的数据量是很小的,速度更快。

    1. db.tcState.aggregate([
    2. { $match : { companyId : 1}},
    3. { $group : { _id : "$tempControllerId",tempControllerId:{$first:"$tempControllerId"}, readingTime:{$first:"$readingTime"}}},
    4. { $sort : { readingTime : -1}},
    5. { $skip : 0},
    6. { $limit : 20}
    7. ]);

    四、分组group

     使用分组,大家要先对aggregate 管道的相关功能有足够了解,这里我就不详细说明了,直接进行我的分析:

    1. db.tcState.aggregate([
    2. { $match : { companyId : 1}},
    3. { $group : { _id : "$tempControllerId",tempControllerId:{$first:"$tempControllerId"}, readingTime:{$first:"$readingTime"}}},
    4. { $sort : { readingTime : -1}},
    5. { $skip : 0},
    6. { $limit : 20}
    7. ]);

     分组,是根据_id 对应的tempControllerId,进行分组的。而_id后边的几个字段,则是你要返回的数据字段,具体返回字段是要你自己写的。

    通过first,则可以得到分组后,每一个终端对应的最新的一条数据。

    但是在这里同样有一个问题,虽然我们通过建立索引,使查询和排序的速度大大提升,但是在分组这一步还是遇到了瓶颈,索引对于分组并没有作用,这就意味着,我们要把索引查询到的终端对应的所有历史数据一起进行分组。这个分组的是很慢的。

    我的解决方案就是:我们其实只需要通过分组留下最新的一条数据,那么我们就可以通过时间字段进行筛选,我们只查询最新的一天或者三天内的数据,通过筛选条件,人为减少分组数据,这样分组的速度就会大大提升。

    readingTime : { "$gte" : ISODate("2022-08-29T00:00:00Z"), "$lt" : ISODate("2022-08-31T00:00:00Z") }

    这就是我认为最核心的一个点,就是通过筛选条件,减少你要操作的数据量。

    1. db.tcState.aggregate([
    2. { $match : { companyId : 1,
    3. readingTime : { "$gte" : ISODate("2022-08-29T00:00:00Z"), "$lt" : ISODate("2022-08-31T00:00:00Z") }
    4. }},
    5. { $group : { _id : "$tempControllerId",tempControllerId:{$first:"$tempControllerId"}, readingTime:{$first:"$readingTime"}}},
    6. { $sort : { readingTime : -1}},
    7. { $skip : 0},
    8. { $limit : 20}
    9. ]
    10. );

    五、最终结果

    通过以上索引、排序、分组的优化,最终我的查询时间缩减到了几十毫秒之内,完全符合业务需求,虽然只是不到2000万的数据量,但是我觉得优化的思路才是非常重要的。

    最后再附上java代码对应的实现方式:

    1. public List searchTcStateList(TcStateDto tcStateDto) {
    2. //组织查询条件集合
    3. Criteria[] criteriaArray = getTcStateCriteriaArray(tcStateDto);
    4. Aggregation aggregation = Aggregation.newAggregation(
    5. match(new Criteria().andOperator(criteriaArray)),
    6. Aggregation.group("tempControllerId")
    7. .first("tempControllerId").as("tempControllerId")
    8. .first("companyId").as("companyId")
    9. .first("tcGroupId").as("tcGroupId")
    10. .first("readingTime").as("readingTime")
    11. .first("tempControllerAddress").as("tempControllerAddress")
    12. .first("runningState").as("runningState")
    13. sort(new Sort(Sort.Direction.DESC, "readingTime"))
    14. )
    15. // 解决排序内存不足问题
    16. .withOptions(AggregationOptions.builder().allowDiskUse(true).build());
    17. // 获取分组后的数据
    18. AggregationResults groupResults
    19. = mongoTemplate.aggregate(aggregation, TcState.class, TcState.class);
    20. List dataList = groupResults.getMappedResults();
    21. return dataList;
    22. }

    1. private Criteria[] getTcStateCriteriaArray(TcStateDto tcStateDto) {
    2. // 定义一个存放条件的集合
    3. List criteriaList = new ArrayList<>();
    4. // 定义一个存放条件的数组(暂时不给长度)
    5. Criteria[] criteriaArray = {};
    6. if (null != tcStateDto) {
    7. if (null != tcStateDto.getCompanyId()) {
    8. criteriaList.add(Criteria.where("companyId").is(tcStateDto.getCompanyId()));
    9. }
    10. //根据温控终端多个查询时,无法使用索引
    11. if ((tcStateDto.getTempIdList() != null && tcStateDto.getTempIdList().size() > 0)) {
    12. Criteria c = Criteria.where("tempControllerId").in(tcStateDto.getTempIdList());
    13. criteriaList.add(c);
    14. }else {
    15. //根据温控终端多个查询时,就不需要终端分组的条件
    16. if ((tcStateDto.getIdList() != null && tcStateDto.getIdList().size() > 0)) {
    17. Criteria c = Criteria.where("tcGroupId").in(tcStateDto.getIdList());
    18. criteriaList.add(c);
    19. } else if (null != tcStateDto.getTcGroupId() && (tcStateDto.getIdList() == null || tcStateDto.getIdList().size() <= 0)) {
    20. criteriaList.add(Criteria.where("tcGroupId").is(tcStateDto.getTcGroupId()));
    21. } else if (null != tcStateDto.getTempControllerId() && (tcStateDto.getTempIdList() == null || tcStateDto.getTempIdList().size() <= 0)) {
    22. criteriaList.add(Criteria.where("tempControllerId").is(tcStateDto.getTempControllerId()));
    23. }
    24. }
    25. if (null != tcStateDto.getReadingTime()) {
    26. criteriaList.add(Criteria.where("readingTime").is(tcStateDto.getReadingTime()));
    27. }
    28. if (StringUtils.isNotBlank(tcStateDto.getStartDate()) && StringUtils.isBlank(tcStateDto.getEndDate())) {
    29. criteriaList.add(Criteria.where("readingTime").gte(DateUtil.stringToDate(tcStateDto.getStartDate())));
    30. }
    31. if (StringUtils.isNotBlank(tcStateDto.getEndDate()) && StringUtils.isBlank(tcStateDto.getStartDate())) {
    32. criteriaList.add(Criteria.where("readingTime").lte(DateUtil.stringToDate(tcStateDto.getEndDate())));
    33. }
    34. if (StringUtils.isNotBlank(tcStateDto.getStartDate()) && StringUtils.isNotBlank(tcStateDto.getEndDate())) {
    35. criteriaList.add(Criteria.where("readingTime").gte(DateUtil.stringToDate(tcStateDto.getStartDate()))
    36. .lte(DateUtil.stringToDate(tcStateDto.getEndDate())));
    37. }
    38. }
    39. // 如果有条件
    40. if (criteriaList.size() > 0) {
    41. // 集合的个数就是数组的长度
    42. criteriaArray = new Criteria[criteriaList.size()];
    43. // 遍历添加到数组中
    44. for (int i = 0; i < criteriaList.size(); i++) {
    45. criteriaArray[i] = criteriaList.get(i);
    46. }
    47. }
    48. return criteriaArray;
    49. }

  • 相关阅读:
    Leedcode 每日一题: 2760. 最长奇偶子数组
    计划跳槽需要做哪些准备?
    企业IP地址管理(IPAM)
    Vue实现刷新当前单个页面 适用于keep-alive【vue3适用】
    Python 考试练习题 3
    还不知道产品帮助中心怎样制作?,来看看这个吧
    关于 async 和 await 两个关键字(C#)【并发编程系列_5】
    自定义Dynamics 365实施和发布业务解决方案 - 1. 准备工作
    「Verilog学习笔记」多功能数据处理器
    java spring boot 数据库密码解密
  • 原文地址:https://blog.csdn.net/xue317378914/article/details/126704554