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

前十四天,我们已经完成 Doris 的定位、架构、部署、表模型、分区分桶、导入、索引、物化视图、高并发点查、宽表与半结构化数据。今天开始把这些能力组合起来,解决真实分析工作中更复杂的问题。
一个经营分析页面通常会同时提出多种计算要求:查看各区域销售 TopN,计算最近七天移动平均,比较用户本次行为与上一次行为,统计新客次日留存,计算多渠道独立访客,展开用户标签,关联订单与会员等级,还要把每笔订单匹配到当时生效的价格快照。每项需求都可以写出 SQL,真正困难的部分在于选择合适的计算表达、数据结构和分布策略,并且让结果在大数据量、高并发环境下保持稳定。
Day 15 建立一套高级分析工具箱,包含六组核心能力:
- 用窗口函数完成组内排名、累计值、移动窗口和前后行比较;
- 用 ARRAY、MAP、STRUCT、JSON/VARIANT 保存并计算嵌套数据;
- 用 Bitmap 完成精确去重、留存和人群集合运算;
- 用 HLL 完成低成本、可合并的近似去重;
- 用正确的 Join 语义和数据分布策略控制网络、内存与数据倾斜;
- 用留存、漏斗、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 会形成可合并的聚合状态。设计阶段先判断数据量、选择性、分组粒度和并发,后续调优才能有明确方向。

二、窗口函数:保留明细的组内计算
Apache Doris 当前窗口函数体系包含专用分析函数和大部分可用于 OVER 的聚合函数。它们在 JOIN、WHERE、GROUP 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 支持 ROWS 和 RANGE。涉及固定前后行数时,显式使用 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_VALUE 和 LAST_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 等类型。

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,副本分布稳定 |

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 证明。
