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 reviewed guide

The reviewed guide body for this track is currently maintained in Chinese.

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


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 会拆分成多个类。


下一步