Day 17|Nereids 优化器:RBO、CBO、统计信息、Join Reorder 与 Runtime Filter

Day 16 沿着一条 SQL 的完整生命周期,从客户端连接一路走到结果返回。今天把视角停留在 FE 的优化阶段,深入理解 Doris 怎样把一棵“语义正确”的逻辑计划,转换成一套在当前数据规模、节点分布和资源条件下更合适的物理执行方案。
优化器决定了很多关键问题:哪张表先参与 Join,哪张表负责构建 Hash Table,是否使用 Broadcast,数据需要怎样 Shuffle,聚合分几级完成,哪些过滤条件能够提前下推,哪些分区和列可以直接裁剪,以及 Runtime Filter 应该从哪里生成、下推到哪里。
优化器能力很强,仍然需要可信输入。统计信息过期、热点值集中、外部表行数未知、数据量突然放大,都可能让估算偏离真实执行。生产调优因此形成一条稳定链路:
理解规则
→ 检查统计信息
→ 阅读 EXPLAIN
→ 执行并采集 Profile
→ 对比估算行数与真实行数
→ 选择 Schema、统计、Hint 或执行侧优化
→ 回归验证并保留回滚方案
今天完成后,你应当能够:
- 解释 Nereids 的整体工作方式,以及 RBO、CBO、Memo 和 Cost Model 各自承担的职责;
- 读懂表级、列级统计信息,知道什么时候需要主动执行
ANALYZE; - 判断 Join Reorder、Build/Probe 选择和 Distribution 是否合理;
- 用
EXPLAIN SHAPE PLAN、EXPLAIN VERBOSE、EXPLAIN MEMO PLAN和 Query Profile 建立证据链; - 判断 Runtime Filter 是否真正减少了扫描、网络和 Join Probe;
- 在统计信息失真、数据倾斜和优化器搜索空间过大时,采取可验证、可回滚的干预措施。
一、优化器解决的核心问题:同一条 SQL 可以有很多种执行方法
SQL 描述业务目标。例如:
SELECT
u.city,
p.category,
SUM(i.quantity * i.unit_price) AS revenue
FROM fact_order o
JOIN dim_user u
ON o.user_id = u.user_id
JOIN fact_order_item i
ON o.order_id = i.order_id
JOIN dim_product p
ON i.product_id = p.product_id
WHERE o.order_date BETWEEN '2026-08-01' AND '2026-08-31'
AND o.status = 'PAID'
GROUP BY u.city, p.category;
这段 SQL 没有规定以下细节:
- 先关联订单和用户,还是先关联订单明细和商品;
dim_user、dim_product是否值得广播到所有 BE;- 两张大表是否要按 Join Key 重新分区;
- 聚合应先在每个 BE 局部完成,再汇总到上层,还是直接集中聚合;
order_date和status能否在 Scan 阶段提前过滤;- 是否能够命中分区裁剪、索引或物化视图;
- Runtime Filter 应使用 IN、Bloom、MinMax,还是由执行器自适应选择;
- 预计读取多少行、发送多少网络数据、占用多少内存。
这些决定会组合成大量候选计划。四张表有多种 Join 顺序,每个 Join 又可能选择不同算法和数据分布方式。候选数量会快速增长,优化器需要在有限时间内搜索、剪枝和选择。
当前 Doris 4.x 的统一查询优化器是 Nereids。官方将它定义为基于 Cascades 框架的现代优化器,通过 RBO 与 CBO 协同生成高效执行计划。Legacy Planner 已从当前主线移除,因此 4.x 课程统一使用 Nereids 术语和计划输出。

可以把优化器看成三个连续动作:
第一步:把表达式改写得更容易执行
第二步:枚举可行的物理方案并估算代价
第三步:输出物理计划,交给 Fragment 与 Pipeline 执行
二、RBO 与 CBO:两类优化力量怎样协作

