第 18 关 · ★★★★

高级分析能力

Window、JOIN、Bitmap、HLL 与复杂类型:攒一个高级 SQL 模式库。

已点亮 · 最佳 分
高级分析能力 第 1 页

Day 15|Doris 高级分析:窗口函数、复杂类型、Bitmap、HLL、Join 与典型分析场景

Day 15 学习总览
Day 15 学习总览

前十四天,我们已经完成 Doris 的定位、架构、部署、表模型、分区分桶、导入、索引、物化视图、高并发点查、宽表与半结构化数据。今天开始把这些能力组合起来,解决真实分析工作中更复杂的问题。

一个经营分析页面通常会同时提出多种计算要求:查看各区域销售 TopN,计算最近七天移动平均,比较用户本次行为与上一次行为,统计新客次日留存,计算多渠道独立访客,展开用户标签,关联订单与会员等级,还要把每笔订单匹配到当时生效的价格快照。每项需求都可以写出 SQL,真正困难的部分在于选择合适的计算表达、数据结构和分布策略,并且让结果在大数据量、高并发环境下保持稳定。

Day 15 建立一套高级分析工具箱,包含六组核心能力:

  1. 用窗口函数完成组内排名、累计值、移动窗口和前后行比较;
  2. 用 ARRAY、MAP、STRUCT、JSON/VARIANT 保存并计算嵌套数据;
  3. 用 Bitmap 完成精确去重、留存和人群集合运算;
  4. 用 HLL 完成低成本、可合并的近似去重;
  5. 用正确的 Join 语义和数据分布策略控制网络、内存与数据倾斜;
  6. 用留存、漏斗、TopN、RFM、时序最近匹配等案例形成完整分析闭环。

这一天的目标是让 SQL 从“能够执行”提升到“语义明确、结果可靠、成本可控、便于验证”。

高级分析能力地图
高级分析能力地图


一、先建立高级分析的选择坐标系

高级分析需求可以先拆成四个问题。

1. 结果是否需要保留明细行

只需要每个城市一行汇总结果时,使用 GROUP BY。每个订单仍需保留,同时增加城市排名、累计金额或上一笔订单金额时,使用窗口函数。

2. 计算对象是标量、嵌套结构还是集合

金额、时间、状态等字段使用普通类型。标签列表适合 ARRAY,动态键值适合 MAP,稳定的嵌套对象适合 STRUCT,字段变化频繁的文档适合 JSON 或 VARIANT。用户集合的交并差、精确 UV 和留存适合 Bitmap,大规模趋势型 UV 可以使用 HLL。

3. 多表关联需要什么结果语义

需要两侧匹配列时使用 INNER JOIN;要保留左表未匹配行时使用 LEFT JOIN;只关心“是否存在”时使用 SEMI JOIN;寻找缺失对象时使用 ANTI JOIN;按时间寻找最近快照时,Doris 4.1+ 可以使用 ASOF JOIN。

4. 性能主要消耗在哪里

窗口函数常见成本来自分区内排序和状态内存;Join 成本来自 Hash Table、Shuffle 和倾斜;复杂类型展开会放大行数;Bitmap/HLL 会形成可合并的聚合状态。设计阶段先判断数据量、选择性、分组粒度和并发,后续调优才能有明确方向。

GROUP BY 与窗口函数
GROUP BY 与窗口函数


二、窗口函数:保留明细的组内计算

Apache Doris 当前窗口函数体系包含专用分析函数和大部分可用于 OVER 的聚合函数。它们在 JOINWHEREGROUP BY 等阶段之后执行,对一个窗口范围进行计算,并为结果集中的每一行生成值。

基本语法如下:

<function>(<arguments>) OVER (
    PARTITION BY <partition_columns>
    ORDER BY <order_columns>
    <window_frame>
)

2.1 PARTITION BY:定义独立计算空间

PARTITION BY user_id 表示每个用户拥有独立窗口。缺少分区条件时,整个结果集会进入同一个窗口。大结果集上的全局排序会产生明显的内存和 Spill 压力。

2.2 ORDER BY:定义组内顺序

排名、累计、LAG、LEAD 都依赖组内顺序。排序字段无法形成唯一顺序时,同值记录之间的先后关系可能变化。生产 SQL 应补充稳定的唯一字段:

