Learn
ClickHouse/14-window-functions

窗口函数与漏斗留存

窗口函数从 ClickHouse 21.x 起完整支持,语法和 MySQL 8 / PostgreSQL 基本一致。但 ClickHouse 还有一批专为漏斗和留存设计的函数,比通用窗口函数快得多。

1. 窗口函数基础

1.1 语法结构

函数名(参数) OVER (
    PARTITION BY 分区列
    ORDER BY 排序列
    ROWS|RANGE BETWEEN 起点 AND 终点
)

和 GROUP BY 的区别:

GROUP BY                       窗口函数
─────────────────────────      ────────────────────────────
1000 行 ──▶ 5 行               1000 行 ──▶ 1000 行
(聚合后行数减少)              (每行都保留,额外加一列)
 
country  pv                    country  page   pv  country_pv
CN       200                   CN       /a     10  200
US       180                   CN       /b     15  200
                               US       /a     8   180

1.2 排名函数

SELECT
    country,
    page,
    count()                                                              AS pv,
    row_number() OVER (PARTITION BY country ORDER BY count() DESC)       AS rn,
    rank()       OVER (PARTITION BY country ORDER BY count() DESC)       AS rnk,
    dense_rank() OVER (PARTITION BY country ORDER BY count() DESC)       AS drnk
FROM events
WHERE country IN ('CN', 'US')
GROUP BY country, page
ORDER BY country, pv DESC
LIMIT 6;
┌─country─┬─page───┬──pv─┬─rn─┬─rnk─┬─drnk─┐
│ CN      │ /p/183 │ 442 │  1 │   1 │    1 │
│ CN      │ /p/427 │ 438 │  2 │   2 │    2 │
│ CN      │ /p/91  │ 438 │  3 │   2 │    2 │
│ US      │ /p/12  │ 441 │  1 │   1 │    1 │
│ US      │ /p/338 │ 435 │  2 │   2 │    2 │
│ US      │ /p/220 │ 431 │  3 │   3 │    3 │
└─────────┴────────┴─────┴────┴─────┴──────┘

三者区别:值相同时 row_number 强行给不同序号,rank 给相同序号但跳号(1,2,2,4),dense_rank 给相同序号且不跳号(1,2,2,3)。

💡取 Top N 优先用 LIMIT BY
-- 窗口函数写法(需要子查询过滤 rn)
SELECT * FROM (
    SELECT country, page, count() AS pv,
           row_number() OVER (PARTITION BY country ORDER BY count() DESC) AS rn
    FROM events GROUP BY country, page
) WHERE rn <= 3;
 
-- LIMIT BY 写法(更简洁、更快)
SELECT country, page, count() AS pv
FROM events GROUP BY country, page
ORDER BY country, pv DESC
LIMIT 3 BY country;

LIMIT BY 不需要为所有行计算序号,只维护每组的 Top N,性能更好。

1.3 偏移函数:lag / lead

ClickHouse 用 lagInFrame / leadInFrame(标准 lag/lead 也支持但语义略有差异):

-- 每天 PV 及环比
SELECT
    day,
    pv,
    lagInFrame(pv)  OVER (ORDER BY day)                     AS prev_pv,
    leadInFrame(pv) OVER (ORDER BY day)                     AS next_pv,
    pv - lagInFrame(pv) OVER (ORDER BY day)                 AS diff,
    round((pv - lagInFrame(pv) OVER (ORDER BY day))
          / lagInFrame(pv) OVER (ORDER BY day) * 100, 2)    AS growth_pct
FROM (
    SELECT toDate(event_time) AS day, count() AS pv
    FROM events GROUP BY day ORDER BY day
)
ORDER BY day
LIMIT 5;
┌────────day─┬────pv─┬─prev_pv─┬─next_pv─┬─diff─┬─growth_pct─┐
│ 2024-06-01 │ 33452 │       0 │   33280 │ 33452│          0 │
│ 2024-06-02 │ 33280 │   33452 │   33391 │ -172 │      -0.51 │
│ 2024-06-03 │ 33391 │   33280 │   33188 │  111 │       0.33 │
│ 2024-06-04 │ 33188 │   33391 │   33507 │ -203 │      -0.61 │
│ 2024-06-05 │ 33507 │   33188 │   33612 │  319 │       0.96 │
└────────────┴───────┴─────────┴─────────┴──────┴────────────┘
⚠️lagInFrame 受窗口帧限制

lagInFrame 只能看到当前窗口帧内的行。默认帧是 RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW,所以 leadInFrame 默认永远返回 0(后面的行不在帧内)。