2.1 RBO:依靠确定性规则快速收敛
RBO,Rule-Based Optimization,按照预定义规则改写逻辑计划。它关注逻辑等价性和确定性收益。典型规则包括:
- 谓词下推:尽早执行过滤,减少上层算子的输入行数;
- 列裁剪:只保留后续真正需要的列,减少解码、内存和网络成本;
- 分区裁剪:根据分区条件排除无关 Partition;
- 常量折叠:提前计算常量表达式,删除恒真或恒假的分支;
- 投影合并:合并重复或无效的 Project;
- Limit 下推:在语义允许时尽早限制行数;
- 子查询去相关:把相关子查询转换为 Join、Semi Join、Anti Join 等关系表达;
- 无用排序消除:删除不会影响最终语义的中间排序;
- 表达式标准化:将等价条件转换成更容易匹配规则和索引的形式。
下面的查询在逻辑上先写 Join,再写过滤:
SELECT u.city, SUM(o.amount)
FROM fact_order o
JOIN dim_user u ON o.user_id = u.user_id
WHERE o.status = 'PAID'
AND o.order_date >= '2026-08-01'
GROUP BY u.city;
RBO 会尝试把订单侧条件移动到 fact_order Scan 附近。过滤后的订单行数降低,Join、Runtime Filter、聚合和 Exchange 都能得到更小输入。
RBO 需要遵守语义边界。Outer Join 的保留侧和 NULL 生成侧处理方式不同,过滤条件移动错误会改变结果。Doris 4.1.2 的修复记录中包含“防止 Runtime Filter 穿过 Outer Join 进行不安全下推”,这类 Patch 说明优化规则必须同时满足性能和正确性。
2.2 CBO:在多个候选方案之间进行成本取舍
CBO,Cost-Based Optimization,会枚举等价的候选计划,估算每个计划的执行成本,选择综合代价更低的方案。它重点处理:
- Join Reorder;
- Hash Join、Nested Loop Join 等物理算法选择;
- Broadcast、Shuffle、Bucket Shuffle、Colocate 等 Distribution 选择;
- 聚合的局部与全局阶段;
- 扫描、计算、网络和内存代价;
- 物理属性,例如数据分布和排序;
- 某些物化视图改写结果是否值得采用。
CBO 依赖三个核心输入:
- 统计信息:行数、NDV、NULL、Min/Max、平均列宽等;
- 数据特征与物理属性:分区、分桶、排序、节点分布、是否同位;
- 成本模型:CPU、IO、网络、内存和算子行为的综合估算。
RBO 与 CBO 经常连续工作。RBO 先减少无关列、无关分区和无效表达式,CBO 再针对收敛后的计划搜索物理实现。规则改写质量会影响 CBO 搜索空间,统计信息质量会影响 CBO 排序候选方案的能力。
2.3 Cost Model 关注哪些资源
“代价”包含多个资源维度的综合估算。理解这些维度,有助于解释一些看似反直觉的计划,也能帮助调优人员判断某个计划究竟在节省扫描、网络、CPU,还是在交换内存和稳定性。
扫描成本与读取行数、读取列宽、压缩解码、索引裁剪和远端文件访问有关。两条计划最终返回相同行数,其中一条需要扫描 2 TB,另一条只扫描 80 GB,后者通常拥有明显优势。
CPU 成本来自表达式计算、Hash、比较、序列化、聚合和排序。字符串 Join Key、复杂表达式和高基数 Group By 会提高单位行成本。
网络成本主要来自 Exchange。Broadcast 的网络量近似为 Build 数据量乘以接收节点数;Shuffle 的网络量取决于参与重分布的数据总量。网络带宽充裕时,适度 Shuffle 可能换来更均衡的并行计算;网络繁忙时,Colocate 和 Bucket Shuffle 的价值会更加明显。
内存成本集中在 Hash Join Build、Hash Aggregate、Sort、窗口函数和缓存结构。计划在内存预算内完成时延迟更低;超出预算后会触发 Spill 或查询终止。
物理属性成本关注数据是否已经具备目标分布和排序。某个候选计划的算子数量更多,但能够复用现有分桶、排序或 Colocate 属性,综合成本仍可能更低。
成本模型输出是相对选择依据,不能替代运行时测量。硬件差异、磁盘冷热、网络争用、数据倾斜和并发负载都会让真实耗时偏离估算。调优人员应把 CBO Cost 看成“选择候选方案的内部评分”,再用 Profile 完成现实校验。
三、Memo:Nereids 怎样管理不断扩张的候选计划
多表 Join 的搜索空间很容易膨胀。六张表已经存在大量顺序和树形组合,再叠加算法与 Distribution,候选数量会继续增长。Cascades 框架通过 Memo 管理等价表达和物理候选。

