Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 07

Table/SQL planner and Calcite

Connect Table/SQL, Calcite planning, optimisation, and Flink physical execution plans.

Reading
60 min
Track
Flink internals
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 60 分钟 前置阅读: doc-01(Job 提交) 下一次阅读: doc-12(Connector)


Table/SQL Planner — Calcite 集成深度解析 图 01

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作用示例
FilterPushDownCalcite谓词下推到 SourceWHERE age > 18 → Source 端过滤
ColumnPruningCalcite只读需要的列SELECT name → 不读 age 列
JoinReorderCalcite基于统计信息重排 Join 顺序小表 × 大表 → 大表 × 小表
FlinkAggregateReduceFunctionsRuleFlink化简聚合函数COUNT(*) + COUNT(col) → 只算 COUNT(*)
FlinkExpandWindowToTumbleWindowRuleFlink扩展窗口为 Tumbling WindowSession→Tumbling 优化
FlinkSinkPushDownFlinkSink 下推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-written100% baseline
开发效率高 (声明式)低 (命令式)
优化空间自动 (规则)手动
可维护性

4. 源码导航(完整版)

文件关键类/方法职责
flink-table/flink-table-planner/.../planner/FlinkPlannerImpl.javaoptimize(), validate()优化器入口
flink-table/flink-table-planner/.../calcite/FlinkCalciteContext.javaCalcite 上下文
flink-table/flink-table-planner/.../plan/nodes/exec/ExecNode 子类物理执行节点
flink-table/flink-table-planner/.../codegen/CodeGenerator.javagenerate(), compile()代码生成器
flink-table/flink-table-planner/.../codegen/JaninoCompiler.javacompile()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 会拆分成多个类。


下一步