要用 leadInFrame 必须显式扩大帧:

leadInFrame(pv) OVER (ORDER BY day ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING)

上面的例子里 next_pv 能出结果,是因为 ClickHouse 24.x 对单纯的 leadInFrame 做了帧自动扩展。但显式写清楚更稳妥。

1.4 窗口帧

-- 7 日滑动平均(最常用的运维指标平滑)
SELECT
    day,
    pv,
    round(avg(pv) OVER (ORDER BY day ROWS BETWEEN 6 PRECEDING AND CURRENT ROW), 1) AS ma7,
    sum(pv) OVER (ORDER BY day ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)   AS cumulative,
    max(pv) OVER (ORDER BY day ROWS BETWEEN 3 PRECEDING AND 3 FOLLOWING)           AS local_max
FROM (
    SELECT toDate(event_time) AS day, count() AS pv FROM events GROUP BY day
)
ORDER BY day
LIMIT 10;
┌────────day─┬────pv─┬─────ma7─┬─cumulative─┬─local_max─┐
│ 2024-06-01 │ 33452 │ 33452.0 │      33452 │     33612 │
│ 2024-06-02 │ 33280 │ 33366.0 │      66732 │     33612 │
│ 2024-06-03 │ 33391 │ 33374.3 │     100123 │     33701 │
│ 2024-06-04 │ 33188 │ 33327.8 │     133311 │     33701 │
│ 2024-06-05 │ 33507 │ 33363.6 │     166818 │     33701 │
│ 2024-06-06 │ 33612 │ 33405.0 │     200430 │     33701 │
│ 2024-06-07 │ 33701 │ 33447.3 │     234131 │     33712 │
│ 2024-06-08 │ 33389 │ 33438.3 │     267520 │     33712 │
└────────────┴───────┴─────────┴────────────┴───────────┘

ROWS 与 RANGE 的区别:

含义例子
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW物理行数前 6 行 + 当前行
RANGE BETWEEN 6 PRECEDING AND CURRENT ROW排序列的值范围排序值在 [当前值-6, 当前值] 内的所有行

有数据缺失的日期时,ROWS 会包含更早的日期,RANGE 则严格按日期算。做时间序列分析要想清楚用哪个。

1.5 命名窗口

多个函数共用同一个窗口定义时:

SELECT
    day,
    pv,
    avg(pv) OVER w7   AS ma7,
    max(pv) OVER w7   AS max7,
    min(pv) OVER w7   AS min7
FROM (SELECT toDate(event_time) AS day, count() AS pv FROM events GROUP BY day)
WINDOW w7 AS (ORDER BY day ROWS BETWEEN 6 PRECEDING AND CURRENT ROW)
ORDER BY day LIMIT 5;

2. 窗口函数 vs GROUP BY 的性能取舍

⚠️窗口函数在 ClickHouse 里相对昂贵

窗口函数需要:

  1. 按 PARTITION BY 分区(哈希或排序)
  2. 每个分区内按 ORDER BY 排序
  3. 逐行滑动计算

这个过程并行度低于 GROUP BY,且内存占用与分区大小成正比。对 10 亿行直接开窗很容易 OOM。

优化策略:

  • 先聚合再开窗。上面的例子都是先 GROUP BY day 把 1 亿行压到 30 行,再开窗。这是标准姿势。
  • 能用 LIMIT BY 就别用 row_number
  • 能用 argMax / groupArray + 数组函数就别用窗口函数

对比示例:

-- 慢:直接在明细上开窗(1 亿行排序)
SELECT user_id, event_time,
       row_number() OVER (PARTITION BY user_id ORDER BY event_time) AS step
FROM events;
-- Elapsed: 42.1 sec, Memory: 8.2 GiB
 
-- 快:用数组代替
SELECT user_id, arrayEnumerate(groupArray(event_time)) AS steps
FROM (SELECT user_id, event_time FROM events ORDER BY user_id, event_time)
GROUP BY user_id;
-- Elapsed: 6.8 sec, Memory: 2.1 GiB

3. windowFunnel:漏斗分析

这是 ClickHouse 的杀手锏函数。

3.1 语法

windowFunnel(时间窗口秒数[, 模式])(时间列, 条件1, 条件2, ..., 条件N)

返回每个用户按顺序完成了前几步(0 到 N)。

-- 分析「浏览 → 点击 → 购买」三步漏斗,时间窗口 1 小时
SELECT
    level,
    count() AS users