可以用两个概念理解 Memo:
- Group:保存一组逻辑等价的表达;
- GroupExpression:保存某种具体算子组合或物理实现。
例如 A JOIN B JOIN C 可能形成以下等价表达:
(A JOIN B) JOIN C
A JOIN (B JOIN C)
(B JOIN C) JOIN A
……
对于同一个 Join,还可能出现:
Hash Join + Broadcast
Hash Join + Shuffle
Hash Join + Bucket Shuffle
Colocate Join
Nested Loop Join
Memo 的价值体现在三点:
- 等价子表达可以复用,减少重复优化;
- 规则能够围绕 Group 继续产生候选;
- 成本上明显缺乏竞争力的候选可以尽早剪枝。
优化器仍然受到规划时间约束。当前官方文档中,nereids_timeout_second 控制最大规划时间,默认值为 30 秒。查询包含大量外部表、深层子查询或复杂自连接时,优化搜索和元数据访问可能同时增加。遇到规划超时时,应先检查 SQL 结构、表数量和外部 Catalog 延迟,再决定是否提高该参数。
SET nereids_timeout_second = 60;
提高超时时间只扩大规划窗口,不会自动改善候选质量。长期复杂 SQL 更适合通过分层建模、物化视图、减少重复子查询和控制 Join 数量降低搜索复杂度。
四、统计信息:CBO 的燃料和生产治理对象
你原有培训材料已经指出,CBO 依靠 Row Count、Min/Max、Null Count 和 NDV 选择 Join 顺序与算法。当前 4.x 官方统计体系进一步给出了完整的采集、查看、缓存、任务和健康度治理方法。

