Apache Calcite 查询编译框架与 Hive 执行引擎深度工程实战:从 SqlNode 校验、Volcano/Cascades 规划器、RelTrait 与 Convention 传播、物化视图改写、Rex 表达式化简到 Hive Tez DAG 与向量化执行的全链路解析

执行摘要

数据基础设施有一个被反复重新发明的轮子:把一个 SQL 字符串变成一台机器能跑的执行计划。Hive、Flink、Druid、Doris、Kylin、Phoenix、Samza 各自造过一遍,最后几乎全部收敛到同一个答案——Apache Calcite。

Calcite 不是一个数据库,也不是一个执行引擎,它是一套可嵌入的查询编译器:把 SQL 解析、校验、关系代数表示、代价优化、物化视图改写、逻辑到物理的转换抽象成一组可以单独替换的 SPI,然后把"最终怎么跑"交还给宿主系统。理解 Calcite,本质上是在理解现代 SQL 引擎的公共骨架。

本文沿着一条真实链路拆开讲:SQL → SqlNode → 校验 → RelNode/RexNode → Volcano 或 Hep 规划器 → Trait/Convention 传播 → 物化视图改写 → 落到 Hive 的 Tez DAG 与向量化执行。最后给出一份生产可调的清单。

一、Calcite 的定位:把"编译"和"执行"切开

传统单体数据库(PostgreSQL、MySQL)里,解析器、优化器、执行器是同一个进程里紧耦合的三段代码。代价是:你想换一个执行引擎,就得把优化器重新写一遍;你想复用优化器,就得把整个数据库搬过来。

Calcite 的切割点是 RelNode(关系代数节点)+ Convention(物理实现约定)。它负责回答"逻辑上该怎么做",宿主系统负责回答"物理上怎么跑":

层归属关键抽象
解析CalciteSqlParser → SqlNode(JavaCC 生成的语法树)
校验CalciteSqlValidator、SqlValidatorCatalogReader、SqlValidatorScope
逻辑计划CalciteSqlToRelConverter → RelNode / RexNode
优化CalciteRelOptPlanner、RelOptRule、RelOptCost、RelMetadataQuery
物理实现宿主Convention、ConverterRule、RelOptTable 实现
执行宿主Hive 的 Operator/Tez、Flink 的 ExecNode、Druid 的 Query

这个拆分带来的工程收益是巨大的:Hive 只需要实现"我的表长什么样"(RelOptHiveTable)和"我的算子怎么跑"(HiveTableScan 等),剩下的谓词下推、列裁剪、聚合下推、连接重排全由 Calcite 的规则库免费提供。

二、第四阶段流水线:Parse → Validate → Optimize → Execute

2.1 解析与校验:比想象中重

SqlNode 只是一棵语法树,它不知道 t.a 是不是合法列名,也不知道 a + b 的结果类型。这一步由 SqlValidatorImpl 完成,核心是两套机制:

  • SqlValidatorNamespace:描述一个关系表达式"产出什么"(行类型、是否可修改)。TableScan、Select、Join 各自有实现。
  • SqlValidatorScope:描述"在这个位置上,一个标识符能看见什么"。SelectScope、JoinScope、AggregatingScope 决定了 SELECT a FROM t GROUP BY b 里 a 为什么非法。

类型推导走的是 SqlOperator.returnTypeInference,最终落到 RelDataTypeFactory 产出的 RelDataType。一个容易被忽略的工程细节:校验阶段会顺手把 SqlNode 树做一次标准化 rewriting(比如把 USING 转成 ON、把 * 展开),这决定了后续 SqlToRelConverter 看到的树已经比用户写的干净得多。

2.2 逻辑关系代数:RelNode 与 RexNode 的二分

这是 Calcite 最关键的建模决策:

  • RelNode 描述关系算子(Scan / Filter / Project / Join / Aggregate / Sort),它们的输入和输出都是"行集合"。
  • RexNode 描述行内标量表达式(RexInputRef 引用列、RexLiteral 常量、RexCall 函数调用、RexSubQuery 子查询、RexOver 窗口函数、RexCorrelVariable 相关变量)。

把表达式和算子分开,意味着表达式改写可以独立于算子树进行。谓词下推、常量折叠、公共子表达式消除全部作用于 RexNode,用 RexShuttle / RexVisitor 遍历即可:

// 用 RexShuttle 把不可下推的函数替换成常量,从而让谓词变成可下推形式
RexNode simplified = rexBuilder.makeCall(
    SqlStdOperatorTable.EQUALS,
    ref,
    rexBuilder.makeLiteral("2026-01-01"));

RexSimplify simplify = new RexSimplify(rexBuilder, predicates, executor);
RexNode result = simplify.simplifyUnknownAsFalse(simplified);