FROM (
    SELECT
        user_id,
        windowFunnel(3600)(
            event_time,
            event_type = 'view',
            event_type = 'click',
            event_type = 'purchase'
        ) AS level
    FROM events
    WHERE event_time >= '2024-06-01' AND event_time < '2024-06-08'
    GROUP BY user_id
)
GROUP BY level
ORDER BY level DESC;
┌─level─┬─users─┐
│     3 │ 18042 │   ← 完成全部三步
│     2 │  9871 │   ← 只完成前两步
│     1 │  7203 │   ← 只完成第一步
│     0 │  4884 │   ← 一步都没完成
└───────┴───────┘

3.2 转化为标准漏斗报表

SELECT
    '1. 浏览'   AS step, countIf(level >= 1) AS users FROM funnel_data
UNION ALL
SELECT
    '2. 点击',  countIf(level >= 2) FROM funnel_data
UNION ALL
SELECT
    '3. 购买',  countIf(level >= 3) FROM funnel_data;

更优雅的写法,用数组:

WITH funnel AS (
    SELECT
        user_id,
        windowFunnel(3600)(event_time,
            event_type = 'view',
            event_type = 'click',
            event_type = 'purchase'
        ) AS level
    FROM events
    WHERE event_time >= '2024-06-01' AND event_time < '2024-06-08'
    GROUP BY user_id
)
SELECT
    step_name,
    users,
    round(users / max(users) OVER () * 100, 2)                        AS total_conv_pct,
    round(users / lagInFrame(users) OVER (ORDER BY step_idx) * 100, 2) AS step_conv_pct
FROM (
    SELECT 1 AS step_idx, '浏览' AS step_name, countIf(level >= 1) AS users FROM funnel
    UNION ALL
    SELECT 2, '点击',   countIf(level >= 2) FROM funnel
    UNION ALL
    SELECT 3, '购买',   countIf(level >= 3) FROM funnel
)
ORDER BY step_idx;
┌─step_name─┬─users─┬─total_conv_pct─┬─step_conv_pct─┐
│ 浏览      │ 35116 │         100.00 │             0 │
│ 点击      │ 27913 │          79.49 │         79.49 │
│ 购买      │ 18042 │          51.38 │         64.64 │
└───────────┴───────┴────────────────┴───────────────┘

3.3 模式参数

模式含义
默认严格按顺序,中间可以有其他事件
'strict_order'严格连续,中间不能有任何其他事件
'strict_deduplication'出现重复事件则中断
'strict_increase'要求时间戳严格递增(同一秒的多个事件只算一个)
-- 严格连续:view 后必须直接是 click,不能有其他事件
windowFunnel(3600, 'strict_order')(event_time, ...)
💡windowFunnel 的时间窗口从第一步开始算

windowFunnel(3600) 表示:从用户完成第一步的时刻起,1 小时内完成的后续步骤才算数。

不是「每两步之间 1 小时」。如果业务上要求「加购后 24 小时内支付」,窗口要设成整个流程的总时长。

3.4 多维漏斗

SELECT
    country,
    device,
    countIf(level >= 1) AS s1_view,
    countIf(level >= 2) AS s2_click,
    countIf(level >= 3) AS s3_purchase,
    round(countIf(level >= 3) / countIf(level >= 1) * 100, 2) AS conv_pct
FROM (
    SELECT
        user_id,
        any(country) AS country,
        any(device)  AS device,
        windowFunnel(3600)(event_time,
            event_type = 'view', event_type = 'click', event_type = 'purchase') AS level
    FROM events
    WHERE event_time >= '2024-06-01' AND event_time < '2024-06-08'
    GROUP BY user_id
)
GROUP BY country, device
ORDER BY conv_pct DESC
LIMIT 5;
┌─country─┬─device──┬─s1_view─┬─s2_click─┬─s3_purchase─┬─conv_pct─┐
│ JP      │ ios     │    2381 │     1902 │        1244 │    52.25 │
│ CN      │ ios     │    2374 │     1888 │        1236 │    52.06 │
│ US      │ web     │    2361 │     1871 │        1221 │    51.71 │
└─────────┴─────────┴─────────┴──────────┴─────────────┴──────────┘

4. retention:留存分析

retention(条件1, 条件2, ..., 条件N)

返回一个 0/1 数组:第一个元素是条件 1 是否成立,后续元素是「条件 1 且 条件 i 都成立」。

