The source notes for this track are currently maintained in Chinese.
预计阅读时间: 60 分钟 前置阅读: doc-01 §4 (优化器概览) 下一次阅读: doc-04 (Pipeline 引擎), doc-09 (Catalog)
1. Cascades 框架回顾
为什么需要 Cascades
传统的 RBO(启发式规则优化器)是"流水线"式的:规则按固定顺序应用一次,无法探索等价计划空间。Cascades 是搜索框架,它将"计划"分解为"等价组",规则不断向等价组中添加新的候选计划,直到无法产生更优计划为止。
核心概念
┌──────────────────┐
│ Memo (计划容器) │
│ ┌──────────────┐ │
│ │ Group 0 │ │ ← 一个等价组(Group)
│ │ ┌──────────┐ │ │
│ │ │Expr A │ │ │ ← GroupExpression A (一个计划)
│ │ │Expr B │ │ │ ← GroupExpression B (等价的计划)
│ │ │Expr C │ │ │ ← GroupExpression C (又一个)
│ │ └──────────┘ │ │
│ └──────────────┘ │
│ ┌──────────────┐ │
│ │ Group 1 │ │
│ │ ... │ │
│ └──────────────┘ │
└──────────────────┘
Group: 语义等价的多个计划组成的集合(如"扫描 t 表的3种方式")
GroupExpression: Group 中的一个具体计划(一个算子树节点)
其 children 是 Group 的引用(非具体 Plan) → 支持组合爆炸式等价
与 Calcite 的 Cascades 区别
| 维度 | Calcite | Nereids |
|---|---|---|
| 语言 | Java (通用) | Java (Doris 专用) |
| Rule 调度 | 固定顺序批次 | Memo 内驱动 (Rule → Pattern Match → Apply) |
| Cost Model | Volcano Cost | CostModelV1 (CPU/IO/Net/Mem) |
| 分布式 | 无(单机) | 支持 Fragment/Exchange 插入 |
源码导航
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
fe/fe-core/.../nereids/NereidsPlanner.java | plan(), analyze(), optimize() | 优化流程编排入口 |
fe/fe-core/.../nereids/memo/Memo.java | copyIn(), getRoot() | 等价计划容器 |
fe/fe-core/.../nereids/memo/Group.java | getLogicalExpressions(), getBestPlan() | 等价组 |
fe/fe-core/.../nereids/memo/GroupExpression.java | getPlan(), getCost() | 组内单个计划 |
fe/fe-core/.../nereids/rules/Rule.java | getPattern(), transform() | 规则接口 |
fe/fe-core/.../nereids/cost/CostCalculator.java | calculate() | 代价计算 |
fe/fe-core/.../nereids/stats/Statistics.java | getRowCount(), getNdv() | 统计信息 |
fe/fe-core/.../nereids/jobs/executor/Optimizer.java | optimize() | Cascades 优化循环 |
2. Rule 体系分类
5 大类 Rule
| 分类 | 包路径 (相对 nereids/rules/) | 职责 | 示例 |
|---|---|---|---|
| analysis | analysis/ | 语义检查, 类型推断 | BindRelation, BindExpression |
| rewrite | rewrite/ | 逻辑等价改写 (RBO) | EliminateOuterJoin, SimplifyAgg |
| exploration | exploration/ | 探索等价计划 (CBO 搜索) | JoinCommute, JoinReorder |
| implementation | implementation/ | 逻辑→物理转换 | LogicalJoin → PhysicalHashJoin |
| expression | expression/ | 表达式级优化 | SimplifyArithmetic, FoldConstant |
Rule 接口
// Rule.java — 核心接口
public abstract class Rule {
// pattern(): 描述这条 Rule 匹配哪种 Plan 形状
public abstract RuleType getRuleType();
public abstract Pattern<? extends Plan> getPattern();
// transform(): 应用 Rule, 产生新的计划(可能多个, 加到 Memo)
public abstract List<Plan> transform(Plan node, CascadesContext context);
}
重点 Rule 清单
| Rule | 分类 | 触发 Pattern | 效果 | 源码位置 |
|---|---|---|---|---|
FilterPushDown | rewrite | Filter(Join) | 将 Filter 推到 Join 之下 | rewrite/ |
ColumnPruning | rewrite | Project(a,b,c) 只用 a | 移除不需要的列 | rewrite/ |
PartitionPruner | rewrite | Filter(part_col) | 裁剪不匹配的分区 | rewrite/ |
JoinReorder | exploration | N 表 Join (N≥3) | 基于统计信息重排 Join | exploration/ |
PushDownAggThroughJoin | rewrite | Agg(Join) | 聚合下推到 Join 之前 | rewrite/ |
MaterializedViewRewrite | exploration | 查询能匹配物化视图 | 用 MV 扫描替代原表扫描 | exploration/ |
EliminateOuterJoin | rewrite | LeftJoin 且右边无引用 | 用 InnerJoin 替代 | rewrite/ |
EliminateLimitZero | rewrite | LIMIT 0 | 直接返回空结果, 不执行 | rewrite/ |
3. Memo 结构与搜索过程