4.1 Doris 当前收集哪些基础统计
Doris 按表和列记录以下基础信息:
| 指标 | 含义 | 对优化器的价值 |
|---|---|---|
row_count |
行数 | 判断表与中间结果规模 |
data_size |
列总数据大小 | 估算扫描、网络和内存成本 |
avg_size_byte |
平均行宽 | 估算 Hash Table、Block 和 Exchange 字节数 |
ndv |
不同值个数 | 估算等值过滤与 Join 基数 |
min / max |
最小值与最大值 | 估算范围过滤与部分 Runtime Filter |
null_count |
NULL 数量 | 估算过滤、Join 和数据倾斜 |
当前基础统计支持 BOOLEAN、整数、浮点数、日期时间和字符串等基础类型。VARIANT、ARRAY、MAP、STRUCT、HLL、BITMAP 等复杂类型会被自动跳过。Day 14 的半结构化宽表因此需要把关键过滤 Path 提升为稳定实体列,或通过 Path 索引承担查询加速;CBO 基础统计无法覆盖所有动态 Path。
4.2 手工采集
-- 异步全量采集
ANALYZE TABLE fact_order;
-- 同步采集,适合实验和上线门禁
ANALYZE TABLE fact_order WITH SYNC;
-- 指定关键列
ANALYZE TABLE fact_order (user_id, order_date, status, amount) WITH SYNC;
-- 抽样 100000 行
ANALYZE TABLE fact_order WITH SAMPLE ROWS 100000;
-- 数据库级采集
ANALYZE DATABASE day17_lab;
同步模式会在采集结束后返回,适合需要立即比较 Explain 的实验。大表全量采集可能占用明显资源,生产环境应在低峰执行,并优先选择参与过滤、Join、分组和排序的关键列。
4.3 自动采集
内部表默认启用自动采样。当前官方文档给出的关键机制包括:
- 后台线程周期检查表和列;
- 表缺少统计信息时会进入候选;
- 表健康度低于阈值时会重新采集,默认健康阈值为 90;
- 内部表发生变化且 24 小时内未采集,也可能触发更新;
- 大表默认抽样 4,194,304 行;
- 自动采集默认排除列数超过 300 的宽表;
- 自动检查间隔默认 5 分钟,但单张表没有“5 分钟内必然完成”的保证。
SHOW VARIABLES LIKE 'enable_auto_analyze';
SHOW VARIABLES LIKE 'table_stats_health_threshold';
SHOW VARIABLES LIKE 'huge_table_default_sample_rows';
SHOW VARIABLES LIKE 'auto_analyze_table_width_threshold';
Day 14 的超宽表需要特别留意 auto_analyze_table_width_threshold。列数超过阈值以后,自动采集可能长期不进入该表。生产治理可以选择提高阈值,也可以只对关键实体列执行手工采集。
4.4 查看采集结果与任务
SHOW TABLE STATS fact_order;
SHOW COLUMN STATS fact_order;
SHOW COLUMN CACHED STATS fact_order;
SHOW ANALYZE;
SHOW AUTO ANALYZE;
SHOW ANALYZE TASK STATUS <job_id>;
持久化统计与 FE Cache 都需要核验。持久化结果存在、当前入口 FE Cache 为空时,计划行为仍可能与预期不同。多 FE 环境还要确认具体执行查询的 FE 已加载目标统计。
4.5 什么时候应主动刷新
建议将以下事件列为强制核验点:
- 首次完成全量导入;
- 大批历史回灌;
- 数据量增长超过原规模一个明显台阶;
- 新分区首次写入大量数据;
- 热点值分布变化;
- 过滤条件命中率发生变化;
- 升级后关键计划出现明显变化;
- 物化视图、分区、分桶或表模型调整;
- Explain 中估算行数与 Profile 真实行数差距过大。
“定期 ANALYZE”只是基础动作。更成熟的做法是把 updated_rows、last_analyze_row_count、关键列 NDV 和业务数据量一起纳入数据任务验收。
4.6 外部 Catalog 的统计边界
Doris 可以为 Hive、Iceberg、JDBC、Paimon 等外部对象提供行数估算或列统计,但不同 Catalog 的能力并不完全一致。当前官方能力矩阵中,Hive 支持手工全量、手工抽样和自动采集;Iceberg 与 JDBC 支持手工全量,自动列统计默认不启用;其他类型需要结合对应 Catalog 文档确认。
外部表规划还会受到元数据和文件信息影响。Hive 可能从 Metastore 的 numRows、totalSize 和文件大小估算行数;Iceberg 可以结合 Snapshot 记录总行数和 Position Delete;JDBC 通常需要向远端数据库执行行数查询。外部统计为空时,SHOW TABLE STATS 可能返回 row_count=-1。这时优化器只能依赖有限信息,跨源 Join 更容易出现错误的 Broadcast 或 Join 顺序。
-- 为外部 Catalog 开启自动分析,需要先确认扫描成本
ALTER CATALOG hive_catalog
SET PROPERTIES ('enable.auto.analyze'='true');
-- 表级策略可以覆盖 Catalog 级策略
ALTER TABLE hive_catalog.ods.orders
SET ('auto_analyze_policy'='enable');
外部表统计采集可能触发大范围文件扫描或远端 SQL,应在 POC 中核验采集耗时、远端压力和数据新鲜度。对于稳定的大表,可以将行数、关键列 NDV 和分区级数据量作为平台元数据定期回灌;对于频繁变化的湖表,应结合 Snapshot 与分区变化安排采集。
4.7 统计采集本身也需要资源治理
统计信息通过 SQL 读取数据,同样会占用 BE 扫描、CPU、内存和 IO。当前官方文档列出了 statistics_simultaneously_running_task_num、statistics_sql_mem_limit_in_bytes、自动分析时间窗口等配置。生产环境应遵循以下方法:
- 核心高频表优先,低频冷表延后;
- 关键过滤列和 Join 列优先,展示型长文本列可以跳过;
- 超大表优先抽样,只有估算误差无法接受时再做全量;
- 自动分析安排在低峰时段;
- 大规模回灌任务结束后,手工同步采集关键列;
- 采集失败时检查任务状态、错误信息和内部统计表 Tablet 健康度;
- 统计结果进入 FE Cache 后再开始 Plan 回归。
统计治理的目标是让计划稳定地接近真实数据分布,同时控制采集成本。
五、Cardinality 与 Selectivity:一个估算误差怎样传导到整棵计划
Cardinality 表示某个算子预计输出多少行;Selectivity 表示过滤后保留的比例。两者会连续传导:
Scan 估算
→ Filter 输出估算
→ Join Build / Probe 估算
→ Join 中间结果估算
→ Aggregate 输入估算
→ Distribution、并行度与内存估算

