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

基于Elasticsearch的队列(Cohort)分析:高效数据提取查询方案求助

Elasticsearch 队列(Cohort)分析:实用查询与思路

我经常帮团队做ES里的队列分析,核心就是追踪特定用户群在时间维度上的行为变化,这里给你分享几个实用的查询示例和优化思路:

核心思路

队列分析的本质是先分组(按首次行为时间/属性),再追踪组内用户的后续行为。在Elasticsearch里,主要靠嵌套聚合来实现,常用的聚合组件包括:

  • date_histogram:按时间窗口(月/周/日)划分队列
  • terms/cardinality:统计用户数量
  • filter:筛选特定行为事件
  • bucket_script:计算留存率等衍生指标

示例1:按首次注册时间分队列,统计月度留存

假设你的事件数据结构包含user_id、event_type(比如register/login)、event_time,我们要统计「每个月注册的用户,在后续每个月的活跃人数」。

完整DSL查询

{
  "size": 0,
  "query": {
    "bool": {
      "filter": [
        {"terms": {"event_type": ["register", "login"]}},
        {"range": {"event_time": {"gte": "2024-01-01", "lte": "2024-06-30"}}}
      ]
    }
  },
  "aggs": {
    "users": {
      "terms": {
        "field": "user_id",
        "size": 10000 // 根据用户量调整,或者用composite聚合分页
      },
      "aggs": {
        "first_register_time": {
          "min": {
            "field": "event_time",
            "filter": {"term": {"event_type": "register"}}
          }
        },
        "cohort_month": {
          "date_histogram": {
            "field": "first_register_time",
            "calendar_interval": "month",
            "format": "yyyy-MM"
          },
          "aggs": {
            // 统计注册当月的登录用户
            "month_0_login": {
              "filter": {
                "bool": {
                  "must": [
                    {"term": {"event_type": "login"}},
                    {"script": {
                      "source": "doc['event_time'].value.getMonth() == doc['first_register_time'].value.getMonth()"
                    }}
                  ]
                }
              },
              "aggs": {"distinct_users": {"cardinality": {"field": "user_id"}}}
            },
            // 统计注册后第1个月的登录用户
            "month_1_login": {
              "filter": {
                "bool": {
                  "must": [
                    {"term": {"event_type": "login"}},
                    {"script": {
                      "source": "doc['event_time'].value.getMonth() == doc['first_register_time'].value.getMonth() + 1"
                    }}
                  ]
                }
              },
              "aggs": {"distinct_users": {"cardinality": {"field": "user_id"}}}
            },
            // 计算留存率(第1个月留存=month_1_login / 注册用户数)
            "month_1_retention": {
              "bucket_script": {
                "buckets_path": {
                  "login_users": "month_1_login.distinct_users",
                  "cohort_users": "_count"
                },
                "script": "params.login_users / params.cohort_users"
              }
            }
          }
        }
      }
    },
    // 按队列月份合并结果
    "cohort_summary": {
      "terms": {
        "field": "cohort_month.keyword"
      },
      "aggs": {
        "total_cohort_users": {"sum": {"field": "_count"}},
        "total_month_1_login": {"sum": {"field": "month_1_login.distinct_users.value"}},
        "avg_month_1_retention": {"avg": {"field": "month_1_retention.value"}}
      }
    }
  }
}

关键说明

  1. 先通过terms聚合拿到每个用户,再用min聚合获取他们的首次注册时间,这是划分队列的核心依据。
  2. 用date_histogram把用户按首次注册月份分到不同队列桶里。
  3. 每个队列桶内,用filter+脚本筛选出对应时间窗口的行为,再用cardinality统计活跃用户数。
  4. 最后用bucket_script计算留存率,直观展示队列的粘性。

示例2:简化版留存率计算(适合大用户量)

如果用户量很大,上面的terms聚合(size=10000)可能内存溢出,推荐用runtime fields提前计算用户的首次行为时间,再结合composite聚合分页:

步骤1:定义runtime field(计算首次注册时间)

{
  "mappings": {
    "runtime": {
      "first_register_time": {
        "type": "date",
        "script": {
          "source": """
            def events = doc['event_type'].stream().filter(v -> v.value == 'register').collect();
            if (events.size() > 0) {
              emit(doc['event_time'].value);
            }
          """
        }
      }
    }
  }
}

步骤2:用composite聚合做队列分析

{
  "size": 0,
  "query": {"term": {"event_type": "login"}},
  "aggs": {
    "cohort": {
      "composite": {
        "size": 1000,
        "sources": [
          {
            "cohort_month": {
              "date_histogram": {
                "field": "first_register_time",
                "calendar_interval": "month",
                "format": "yyyy-MM"
              }
            }
          },
          {
            "activity_month": {
              "date_histogram": {
                "field": "event_time",
                "calendar_interval": "month",
                "format": "yyyy-MM"
              }
            }
          }
        ]
      },
      "aggs": {
        "active_users": {"cardinality": {"field": "user_id"}}
      }
    }
  }
}

这个查询会输出每个队列(cohort_month)在每个活跃月份(activity_month)的用户数,你可以在客户端进一步计算留存率。


优化建议

  • 提前过滤数据:用bool.filter先筛选出需要的事件类型和时间范围,减少聚合的数据量。
  • 选择合适的时间间隔:用calendar_interval(比如month/week)比fixed_interval更符合业务逻辑,避免跨月/周的问题。
  • 避免超大size的terms聚合:用户量超过10w时,改用composite聚合分页,或者用data frame transforms提前预计算队列数据。
  • 利用索引别名:如果数据按天/月分片,用索引别名只查询相关时间段的索引,提升查询速度。

内容的提问来源于stack exchange,提问作者Nishant Dixit

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 10:39:10