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(物理实现约定)。它负责回答"逻辑上该怎么做",宿主系统负责回答"物理上怎么跑":
| 层 | 归属 | 关键抽象 |
|---|---|---|
| 解析 | Calcite | SqlParser → SqlNode(JavaCC 生成的语法树) |
| 校验 | Calcite | SqlValidator、SqlValidatorCatalogReader、SqlValidatorScope |
| 逻辑计划 | Calcite | SqlToRelConverter → RelNode / RexNode |
| 优化 | Calcite | RelOptPlanner、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(),由RelMdRowCounthandler 提供。
搜索过程:把根节点按目标 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:
RelNode.deriveRowCount/RelTraitSet由子节点向父节点传播:LogicalSort要求输入有RelCollation,子节点若是HiveSort就自带;- 若子节点无法满足,planner 生成
AbstractConverter占位; 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);
工程上要抓三个点:
- 行数与代价必须准确,否则改写后的 plan 在 Volcano 里被 best-cost 比较淘汰,等于白改;
- 列裁剪:视图通常宽,改写的收益取决于能否只扫描少量列,列式格式(ORC/Parquet)下这一点被显著放大;
- 新鲜度:Calcite 不负责视图刷新,Hive 侧需要
hive.materializedview.rewriting与重建任务配合,过期视图必须被禁用,否则会产出错误结果。
六、落到 Hive:从 RelNode 到 Tez DAG 与向量化执行
Hive 的 CBO 入口是 CalcitePlanner(继承 SemanticAnalyzer),它:
- 用
RelOptHiveTable包装表与分区信息——分区裁剪在这里完成,getRowCount()直接读 pruned partition list 的统计信息; - 构造 Hive 自己的 RelNode 家族(
HiveTableScan、HiveFilter、HiveJoin、HiveAggregate、HiveSortLimit),全部打上HiveRel.CONVENTION; - 走
HiveVolcanoPlanner优化,代价与行数由HiveRelMdRowCount/HiveRelMdSelectivity提供(走 Hive metastore 的列统计:NDV、min/max、histogram); - 把最终 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 语义 | 触发条件 |
|---|---|
BROADCAST | map join、hive.auto.convert.join 且小表低于阈值 |
SCATTER_GATHER | 常规 shuffle join / group by |
ONE_TO_ONE | 同并发、无 shuffle |
CUSTOM_SIMPLE_EDGE / CUSTOM_EDGE | SMB 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 拖尾 | 热点 key | hive.groupby.skewindata=true、hive.optimize.skewjoin=true,或加盐打散 |
| 向量化没生效 | 不兼容算子 / UDF | EXPLAIN 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 的设计提炼成可迁移的原则,其实只有三条:
- 用"逻辑代数 + 可插拔物理约定"解耦编译与执行。这决定了同一个优化器可以服务批处理、流处理、OLAP 三种完全不同的运行时。今天 Flink 的 Table Planner、Doris 的优化器、Paimon 的查询层,本质上都是这条思路的再实现。
- 用代价模型 + importance 剪枝替代规则顺序的玄学。规则该不该应用、什么时候应用,交给
RelOptCost与RuleQueue的 importance 决定,而不是靠人排规则顺序。Hive 早期版本规则硬编码顺序所积累的历史包袱,正是这条原则的反面教材。 - Trait 是物理属性的唯一真源。分布、排序、实现引擎全部编码为 trait,让"哪些操作可以免 shuffle/免排序"变成一个可被规划器推理的命题,而不是散落在各处的 if-else。
真正理解查询编译器,不是为了背下 CoreRules.FILTER_INTO_JOIN 这个名字,而是理解一件事:当"正确性"由关系代数等价变换保证之后,所有的工程精力都可以投入到"如何在不可预测的负载下保持可预测的性能"这一件事上。这套思路,从 Calcite 到 Flink 到任何一个新的 OLAP 引擎,一字未改。

发表评论 取消回复