ORDER BY event_time, event_id

这条规则对留存、会话切分、首末事件识别十分重要。只用秒级时间字段排序,而同一秒存在多条事件时,结果可能缺乏确定性。

2.3 FRAME:定义当前行看到的范围

常见写法包括:

ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW

用于累计值;

ROWS BETWEEN 6 PRECEDING AND CURRENT ROW

用于七条记录移动窗口;

ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING

用于当前位置附近三行计算。

Doris 支持 ROWSRANGE。涉及固定前后行数时,显式使用 ROWS。带 ORDER BY 且省略 Frame 时,默认窗口通常会扩展到当前排序值,重复排序值可能一起进入 RANGE。对于金额累计、移动平均等核心口径,建议把 Frame 写完整。

窗口函数三个坐标
窗口函数三个坐标

2.4 排名函数怎样选择

函数 同值处理 适用场景
ROW_NUMBER() 每行唯一序号 每组取一条、TopN 明细、去重保留最新记录
RANK() 同值同名次,后续名次跳号 竞赛排名、并列名次保留空位
DENSE_RANK() 同值同名次,后续连续 价格档位、等级排名
NTILE(n) 分成近似等量的 n 组 用户分层、四分位、实验分桶
CUME_DIST() 返回累计分布比例 百分位位置、分布分析

每个品类取 GMV 最高的十个商品,可以用 CTE 配合 ROW_NUMBER()

WITH ranked AS (
    SELECT
        category_id,
        product_id,
        SUM(pay_amount) AS gmv,
        ROW_NUMBER() OVER (
            PARTITION BY category_id
            ORDER BY SUM(pay_amount) DESC, product_id
        ) AS rn
    FROM fact_order
    WHERE order_date BETWEEN '2026-08-01' AND '2026-08-31'
    GROUP BY category_id, product_id
)
SELECT category_id, product_id, gmv
FROM ranked
WHERE rn <= 10;

2.5 累计、移动与前后行比较

累计 GMV:

SELECT
    order_date,
    daily_gmv,
    SUM(daily_gmv) OVER (
        ORDER BY order_date
        ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
    ) AS cumulative_gmv
FROM agg_order_day;

七日移动平均:

AVG(daily_gmv) OVER (
    ORDER BY order_date
    ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
)

与上一日比较:

LAG(daily_gmv, 1, 0) OVER (ORDER BY order_date)

下一次事件时间:

LEAD(event_time) OVER (
    PARTITION BY user_id
    ORDER BY event_time, event_id
)

窗口函数经常需要排序。查询数据已经按窗口分区键和排序键组织时,系统更容易获得稳定并行度;单个超大用户、超大设备或空值热点会形成超大分区,导致某个执行实例承担大量数据。Profile 中要重点观察 Sort、Window、PeakMemory、Spill 和各 Instance 的行数差异。

2.6 FIRST_VALUE、LAST_VALUE 与默认窗口边界

FIRST_VALUELAST_VALUE 看起来很直观,实际使用时经常被默认 Frame 影响。带 ORDER BY 的窗口如果没有显式声明结束边界,当前行通常只能看到从分区起点到当前排序值的范围。此时 LAST_VALUE 返回的可能是“当前窗口最后一条”,并非整个分区最终一条。

获取用户整段历史中的最后一次渠道,可以明确写成:

LAST_VALUE(channel) OVER (
    PARTITION BY user_id
    ORDER BY event_time, event_id
    ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
)

这段 SQL 会让每一行都看到整个用户分区。窗口扩大后,状态维护和排序成本也会增加。只需要每个用户最终一条记录时,通常使用 ROW_NUMBER() ... ORDER BY event_time DESC 后过滤 rn = 1,结果行数更小,语义也更清楚。

2.7 用窗口函数完成去重、间隔和会话切分

CDC、埋点和日志系统可能重复写入同一业务事件。可以按业务唯一键分区,再按版本或接收时间倒序编号:

WITH deduplicated AS (
    SELECT
        *,
        ROW_NUMBER() OVER (
            PARTITION BY event_id
            ORDER BY source_version DESC, ingest_time DESC
        ) AS rn
    FROM raw_event
)
SELECT *
FROM deduplicated
WHERE rn = 1;