Memo 核心数据结构
// Memo.java — 计划容器
public class Memo {
// Group 表(按 GroupId 索引)
private final List<Group> groups;
// 根 Group (最外层查询的等价组)
private Group root;
// 添加一个新的 GroupExpression 到指定的 Group
public GroupExpression copyIn(GroupExpression groupExpr, Group target, Plan copy);
}
// Group.java — 等价组
public class Group {
private final int groupId;
private final List<GroupExpression> logicalExpressions; // 逻辑计划
private final List<GroupExpression> physicalExpressions; // 物理计划
private Statistics statistics; // 统计信息
private double bestCost; // 当前最优 Cost
private PhysicalProperties bestProperties; // 最优物理属性
}
// GroupExpression.java — 一个具体计划
public class GroupExpression {
private final Plan plan;
private final List<Group> children; // 子节点是 Group 引用!
private final double cost;
}
Explore vs Optimize
Phase 1: EXPLORATION (探索等价计划)
→ 遍历 Memo 中的所有 Group
→ 为每个 GroupExpression 尝试匹配所有 exploration Rule
→ Rule 产生的新 Plan → copyIn 到 Memo (可能形成新 Group)
→ 循环直到没有新 Plan 产生 (达到 fixed point)
Phase 2: OPTIMIZATION (选择最优计划)
→ 为每个 Group 计算所有 GroupExpression 的 Cost
→ 选择 Cost 最小的 GroupExpression 作为该 Group 的最优计划
→ 自底向上累加: child group 的最优 cost + 当前算子的 cost
→ 根 Group 的最优计划 = 最终物理计划
4. Cost Model 与统计信息
Cost 计算公式
TotalCost = CPU_Cost + IO_Cost + Network_Cost + Memory_Cost
CPU_Cost = 处理行数 × cpu_weight (默认 1.0)
IO_Cost = 扫描字节数 ÷ io_speed × io_weight
Net_Cost = Shuffle 字节数 × net_weight
Mem_Cost = 峰值内存 × mem_weight
统计信息来源
| 统计信息 | 来源 | 更新频率 |
|---|---|---|
| 行数 (row count) | FE 定期从 BE 收集 | 每 10min (可配) |
| NDV (Distinct count) | BE 采样 | 同上 |
| NULL 比例 | BE 采样 | 同上 |
| 直方图 (Histogram) | BE 采样 (可选) | 同上 |
| 数据大小 (bytes) | BE 上报 | 同上 |
CostCalculator
// CostCalculator.java — 计算计划的 Cost
public class CostCalculator {
// 计算一个物理计划的 Cost (自底向上, 缓存 Child Cost)
public static Cost calculate(PhysicalPlan plan, CascadesContext context) {
Cost childCost = plan.children()
.map(c -> calculate(c, context))
.reduce(Cost.zero(), Cost::add);
Cost selfCost = calculateSelf(plan); // CPU/IO/Net 按公式计算
return childCost.add(selfCost);
}
}
5. 分布式执行计划生成——Fragment/Exchange 插入
Nereids 的最后一步:为物理计划插入 Exchange 和 Fragment 边界。
Fragment 插入规则
Physical Plan:
PhysicalHashAgg (Global)
└── PhysicalOlapScan (orders)
InsertExchange:
→ 检查: Agg 是 Global 还是 Local?
→ Global Agg → 需要子节点的数据按 GROUP BY key shuffle
→ Insert Exchange (SHUFFLE by status)
PhysicalHashAgg (Global)
└── Exchange (SHUFFLE by status)
└── PhysicalHashAgg (Local) ← 新增 (Local Agg 减少数据量)
└── PhysicalOlapScan (orders)
Insert Fragment:
→ Exchange 边界 = Fragment 边界
Fragment 0: Global HashAgg
Fragment 1: Exchange → Local HashAgg → OlapScan
6. Profile/Explain 结果解读
EXPLAIN 输出示例
EXPLAIN SELECT COUNT(*) FROM orders WHERE dt='2026-07' GROUP BY status;
-- Nereids Explain 输出:
PLAN FRAGMENT 0
OUTPUT EXPRS: count(*)
PARTITION: HASH_PARTITIONED: status
STREAM DATA SINK
EXCHANGE ID: 01
HASH_PARTITIONED: status
PLAN FRAGMENT 1
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 01
HASH_PARTITIONED: status
1:AGGREGATE (update finalize)
| output: count(*)
| group by: status
|
0:OlapScanNode
TABLE: orders
PREAGG: ON
partitions=7/7
rollup: orders
TabletRatio=20/20
关键指标解读
| 指标 | 含义 | 优化方向 |
|---|---|---|
PREAGG: OFF | 存储层预聚合未生效, 所有行传回 BE | 检查表设计, 考虑开启聚合模型 |
partitions=N/M | 只扫描 N/M 个分区 (分区裁剪) | 分区合理则 N<<M |
TabletRatio | 扫描的 Tablet 比例 | 检查数据是否均衡分布 |
cardinality | Nereids 估计的行数 | 如偏离实际, 需更新统计信息 |
predicates | 下推的谓词 | 越多越好 (存储层过滤) |
7. 常见问题 / 面试题
Q1: Nereids 和旧 RBO 优化器的性能对比? A: Nereids 在复杂查询(Tpch Q5/Q7/Q9 等)上有 2-5x 的优化, 主要来自 JoinReorder + 统计信息驱动的 Cost Model。简单点查(单表+简单 filter)上两者几乎相同。Nereids 从 Doris 2.0 起成为默认优化器。
Q2: Cascades 的"定点搜索"太慢怎么办? A: (1) RuleSet 按查询类型裁剪 (如单表查询不需要 JoinReorder); (2) Rule 的 pattern() 方法快速过滤; (3) Explore 轮次限制 (默认 10 轮); (4) Cost Improvement 最小阈值 (新计划不显著优化则跳过)。
Q3: 为什么 JoinReorder 不是对所有 Join 都生效? A: 当 N>=N_threshold (默认 10) 时, Join 的排列组合爆炸 (NP-hard), 穷举不可行。Nereids 使用基于统计信息的贪心+动态规划 (DPccp) 组合算法, 在合理的搜索空间中找到近似最优。
Q4: 物化视图改写的触发条件? A: 查询必须满足: (1) 查询的 Scan 表与物化视图的源表匹配; (2) 查询的 Filter 被 MV 的 WHERE 包含(或更严格); (3) 查询的 Project 列在 MV 中; (4) 查询的 Group By + Agg 与 MV 一致或可派生。所有条件满足后, Nereids 比较 Cost (扫描 MV vs 扫描原始表) 决定是否改写。
Q5: 如何为一条慢查询排查"是优化器选错了计划"还是"执行引擎性能差"?
A: 使用 EXPLAIN 看计划是否合理 (Scan 量、Join 顺序、Partition 裁剪)。如果计划看起来正确但查询慢→执行引擎/存储问题。如果计划不优→可能是统计信息过期、缺失直方图、或优化器规则的 bug。使用 ANALYZE TABLE 更新统计信息重试。

下一步
- 执行引擎: doc-04-pipeline-execution.md——优化器产出的计划如何被执行
- 元数据: doc-09-catalog-journal.md——Catalog 和统计信息管理