RexSimplify 是生产中最常用的工具之一:它做常量折叠、IS NULL 推导、OR 分支去重、范围合并,Hive 的 HiveRexExecutorImpl 就是把 RexCall 编译成 Java 代码后反射执行,用来在优化期算出确定性表达式的常量值。

2.3 去相关子查询

SqlToRelConverter 会先保留 LogicalCorrelate 表示相关子查询,随后由 RelDecorrelator + SubQueryRemoveRule 把它重写为普通 Join + Aggregate。这一步是 CBO 能否生效的分水岭:只要留下 Correlate,代价模型基本失效,因为相关子查询的行数估计无法用常规 selectivity 传播。

三、两套规划器:Volcano 与 Hep 的分工

Calcite 同时提供两个规划器,生产系统几乎总是串着用。

3.1 HepPlanner:启发式、线性、可控

HepPlanner 按 HepProgram 定义的顺序,用 HepMatchOrder(BOTTOM_UP、TOP_DOWN、DEPTH_FIRST、ARBITRARY)遍历 DAG,对每个匹配点应用规则,不做代价比较。它快、确定性强、不会爆炸,适合做"必然正确"的改写:

HepProgram program = new HepProgramBuilder()
    .addRuleInstance(CoreRules.FILTER_INTO_JOIN)
    .addRuleInstance(CoreRules.PROJECT_MERGE)
    .addRuleInstance(CoreRules.FILTER_PROJECT_TRANSPOSE)
    .addMatchOrder(HepMatchOrder.BOTTOM_UP)
    .build();
Program hep = Programs.of(program, false, DefaultRelMetadataProvider.INSTANCE);

注意 HepMatchLimit 默认是 Integer.MAX_VALUE,在自引用规则(如 PROJECT_MERGE 与 PROJECT_REDUCE_EXPRESSIONS 互相触发)的组合下可能循环,需要显式设限。

3.2 VolcanoPlanner:Cascades 风格的动态规划

VolcanoPlanner 是核心。它维护:

  • RelSet:一组逻辑等价的表达式(去掉 trait 后相同)。
  • RelSubset:RelSet 中具有某个确定 RelTraitSet 的子集,是优化器真正的搜索节点。
  • RuleQueue:按 Importance 排序的待应用规则队列,importance 由子节点的代价与子树"是否被父节点需要"反推,这就是 Cascades 的自顶向下剪枝。
  • VolcanoCost:三元组(rowCount、cpu、io)。注意 rowCount 在比较中只是一个权重,不是真实行数;真实行数来自 RelMetadataQuery.getRowCount(),由 RelMdRowCount handler 提供。

搜索过程:把根节点按目标 trait 注册成 RelSubset,不断 pop 队列里 importance 最高的 RelOptRuleCall,onMatch 里 call.transformTo(newRel) 把新表达式注册进对应的 RelSubset,若该 subset 的 best plan 被更新,则向上传播 importance。直到队列空或达到 VolcanoPlanner.setNoneConventionHasInfiniteCost / 迭代上限。

生产里最关键的一个旋钮是规则集合的大小。Volcano 的搜索空间对 join 数量是指数级的:

VolcanoPlanner planner = new VolcanoPlanner();
planner.addRelTraitDef(ConventionTraitDef.INSTANCE);
planner.addRelTraitDef(RelCollationTraitDef.INSTANCE);
planner.setTopDownOpt(true);                    // 1.26+ 开启真正的 top-down 搜索
CalciteSystemProperty.BUSH_JOIN_THRESHOLD.value(); // 超过阈值切换到 bushy/启发式

Hive 用的是自己的 HiveVolcanoPlanner,并显式禁用了部分在 OLAP 场景下容易选错的规则(例如某些 semijoin 变体),同时把 join 重排限制在 hive.cbo.enable 打开且统计信息可用时。

四、Trait 与 Convention:跨引擎的物理属性传播

Convention 是"这个 RelNode 用什么引擎/什么形式实现"的标记。RelCollation 表示排序,RelDistribution 表示分布。三者都是 RelTrait,由各自的 RelTraitDef 注册到 planner。

优化器处理物理需求的机制是trait propagation + enforcement:

  1. RelNode.deriveRowCount/RelTraitSet 由子节点向父节点传播:LogicalSort 要求输入有 RelCollation,子节点若是 HiveSort 就自带;
  2. 若子节点无法满足,planner 生成 AbstractConverter 占位;
  3. ExpandConversionRule 通过 ConverterRule 把逻辑算子转成满足 trait 的物理算子,或在无 converter 时插入 enforcer(例如插入一个 Sort/Exchange)。

这就是同一套优化规则能在"Hive 上跑"和"Flink 流上跑"之间复用的原因——规则只操作逻辑 RelNode,物理差异全部编码在 Convention 里:

public interface HiveRel {
  Convention CONVENTION = new Convention.Impl("HIVE", HiveRel.class);
}

// ConverterRule:把 LogicalFilter 转成 HiveFilter,前提是输入已经是 HIVE convention
public class HiveFilterRule extends ConverterRule {
  public static final HiveFilterRule INSTANCE =
      Config.INSTANCE
          .withConversion(LogicalFilter.class, Convention.NONE,
                          HiveRel.CONVENTION, "HiveFilterRule")
          .withRuleFactory(HiveFilterRule::new)
          .as(Config.class)
          .toRule(HiveFilterRule.class);

  @Override public RelNode convert(RelNode rel) {
    LogicalFilter filter = (LogicalFilter) rel;
    return new HiveFilter(filter.getCluster(),
        convert(filter.getTraitSet(), HiveRel.CONVENTION),
        convert(filter.getInput(), HiveRel.CONVENTION),
        filter.getCondition());
  }
}

Trait 缺失是生产中最常见的"优化器突然变笨"根因:如果一个 RelOptTable 实现没有声明它已经按分桶列 hash 分布(RelDistribution.Hash),优化器就不知道 join 可以免 shuffle,会硬塞一次 ReduceSink,代价直接翻倍。

五、物化视图改写:SubstitutionVisitor 与 UnifyRule

Calcite 的物化视图改写不是"匹配 SQL 文本",而是在关系代数层面做结构统一(unification):

  • 用户注册 RelOptMaterialization(viewTable, viewQueryRel, ...) 到 planner;
  • 改写规则(MaterializedViewScanRule 家族)遍历查询中的 Scan/Aggregate/Join;
  • 底层 SubstitutionVisitor 用一组 UnifyRule(UnifyFilterToScan、UnifyProjectToScan、UnifyAggregateToScan、UnifyJoinToScan)尝试把查询片段与视图定义统一;
  • 统一成功后用 RelBuilder 构造"视图扫描 + 补偿算子(compensation)",补偿部分通常由 RexShuttle 做列名与表达式重映射。
RelOptMaterialization mat =
    new RelOptMaterialization(viewRel, viewRel, null, viewTableName);
planner.addMaterialization(mat);

工程上要抓三个点:

  1. 行数与代价必须准确,否则改写后的 plan 在 Volcano 里被 best-cost 比较淘汰,等于白改;
  2. 列裁剪:视图通常宽,改写的收益取决于能否只扫描少量列,列式格式(ORC/Parquet)下这一点被显著放大;
  3. 新鲜度:Calcite 不负责视图刷新,Hive 侧需要 hive.materializedview.rewriting 与重建任务配合,过期视图必须被禁用,否则会产出错误结果。

六、落到 Hive:从 RelNode 到 Tez DAG 与向量化执行

Hive 的 CBO 入口是 CalcitePlanner(继承 SemanticAnalyzer),它:

  1. 用 RelOptHiveTable 包装表与分区信息——分区裁剪在这里完成,getRowCount() 直接读 pruned partition list 的统计信息;
  2. 构造 Hive 自己的 RelNode 家族(HiveTableScan、HiveFilter、HiveJoin、HiveAggregate、HiveSortLimit),全部打上 HiveRel.CONVENTION;
  3. 走 HiveVolcanoPlanner 优化,代价与行数由 HiveRelMdRowCount / HiveRelMdSelectivity 提供(走 Hive metastore 的列统计:NDV、min/max、histogram);
  4. 把最终 RelNode 反向转成 Hive 的 Operator 流水线。

6.1 Operator 流水线 → Tez Vertex

Hive 的执行计划是一条 Operator 链:TableScanOperator → FilterOperator → SelectOperator → (JoinOperator | GroupByOperator) → ReduceSinkOperator → ... → FileSinkOperator。

TezCompiler 负责切图:ReduceSinkOperator 就是切分点。它把两侧的算子分到不同的 Tez Vertex,并依据 ReduceSink 的分区键、join 类型、是否 sort merge 决定 Edge 属性:

Edge 语义触发条件
BROADCASTmap join、hive.auto.convert.join 且小表低于阈值
SCATTER_GATHER常规 shuffle join / group by
ONE_TO_ONE同并发、无 shuffle
CUSTOM_SIMPLE_EDGE / CUSTOM_EDGESMB join、bucket map join

切完之后由 TezTask 提交给 YARN,每个 Vertex 的并行度来自 hive.tez.* 的 sizing 估算(hive.tez.min.partition.factor / max.partition.factor),这也是"reducer 数量凭感觉"的根源——它其实由输入字节数除以 hive.exec.reducers.bytes.per.reducer 估算。

6.2 向量化执行:为什么能快 3~5 倍