假设订单表有 1 亿行,country 的 NDV 为 100,并且优化器暂时按照较均匀分布估算:
country = 'CN' 的选择率 ≈ 1 / 100
估算输出 ≈ 100 万行
真实数据中,中国订单可能占 45%。此时实际输出为 4500 万行,估算误差达到 45 倍。这个误差会继续影响:
- 是否把过滤后的订单侧当成“小表”;
- Hash Join Build 侧选择;
- Broadcast 是否安全;
- Hash Table 内存估算;
- Runtime Filter 的候选类型;
- 后续聚合与 Exchange 的规模。
基础统计通常包含 NDV、Min/Max 和 NULL,不等同于完整数据分布模型。严重热点、值相关性和多列联合分布仍可能超出均匀假设。生产诊断要同时查看两类数字:
Explain 中的 cardinality / estimated rows
Profile 中的 InputRows / BuildRows / ProbeRows / RowsReturned
当估算和真实结果相差一个到多个数量级,应先检查统计新鲜度和数据分布,再决定 Hint。直接固定 Join 顺序可能掩盖统计治理问题,并在下一次数据变化后形成新的坏计划。
5.6 怎样检查统计质量
统计信息存在并不代表质量足够。建议从完整性、新鲜度和合理性三个方向检查。
完整性关注核心列是否都有统计。事实表至少覆盖分区列、常用过滤列、Join Key、Group By 列和高频排序列。维表重点覆盖主键、外键、状态、类型和业务过滤列。SHOW COLUMN STATS 为空时,先确认列类型是否支持收集,再检查任务状态和内部统计表健康度。
新鲜度关注 updated_rows、updated_time、new_partition 和 last_analyze_version。大批量写入后,行数变化比例低于阈值也可能改变热点分布,例如新增数据全部集中到一个渠道或状态。此时表健康度看起来尚可,关键列 NDV 和 Min/Max 已经失真,需要按业务事件触发收集。
合理性关注统计值是否符合业务常识:
SELECT COUNT(*) AS rows,
NDV(user_id) AS user_ndv,
MIN(dt) AS min_dt,
MAX(dt) AS max_dt,
SUM(user_id IS NULL) AS null_users
FROM fact_order;
将人工校验结果与 SHOW COLUMN STATS 对比,可以识别采样偏差、历史统计残留和异常注入。官方 FAQ 提供的思路是对 Count、NDV、Min、Max 做独立 SQL 验证。
5.7 外部表统计需要单独设计
Hive、Iceberg、JDBC 等外部表的统计来源和维护成本不同。Hive 可以使用 Metastore 参数、文件大小估算和主动采集;Iceberg 可以从 Snapshot 元数据获得行数;JDBC 可能需要访问远端数据库执行 Row Count。外部 Catalog 默认策略通常更保守,以免周期性扫描大量历史文件或给 TP 数据库增加压力。
跨源 Join 中,一侧统计缺失会影响整个计划。生产设计应记录外部表的统计来源、刷新周期、远端权限和失败降级方式。关键跨源报表还应准备落地到 Doris 内表或异步物化视图的替代路径。
