Doris architecture

ANALYTICAL DATABASE / SOURCE READING / LESSON 03

Nereids optimizer

Study Doris logical planning, rules, costing, and physical plan generation through the Nereids Cascades framework.

Reading
60 min
Track
Doris architecture
Source
Chinese source notes

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 区别

维度CalciteNereids
语言Java (通用)Java (Doris 专用)
Rule 调度固定顺序批次Memo 内驱动 (Rule → Pattern Match → Apply)
Cost ModelVolcano CostCostModelV1 (CPU/IO/Net/Mem)
分布式无(单机)支持 Fragment/Exchange 插入

源码导航

文件关键类/方法职责
fe/fe-core/.../nereids/NereidsPlanner.javaplan(), analyze(), optimize()优化流程编排入口
fe/fe-core/.../nereids/memo/Memo.javacopyIn(), getRoot()等价计划容器
fe/fe-core/.../nereids/memo/Group.javagetLogicalExpressions(), getBestPlan()等价组
fe/fe-core/.../nereids/memo/GroupExpression.javagetPlan(), getCost()组内单个计划
fe/fe-core/.../nereids/rules/Rule.javagetPattern(), transform()规则接口
fe/fe-core/.../nereids/cost/CostCalculator.javacalculate()代价计算
fe/fe-core/.../nereids/stats/Statistics.javagetRowCount(), getNdv()统计信息
fe/fe-core/.../nereids/jobs/executor/Optimizer.javaoptimize()Cascades 优化循环

2. Rule 体系分类

5 大类 Rule

分类包路径 (相对 nereids/rules/)职责示例
analysisanalysis/语义检查, 类型推断BindRelation, BindExpression
rewriterewrite/逻辑等价改写 (RBO)EliminateOuterJoin, SimplifyAgg
explorationexploration/探索等价计划 (CBO 搜索)JoinCommute, JoinReorder
implementationimplementation/逻辑→物理转换LogicalJoin → PhysicalHashJoin
expressionexpression/表达式级优化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效果源码位置
FilterPushDownrewriteFilter(Join)将 Filter 推到 Join 之下rewrite/
ColumnPruningrewriteProject(a,b,c) 只用 a移除不需要的列rewrite/
PartitionPrunerrewriteFilter(part_col)裁剪不匹配的分区rewrite/
JoinReorderexplorationN 表 Join (N≥3)基于统计信息重排 Joinexploration/
PushDownAggThroughJoinrewriteAgg(Join)聚合下推到 Join 之前rewrite/
MaterializedViewRewriteexploration查询能匹配物化视图用 MV 扫描替代原表扫描exploration/
EliminateOuterJoinrewriteLeftJoin 且右边无引用用 InnerJoin 替代rewrite/
EliminateLimitZerorewriteLIMIT 0直接返回空结果, 不执行rewrite/

3. Memo 结构与搜索过程

Nereids 优化器深度解析 图 01

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 比例检查数据是否均衡分布
cardinalityNereids 估计的行数如偏离实际, 需更新统计信息
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 更新统计信息重试。


Nereids Cascades 框架: Memo / Rule / Cost

下一步