
1. Calcite 集成: SqlParser → SqlNode → RelNode
1.1 SQL 编译器管道
SQL: INSERT INTO sink SELECT url, COUNT(*) FROM source GROUP BY url
Phase 1: Parse (Calcite)
SqlParser.parse(sql) → SqlNode AST
SqlNode 树:
SqlInsert
├── target: SqlIdentifier("sink")
└── source: SqlSelect
├── selectList: SqlIdentifier("url"), SqlCall("COUNT", SqlIdentifier("*"))
├── from: SqlIdentifier("source")
└── groupBy: SqlIdentifier("url")
Phase 2: Validate (Calcite + Flink)
SqlValidator.validate(SqlNode) → RelNode (Logical Plan)
RelNode 树:
LogicalSink
└── LogicalAggregate(group={url}, agg#0=COUNT())
└── LogicalTableScan(table=source)
Phase 3: Optimize (Calcite Rules + Flink Extensions)
FlinkPlannerImpl.optimize(RelNode)
应用规则: FilterPushDown → Filter → Source
ColumnPruning → 只读需要的列
JoinReorder → 小表在左
SubQueryRemove → 移除子查询
→ Optimized RelNode
Phase 4: Physical (Flink)
translateToExecNode(Optimized RelNode)
→ ExecNode 树: StreamExecSink → StreamExecGroupAggregate → StreamExecTableSourceScan
Phase 5: Plan (Flink)
translateToPlan(ExecNode)
→ Transformation → StreamGraph → JobGraph
1.2 流批统一的 RelNode → ExecNode 转换
流模式 (Stream): 批模式 (Batch):
LogicalAggregate LogicalAggregate
│ │
StreamExecGroupAggregate BatchExecGroupAggregate
(增量聚合,StateBackend) (全量聚合,Sort-based)
│ │
OneInputTransformation OneInputTransformation
2. 优化规则体系
2.1 Flink 扩展的关键规则
| 规则 | Calcite/Flink | 作用 | 示例 |
|---|---|---|---|
FilterPushDown | Calcite | 谓词下推到 Source | WHERE age > 18 → Source 端过滤 |
ColumnPruning | Calcite | 只读需要的列 | SELECT name → 不读 age 列 |
JoinReorder | Calcite | 基于统计信息重排 Join 顺序 | 小表 × 大表 → 大表 × 小表 |
FlinkAggregateReduceFunctionsRule | Flink | 化简聚合函数 | COUNT(*) + COUNT(col) → 只算 COUNT(*) |
FlinkExpandWindowToTumbleWindowRule | Flink | 扩展窗口为 Tumbling Window | Session→Tumbling 优化 |
FlinkSinkPushDown | Flink | Sink 下推 | DynamicFiltering 下推到 Sink |
2.2 谓词下推示例
原始 SQL:
SELECT * FROM kafka_source WHERE status = 'active' AND age > 18
优化前:
LogicalFilter(status='active' AND age>18)
→ LogicalTableScan(kafka_source)
Kafka Source 读取全部数据 → Filter 在 Flink 端计算
优化后 (FilterPushDown):
LogicalTableScan(kafka_source, pushedFilter='age>18')
Kafka Source 读取时应用 age>18 过滤 (如果 Connector 支持)
→ Flink 端只计算 status='active'
3. CodeGen — 代码生成
3.1 为什么需要 CodeGen
Flink 使用 Janino 编译器 (轻量级 Java 运行时编译器) 将 RelNode 转为 Java 字节码,避免解释执行开销。
// 生成的代码等价于:
public class GeneratedAggFunction {
public void processElement(RowData input) {
// 计算 key
int urlHash = input.getString(0).hashCode();
// 从 State 读取累加器
Long count = (Long) accState.value();
if (count == null) count = 0L;
count++;
// 写回 State
accState.update(count);
// 输出
out.collect(RowData.of(input.getString(0), count));
}
}
// 手写 DataStream 代码等价于:
keyedStream.process(new KeyedProcessFunction() {
private ValueState<Long> countState;
public void processElement(...) {
Long count = countState.value();
countState.update(count + 1);
out.collect(...);
}
});
3.2 CodeGen vs 手写 DataStream
| 维度 | CodeGen (SQL) | 手写 DataStream |
|---|---|---|
| 性能 | ~95-105% of hand-written | 100% baseline |
| 开发效率 | 高 (声明式) | 低 (命令式) |
| 优化空间 | 自动 (规则) | 手动 |
| 可维护性 | 高 | 中 |
4. 源码导航(完整版)
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-table/flink-table-planner/.../planner/FlinkPlannerImpl.java | optimize(), validate() | 优化器入口 |
flink-table/flink-table-planner/.../calcite/FlinkCalciteContext.java | — | Calcite 上下文 |
flink-table/flink-table-planner/.../plan/nodes/exec/ | ExecNode 子类 | 物理执行节点 |
flink-table/flink-table-planner/.../codegen/CodeGenerator.java | generate(), compile() | 代码生成器 |
flink-table/flink-table-planner/.../codegen/JaninoCompiler.java | compile() | Janino 编译 |
flink-table/flink-table-planner/.../optimize/FlinkStreamPrograms.java | — | 优化规则集合 |
flink-table/flink-table-planner/.../rules/ | — | 优化规则定义 |
5. 常见问题 / 面试题
Q1: Flink SQL 和 DataStream API 的性能差异?
A: SQL 经过 Calcite 优化器(谓词下推、列裁剪、Join 重排) + CodeGen(Janino 编译),通常和手写 DataStream 性能相当。CodeGen 生成的代码是直接的无函数调用路径(inline),在某些场景下甚至比手写 DataStream 更快。但高度定制化的逻辑(如复杂状态机)DataStream 更灵活。
Q2: 为什么 Flink 选 Calcite 而不是自研优化器?
A: Calcite 是 Apache 顶级项目,Hive/Spark/Kylin/Druid 等都基于它。自研优化器的开发成本极高(需要完整的 SQL Parser、Validator、优化规则框架、Cost Model)。Calcite 的 Rule 体系易于扩展——Flink 只需要添加流处理特有的规则(FlinkFilterPushDown, FlinkAggregateReduceFunctionsRule 等)。
Q3: CodeGen 有什么局限性?
A: (1) Janino 编译器限制——不支持完整的 Java 8+ 特性(如 Lambda 表达式);(2) 生成的代码长度有限制("code length > 64KB" → 编译失败);(3) 生成的代码难以调试(Source 在 JAR 外,断点不易打)。对于超长方法,Flink 会拆分成多个类。
下一步
- Connector: doc-12-connector-framework.md