这类 SQL 适合离线修复、质量核验和临时回溯。长期持续去重仍要优先依赖上游幂等、Unique Key 和 Sequence Column,避免每次查询都重新排序全量历史数据。

会话切分可以通过 LAG 计算相邻事件间隔。当间隔超过三十分钟时标记新会话,再对标记做累计求和:

WITH gaps AS (
    SELECT
        user_id,
        event_time,
        event_id,
        CASE
            WHEN LAG(event_time) OVER (
                PARTITION BY user_id
                ORDER BY event_time, event_id
            ) IS NULL THEN 1
            WHEN TIMESTAMPDIFF(
                MINUTE,
                LAG(event_time) OVER (
                    PARTITION BY user_id
                    ORDER BY event_time, event_id
                ),
                event_time
            ) > 30 THEN 1
            ELSE 0
        END AS new_session
    FROM fact_user_event
)
SELECT
    *,
    SUM(new_session) OVER (
        PARTITION BY user_id
        ORDER BY event_time, event_id
        ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
    ) AS session_no
FROM gaps;

生产 SQL 可以先在一层 CTE 中计算一次 LAG,减少重复表达式。会话规则还要明确跨天处理、客户端离线补报、时区和乱序容忍范围。


三、复杂类型:让嵌套业务对象保留原有结构

Doris 支持 ARRAY、MAP、STRUCT、JSON 和 VARIANT。复杂类型适合日志、埋点、用户画像、订单明细和配置数据。选择时可以使用一条简单原则:结构稳定且字段有明确含义时使用普通列或 STRUCT;多值有序集合使用 ARRAY;动态键值使用 MAP;字段路径持续变化时使用 JSON/VARIANT。

3.1 ARRAY:标签、品类和多值属性

CREATE TABLE user_profile (
    user_id BIGINT,
    tags ARRAY<STRING>,
    scores ARRAY<INT>
)
DUPLICATE KEY(user_id)
DISTRIBUTED BY HASH(user_id) BUCKETS 16;

常用计算:

array_contains(tags, 'vip')
array_join(tags, ',')
array_map(x -> x * 2, scores)
array_filter(x -> x >= 80, scores)

ARRAY 适合保留一行内的多值集合。需要按标签聚合时,可以展开为多行。

3.2 MAP:动态属性和配置项

attrs MAP<STRING, STRING>

常用访问方式:

attrs['os']
map_keys(attrs)
map_values(attrs)
element_at(attrs, 'channel')

MAP 的键和值类型固定,适合字段名动态、值类型一致的场景。值类型差异很大时,VARIANT 的兼容性更好。

3.3 STRUCT:稳定的嵌套对象

device STRUCT<os:STRING, brand:STRING, version:STRING>

读取:

device.os

STRUCT 对字段名称和类型进行明确约束,适合设备、地址、坐标等稳定复合对象。

3.4 LATERAL VIEW 与 UNNEST

Doris 4.0 可以通过 LATERAL VIEW EXPLODE 展开数组:

SELECT user_id, tag
FROM user_profile
LATERAL VIEW EXPLODE(tags) t AS tag;

Doris 4.1 引入了更接近 PostgreSQL、Trino 习惯的 UNNEST

SELECT user_id, tag
FROM user_profile, UNNEST(tags) AS t(tag);

展开操作会放大行数。输入一亿行、平均每行二十个标签,理论输出可达到二十亿行。优化顺序应当是:先用分区、普通列和数组函数过滤,再展开;对异常超长数组设置质量规则;只选择后续真正需要的列。

复杂类型与展开
复杂类型与展开


四、Join:语义正确之后,再控制分布式成本

Join 首先是一项结果集合操作。Doris 当前支持 INNER、LEFT、RIGHT、FULL、CROSS、LEFT/RIGHT SEMI、LEFT/RIGHT ANTI、NULL AWARE LEFT ANTI 等类型。

Join 语义矩阵
Join 语义矩阵

4.1 SEMI 与 ANTI 的价值

只需要筛选“购买过商品 A 的用户”,结果不需要商品表字段时,SEMI JOIN 更贴合语义:

SELECT u.user_id
FROM dim_user u
LEFT SEMI JOIN fact_order o
    ON u.user_id = o.user_id
WHERE o.product_id = 1001;

