0%

Flink-SQL原理之Calcite

Flink SQL是Flink API的最顶层抽象,在使用Flink SQL的时候,是否被其便捷和高效惊艳到?

到底一条SQL语句是如何运行起来的?提到Flink SQL,离不开SQL引擎框架 – Calcite。Calcite是面向Hadoop的查询引擎,提供了SQL解析、优化、多重数据源查询的基础框架。Flink借助Calcite实现了SQL解析、优化和graph生成。

img

Flink提供了三层API抽象,每一层API都是便捷性和变现力之前的权衡,来应对不同的计算场景。

Flink SQL是最顶层的抽象,这层抽象在语义和程序表达式上都类似于Table API,但是程序实现都是SQL表达式:

  1. SQL和Table API都是遵循关系模型:元数据和操作
    • 元数据类似于关系型数据库中的schema
    • 操作类似于关系型数据库中的操作, 如select、project、join、group-by和aggregate
  2. SQL和Table API都是声明式定义,而不会执行执行的具体代码
  3. 简洁,通过UDF扩展,但比core API的表达能力差
  4. 执行之前,通过优化器中的优化规则对用户编写的表达式进行优化
  5. TableEnvironment是Flink SQL和Table API的入口,可以无缝切换到DataStream/DataSet, 允许混用。

Calcite

Calcite简介

是一个动态数据的管理框架,可以用来构建数据库系统的语法解析模块

  • 不包含数据存储、数据处理等功能
  • 可以通过编写 Adaptor 来扩展功能,以支持不同的数据处理平台
  • Flink SQL 使用并对其扩展以支持 SQL 语句的解析和验证

Calcite提供了SQL parser、SQL validation、Query optimizer、SQL generator和Data federator

查询的执行过程

分四步:

img

  1. Parse(SQL -> SqlNode ): 使用JavaCC生成的parser来转换查询
  2. validate(SqlNode -> SqlNode ): 通过元数据验证查询
  3. optimize: 逻辑计划优化和转换为物理计划
    • 语义分析(SqlNode -> RelNode/RexNode ): 根据 SqlNode 及元信息构建 RelNode 树,也就是最初版本的逻辑计划(Logical Plan);
    • 逻辑计划优化(RelNode -> RelNode ): 优化器的核心,根据前面生成的逻辑计划按照相应的规则(Rule)进行优化;
  4. execute: 物理计划转换为应用框架的执行逻辑(如Flink的graph)

Catalog

定义Calcite查询的命名空间:

Schema : schematable的集合,可以任意嵌套

Table: 代表单个数据集,字段定义为RelDataType

RelDataType: 代表数据集中的字段,支持所有的SQL数据类型,包括结构体

Statistic: 提供用于优化的表统计信息, 如行数、分布信息、是否为key

image-20210418173928391

SQL parser

  • LL(K) parser通过JavaCC(Java Compiler Compiler)写的
  • 输入的查询转换为AST(abstract syntax tree)
  • SqlNode代表Token
  • SqlNode可以通过unparse方法转换会SQL

image-20210418181320343

SqlNode代表AST中的一个节点

image-20210418181402108

SqlDialect代表特定数据库的方言规则

Query optimizer

查询计划(Query plans)代表执行一个查询必须的步骤

image-20210418184009836

Query优化:

  • 优化逻辑计划
  • 目标通常是尽量减少计划中必须在早期处理的数据量
  • 将逻辑计划转换为物理计划
  • 物理计划与引擎有关,代表了物理执行过程

常见的优化方法:

RBO 规则名称
列裁剪 column_prune Prune unused fields
子查询去关联 decorrelate
子查询转换为join Convert subqueries to joins
聚合消除 aggregation_eliminate
投影消除 projection_eliminate
最大最小消除 max_min_eliminate
谓词下推 predicate_push_down
外连接消除 outer_join_eliminate
分区裁剪 partition_processor
聚合下推 aggregation_push_down
TopN 下推 topn_push_down
Join 重排序 join_reorder