-- 7 日留存
SELECT
    sum(r[1]) AS day0,
    sum(r[2]) AS day1,
    sum(r[3]) AS day2,
    sum(r[4]) AS day3,
    sum(r[5]) AS day7,
    round(sum(r[2]) / sum(r[1]) * 100, 2) AS d1_pct,
    round(sum(r[3]) / sum(r[1]) * 100, 2) AS d2_pct,
    round(sum(r[4]) / sum(r[1]) * 100, 2) AS d3_pct,
    round(sum(r[5]) / sum(r[1]) * 100, 2) AS d7_pct
FROM (
    SELECT
        user_id,
        retention(
            toDate(event_time) = '2024-06-01',
            toDate(event_time) = '2024-06-02',
            toDate(event_time) = '2024-06-03',
            toDate(event_time) = '2024-06-04',
            toDate(event_time) = '2024-06-08'
        ) AS r
    FROM events
    WHERE event_time >= '2024-06-01' AND event_time < '2024-06-09'
    GROUP BY user_id
);
┌─day0──┬─day1──┬─day2──┬─day3──┬─day7─┬─d1_pct─┬─d2_pct─┬─d3_pct─┬─d7_pct─┐
│ 22104 │ 12833 │  9421 │  7688 │ 5102 │  58.06 │  42.62 │  34.78 │  23.08 │
└───────┴───────┴───────┴───────┴──────┴────────┴────────┴────────┴────────┘

4.1 完整留存矩阵

生成运营最爱看的三角矩阵:

WITH
    toDate('2024-06-01') AS start_date,
    7                    AS days
SELECT
    cohort_day,
    cohort_size,
    arrayMap(i -> round(retained[i] / cohort_size * 100, 1), range(1, days + 1)) AS retention_pct
FROM (
    SELECT
        cohort_day,
        count()                                    AS cohort_size,
        arrayMap(d -> countIf(has(active_days, d)), range(0, days)) AS retained
    FROM (
        SELECT
            user_id,
            min(toDate(event_time))                                          AS cohort_day,
            groupUniqArray(dateDiff('day', min(toDate(event_time)) OVER (PARTITION BY user_id), toDate(event_time))) AS active_days
        FROM events
        WHERE event_time >= start_date
        GROUP BY user_id
    )
    GROUP BY cohort_day
)
ORDER BY cohort_day;

实际生产中更简单直接的写法是用 bitmap(第 12 章讲过):

WITH daily_users AS (
    SELECT
        toDate(event_time) AS day,
        groupBitmapState(toUInt32(user_id)) AS users
    FROM events
    WHERE event_time >= '2024-06-01' AND event_time < '2024-06-08'
    GROUP BY day
)
SELECT
    a.day                                            AS cohort_day,
    dateDiff('day', a.day, b.day)                    AS day_offset,
    bitmapCardinality(a.users)                       AS cohort_size,
    bitmapAndCardinality(a.users, b.users)           AS retained,
    round(bitmapAndCardinality(a.users, b.users) / bitmapCardinality(a.users) * 100, 2) AS pct
FROM daily_users AS a
CROSS JOIN daily_users AS b
WHERE b.day >= a.day
ORDER BY cohort_day, day_offset
LIMIT 10;
┌─cohort_day─┬─day_offset─┬─cohort_size─┬─retained─┬────pct─┐
│ 2024-06-01 │          0 │       22104 │    22104 │ 100.00 │
│ 2024-06-01 │          1 │       22104 │    12833 │  58.06 │
│ 2024-06-01 │          2 │       22104 │     9421 │  42.62 │
│ 2024-06-01 │          3 │       22104 │     7688 │  34.78 │
│ 2024-06-02 │          0 │       22019 │    22019 │ 100.00 │
│ 2024-06-02 │          1 │       22019 │    12780 │  58.04 │
└────────────┴────────────┴─────────────┴──────────┴────────┘

bitmap 方式的优势:日级 bitmap 可以预计算存进 AggregatingMergeTree,之后任意日期组合的留存都是毫秒级的位运算。这是生产环境的标准方案。

5. sequenceMatch 与 sequenceCount

比 windowFunnel 更灵活的序列模式匹配:

-- 匹配「view 后接 purchase」的模式
SELECT
    user_id,
    sequenceMatch('(?1)(?2)')(event_time,
        event_type = 'view',
        event_type = 'purchase') AS matched,
    sequenceCount('(?1)(?2)')(event_time,
        event_type = 'view',
        event_type = 'purchase') AS match_count
FROM events
GROUP BY user_id
HAVING matched = 1
LIMIT 5;

模式语法:

语法含义
(?1)第 1 个条件
.*任意事件
(?t>3600)两个事件间隔大于 3600 秒
(?t<=60)间隔不超过 60 秒
-- view 后 60 秒内 purchase(快速决策用户)
sequenceMatch('(?1)(?t<=60)(?2)')(event_time,
    event_type = 'view', event_type = 'purchase')
 
-- view 后经过任意事件,最终 purchase
sequenceMatch('(?1).*(?2)')(event_time,
    event_type = 'view', event_type = 'purchase')

6. 综合实战:转化分析看板

WITH
    '2024-06-01' AS d_from,
    '2024-06-08' AS d_to,
    user_funnel AS (
        SELECT
            user_id,
            any(country) AS country,
            any(device)  AS device,
            windowFunnel(86400)(event_time,
                event_type = 'view',
                event_type = 'click',
                event_type = 'purchase'
            ) AS level,
            count()                AS total_events,
            sum(duration_ms)       AS total_ms
        FROM events
        WHERE event_time >= d_from AND event_time < d_to
        GROUP BY user_id
    )
SELECT
    country,
    count()                                                 AS users,
    countIf(level >= 1)                                     AS viewed,
    countIf(level >= 2)                                     AS clicked,
    countIf(level >= 3)                                     AS purchased,
    round(countIf(level >= 2) / countIf(level >= 1) * 100, 2) AS view_to_click,
    round(countIf(level >= 3) / countIf(level >= 2) * 100, 2) AS click_to_buy,
    round(countIf(level >= 3) / count() * 100, 2)           AS overall_conv,
    round(avg(total_events), 1)                             AS avg_events,
    round(avgIf(total_ms, level >= 3) / 1000, 1)            AS buyer_avg_sec,
    bar(countIf(level >= 3), 0, 4000, 20)                   AS chart
FROM user_funnel
GROUP BY country
ORDER BY overall_conv DESC;
┌─country─┬─users─┬─viewed─┬─clicked─┬─purchased─┬─view_to_click─┬─click_to_buy─┬─overall_conv─┬─avg_events─┬─buyer_avg_sec─┬─chart────────────────┐
│ JP      │  7942 │   7031 │    5584 │      3612 │         79.42 │        64.69 │        45.48 │        6.2 │          18.4 │ ██████████████████   │
│ CN      │  7918 │   7009 │    5561 │      3594 │         79.34 │        64.63 │        45.39 │        6.2 │          18.3 │ █████████████████    │
│ US      │  7901 │   6988 │    5542 │      3578 │         79.31 │        64.56 │        45.29 │        6.1 │          18.5 │ █████████████████    │
└─────────┴───────┴────────┴─────────┴───────────┴───────────────┴──────────────┴──────────────┴────────────┴───────────────┴──────────────────────┘
🎯练习
  1. 用窗口函数计算每天 PV 的 7 日滑动平均和累计值,对比 ROWS 和 RANGE 两种帧定义在有日期缺失时的差异。
  2. 用 windowFunnel(3600) 做「view → click → purchase」漏斗,然后把窗口改成 60 秒,观察各步骤人数变化,理解时间窗口对漏斗的影响。
  3. 分别用 strict_order 和默认模式跑同一个漏斗,解释两者数字差异的原因。
  4. 用 groupBitmapState 预计算每日活跃用户 bitmap,存进一张 AggregatingMergeTree 表,然后基于它算出完整的 7 日留存矩阵。对比直接从明细算的耗时。
  5. 用 sequenceMatch('(?1)(?t<=300)(?2)') 找出「浏览后 5 分钟内购买」的快速转化用户,统计他们占全部购买用户的比例。

小结

  • 窗口函数语法与标准 SQL 一致,row_number / rank / dense_rank 做排名,lagInFrame / leadInFrame 做偏移
  • ROWS 按物理行数,RANGE 按排序列的值范围,有数据缺失时结果不同
  • 窗口函数在 ClickHouse 里并行度低、内存占用高,永远先聚合再开窗
  • 取 Top N 优先用 LIMIT n BY,比 row_number 快
  • windowFunnel(窗口秒数)(时间列, 条件...) 返回每个用户完成的步数,是漏斗分析的标准工具
  • 时间窗口从第一步开始算,不是步骤之间
  • retention 返回 0/1 数组做留存,但生产环境更推荐用 groupBitmap 预计算 + 位运算
  • sequenceMatch / sequenceCount 支持带时间约束的复杂序列模式
  • 下一章讲 JOIN 与字典,解决多表关联问题 →