寻找三十天未复购用户可以使用 ANTI JOIN。相比先做普通 JOIN 再 DISTINCT,SEMI/ANTI 能减少不必要的右表列输出和结果膨胀。

4.2 Hash Join 与 Nest Loop Join

等值条件通常使用 Hash Join。右侧数据构建 Hash Table,左侧数据作为 Probe 流输入。右侧过大时,内存消耗会快速增加。

非等值条件、笛卡尔积等场景可能使用 Nest Loop Join,计算量通常更高。Cross Join 的结果规模等于两侧行数乘积,遗漏 ON 条件会在很短时间内放大数据量。

4.3 四种 Hash Join 分布策略

Doris 官方将分布策略分为 Broadcast、Partition Shuffle、Bucket Shuffle 和 Colocate:

策略 网络量 典型条件
Broadcast N × T(R) 右表较小,可以复制到每个执行节点
Partition Shuffle T(S) + T(R) 两侧都按 Join Key 重分布
Bucket Shuffle T(R) 一侧已按 Join Key 分桶,只移动另一侧
Colocate 0 两表同分桶键、同桶数、同 Colocate Group,副本分布稳定

Join 数据分布策略
Join 数据分布策略

Broadcast 对小维表非常有效,维表增长后会增加网络和每个节点的 Hash Table 内存。Bucket Shuffle 对物理表有分区等约束,直接扫描分区表时通常需要通过分区裁剪收敛到单个分区。Colocate 适合固定、高频的大表关联,扩容、均衡和副本修复会增加运维约束;Colocate Group 处于不稳定状态时,查询可能回退到常规 Join。

4.4 Runtime Filter

Hash Join 的 Build 侧可以生成 IN、MinMax 或 Bloom 等 Runtime Filter,并发送到 Probe 侧 Scan。大量不匹配数据可以在读取阶段被过滤,从而减少扫描、网络和 Join 计算。

验证时同时查看 Explain 和 Profile:

  • Join 节点是否生成 Runtime Filter;
  • Scan 节点是否接收并应用;
  • 等待时间是否过长;
  • 被过滤行数是否足够大;
  • 最终 Rows、Bytes 和 Join Time 是否下降。

4.5 Join 倾斜

Join Key 中存在大量 NULL、空字符串、默认值或超级热点 ID 时,某个分区会远大于其他分区。典型处理方法包括:

  • 过滤没有业务含义的 NULL 和空值;
  • 将热点值单独分支计算;
  • 为热点键增加 Salt,再在下游聚合;
  • 重新评估分桶键和维度模型;
  • 收集准确统计信息,让 CBO 选择合理 Join 顺序和分布方式。

4.6 NULL、类型转换与 Join 口径

Join Key 两侧类型不同会触发隐式转换,可能影响 Runtime Filter、Hash 计算和结果精度。生产建模应让事实表与维表的键使用相同类型、相同编码和相同业务字典。字符串数字与整数键混用时,建议在数据接入阶段完成标准化。

NULL = NULL 的结果并不成立。需要把 NULL 视为同一业务值时,应使用 Null-safe Equal <=>,并确认该语义符合业务口径。大量 NULL 参与分布式 Hash 会形成热点,很多场景更适合先过滤无效键,再单独统计缺失数据。

外连接的过滤条件位置也会改变结果。LEFT JOIN 右表条件写在 WHERE 中,可能把右表 NULL 行过滤掉,最终效果接近 INNER JOIN。需要保留左表全部记录时,右表约束通常放在 ON 条件中:

SELECT u.user_id, o.order_id
FROM dim_user u
LEFT JOIN fact_order o
  ON u.user_id = o.user_id
 AND o.order_status = 'paid';

4.7 Hint 的使用边界

Nereids CBO 会依据统计信息选择 Join 顺序和分布策略。统计信息过期、数据分布高度特殊或 POC 需要验证假设时,可以使用 Broadcast、Shuffle、Leading 等 Hint 做对照实验。Hint 应当记录使用原因、目标版本和回归条件。

长期依赖 Hint 会增加数据增长和版本升级后的维护成本。更稳妥的治理方法包括及时 ANALYZE TABLE、修正表模型、调整分区分桶、拆分热点键、建设物化视图。Hint 适合作为诊断和临时控制手段,最终方案仍需用 Profile 证明。


登录后可阅读本文完整内容。