image-20210418190335990

核心概念

  1. 关系代数(Relational algebra):即关系表达式。它们通常以动词命名,例如 Sort, Join, Project, Filter, Scan, Sample.
  2. 行表达式(Row expressions):例如 RexLiteral (常量), RexVariable (变量), RexCall (调用) 等,例如投影列表(Project)、过滤规则列表(Filter)、JOIN 条件列表和 ORDER BY 列表、WINDOW 表达式、函数调用等。使用 RexBuilder 来构建行表达式。
  3. 表达式有各种特征(Trait):使用 Trait 的 satisfies() 方法来测试某个表达式是否符合某 Trait 或 Convention.
  4. 转化特征(Convention):属于 Trait 的子类,用于转化 RelNode 到具体平台实现(可以将下文提到的 Planner 注册到 Convention 中). 例如 JdbcConvention,FlinkConventions.DATASTREAM 等。同一个关系表达式的输入必须来自单个数据源,各表达式之间通过 Converter 生成的 Bridge 来连接。
  5. 规则(Rules):用于将一个表达式转换(Transform)为另一个表达式。它有一个由 RelOptRuleOperand 组成的列表来决定是否可将规则应用于树的某部分。

Planner(规划器)

规划器(Planner) :即请求优化器,它可以根据一系列规则和成本模型(例如基于成本的优化模型 VolcanoPlanner、启发式优化模型 HepPlanner)来将一个表达式转为语义等价(但效率更优)的另一个表达式。

HepPlanner(启发式优化模型)

  • 与Spark优化器类似的启发式优化器
  • 启发式优化比CBO要快速
  • 如果规则对计划做出相反的改变,则存在无限递归的风险

VolcanoPlanner(基于成本的优化模型)

  • 遍历所有的规则,选择代价最小的计划
  • 代价是通过关系表达式提供的
  • 并不是所有可能的计划都会计算
  • 当经过指定的迭代未显著提升将停止优化
  • 代价包括行数、I/O和CPU
  • Statistics用来提高代价评估的准确性
  • Calcite提供了工具来在统计资源消耗

Calcite中Flink中的重要作用

在这里插入图片描述

在Flink中,Calcite扮演着重要的角色:

  • 以Calcite Catalog为核心,上面承载了Table/SQL API
  • Flink SQL和Table的代码最后生成Calcite Logic Plan(SqlNode Tree)
  • 随后验证、优化为 RelNode 树,
  • 最终通过 Rules(规则)和 Convention(转化特征)生成具体的 DataSet Plan(批处理)或 DataStream Plan(流处理),即 Flink 算子构成的处理逻辑。

在这里插入图片描述

Table / SQL API 的编程框架如下:

  • 通过 TableEnvironment 配置 CalciteConfig 对象,自动设置 SQL & Table API 默认处理参数。

  • 使用 registerTableSource() 来将一个 TableSource 注册到 rootSchema. 后续可以通过 scan() 获取此 Table 并调用各种 Table API 进行处理。

  • 接下可以调用 sqlQuery() 和 sqlUpdate() 方法来使用 SQL 语句进行数据处理。

Planner接口: 解析SQL,转换为Transformation

Executor接口: 将Planner转换的Transformation生成streamGraph并执行

Parser接口: 负责SQL解析,parse方法将SQL语句转换为Operation数组

  • 通过Calcite将SQL解析为SqlNode
  • 根据SqlNode的类型,将SqlNode转换为Operation数组

DDL语句的转换过程:

SqlNode转换为RelNode:

  • 推断Table类型
  • 推断计算列
  • 推断watermark分配

SQL转换:

  1. Operation->RelNode
  2. 优化RelNode
  3. RelNode->ExecNode
  4. ExecNode->Transformation算子

死码



[参考文献]

  1. [Parsing database Query with Apache Calcite](Parsing database Query with Apache Calcite - Knoldus Blogs)
  2. Apache Calcite: One planner fits all (slideshare.net)