Hive 的行式路径每一行都要走一遍 ObjectInspector 反射取值 + 表达式对象方法调用,函数调用开销与分支预测失败占了大头。向量化路径把数据组织成 VectorizedRowBatch:

  • 默认 1024 行一批,每列是一个 ColumnVector(LongColumnVector、DoubleColumnVector、BytesColumnVector、DecimalColumnVector、StructColumnVector、ListColumnVector);
  • 表达式由 Vectorizer(配合 VectorizationContext)编译成 VectorExpression——evaluate(VectorizedRowBatch) 一次处理整批,没有虚调用,循环体可被 JIT 内联;
  • Filter 不生成新批次,而是写 selected 索引数组并设 size,实现延迟物化;
  • null 用独立 bitmap(isNull)而非对象引用,避免装箱。
<property><name>hive.vectorized.execution.enabled</name><value>true</value></property>
<property><name>hive.vectorized.execution.reduce.enabled</name><value>true</value></property>
<property><name>hive.vectorized.execution.mapjoin.native.enabled</name><value>true</value></property>

一个实战教训:向量化一旦中途中断(某个算子不支持向量化,如部分 UDF、复杂类型某些操作),Hive 会在该处插入 Vectorized→Row 的适配器,前后各付一次转换成本。EXPLAIN VECTORIZATION 里的 Execution mode: vectorized 必须逐 Vertex 看,只要有一个 Vertex 显示 row 且数据量很大,优化收益就被吃掉了。

6.3 LLAP:把 IO 与执行常驻化

Tez 每次查询都要启动新容器、重新打开 ORC 文件、重新读 footer。LLAP(Live Long and Process)用常驻的 LlapDaemon 解决三件事:常驻 executor 免启动开销、off-heap 缓存复用热数据、IO 层旁路直接读 ORC 并下推谓词到 stripe 级。对交互式场景(BI、Ad-hoc)效果显著,但对大批量 ETL 收益有限——因为容器复用省下的几秒在小时级作业里无意义。

七、生产调优清单

症状根因参数 / 动作
join 顺序离谱、走成笛卡尔积无统计信息hive.cbo.enable=true、ANALYZE TABLE ... COMPUTE STATISTICS FOR COLUMNS
全表扫描未裁剪分区谓词未下推 / 分区列上套了函数hive.optimize.ppd=true;把 date_format(dt) 改成 dt between
reducer 数过多或过少sizing 估算失真hive.exec.reducers.bytes.per.reducer、hive.tez.min/max.partition.factor
数据倾斜,个别 task 拖尾热点 keyhive.groupby.skewindata=true、hive.optimize.skewjoin=true,或加盐打散
向量化没生效不兼容算子 / UDFEXPLAIN VECTORIZATION 逐 Vertex 查,改写或换 UDF 实现
小表 join 仍走 shuffle阈值和统计不准hive.auto.convert.join=true + hive.auto.convert.join.noconditionaltask.size
物化视图不生效代价/新鲜度校准行数、确保视图已重建、hive.materializedview.rewriting=true

排查手段上,Hive 的 EXPLAIN 有几个常被忽略的变体:EXPLAIN CBO(看 Calcite 输出的 RelNode 树)、EXPLAIN VECTORIZATION(看向量化覆盖度)、EXPLAIN FORMATTED(含统计信息与行数估算)、EXPLAIN EXTENDED(含 Operator 的 schema 与表达式)。先看 CBO 计划再看 Tez DAG,能立刻区分"优化器选错了"和"执行引擎跑歪了"。

八、结论:Calcite 真正的设计遗产

把 Calcite 的设计提炼成可迁移的原则,其实只有三条:

  1. 用"逻辑代数 + 可插拔物理约定"解耦编译与执行。这决定了同一个优化器可以服务批处理、流处理、OLAP 三种完全不同的运行时。今天 Flink 的 Table Planner、Doris 的优化器、Paimon 的查询层,本质上都是这条思路的再实现。
  2. 用代价模型 + importance 剪枝替代规则顺序的玄学。规则该不该应用、什么时候应用,交给 RelOptCost 与 RuleQueue 的 importance 决定,而不是靠人排规则顺序。Hive 早期版本规则硬编码顺序所积累的历史包袱,正是这条原则的反面教材。
  3. Trait 是物理属性的唯一真源。分布、排序、实现引擎全部编码为 trait,让"哪些操作可以免 shuffle/免排序"变成一个可被规划器推理的命题,而不是散落在各处的 if-else。

真正理解查询编译器,不是为了背下 CoreRules.FILTER_INTO_JOIN 这个名字,而是理解一件事:当"正确性"由关系代数等价变换保证之后,所有的工程精力都可以投入到"如何在不可预测的负载下保持可预测的性能"这一件事上。这套思路,从 Calcite 到 Flink 到任何一个新的 OLAP 引擎,一字未改。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部