flink sql与calcite

[toc]

基于Flink1.4.2版本分析flink与calcite结合构建的flink sql模块。

Flink SQL是现在Flink社区中着重发展的一个模块,我理解主要原因是因为

  1. SQL是一门发展很有的通用的描述性语言,接入门槛较低
  2. 有希望在sql层面实现流批计算的统一
  3. 能够通过sql优化器内置优化能力,避免需要每个用户方需要理解低阶任务的调优,屏蔽实现细节

概述

Flink SQL的总体执行流程为:

  • SELECT查询语句经过caclite parse成SqlNode
  • SqlNode经过validate校验
  • SqlNode经过calcite转化为relNode
  • Insert语句将relNode经过calcite的优化和转化成FlinkRelNode
  • 将相应的FlinkRelNode和codeGen生成的Function结合生成相应的执行算子

下面以一个查询sql来讲解整体流程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
val stream = env
.fromCollection(data)
.assignTimestampsAndWatermarks(
new TimestampAndWatermarkWithOffset[(Long, String, String)](0L))
val table = stream.toTable(tEnv, 'a, 'b, 'c, 'rowtime.rowtime)

tEnv.registerTable("T1", table)

val sqlQuery = "SELECT c, COUNT(*), COUNT(1), COUNT(b) FROM T1 " +
"GROUP BY TUMBLE(rowtime, interval '5' SECOND), c"

val result = tEnv.sqlQuery(sqlQuery).toAppendStream[Row]
result.addSink(new StreamITCase.StringSink[Row])

env.execute()

parse

SqlNode

1
2
3
val parser: SqlParser = SqlParser.create(sql, parserConfig)
val sqlNode: SqlNode = parser.parseStmt
sqlNode

SqlNode

SqlNode表示的是一颗sql解析树,由于Flink暂时只支持SELECT查询,所以我们这里得到的其实是一个SqlSelect实例,SqlSelect是一个SqlCall,SqlCall继承自SqlNode,每一个无叶子节点的节点就是一个Sqlcall,常见的SqlNode的子类就是,SqlKind是所有SqlNode类型的枚举类:

  • SqlCall 表示一个树的无叶子节点的调用,例如图中的Count(*)
  • SqlNodeList 表示SqlNode的集合
  • SqlIdentifer 表示某个标识符

SqlOperator

SqlNode的成员方法:

1
2
3
4
5
6
7
8
public List<SqlNode> getOperandList() {
return ImmutableNullableList.of(keywordList, selectList, from, where,
groupBy, having, windowDecls, orderBy, offset, fetch);
}

public SqlOperator getOperator() {
return SqlSelectOperator.INSTANCE;
}

getOperator返回的是这个是个什么操作,operands得到的运算对象。每一个SqlNode是由作用于一系列SqlNode的SqlOperator组成,SqlFunction也是一种SqlOperator. 这里SqlSelect node是SqlSelectOperator作用于以下的SqlNode节点

1
2
3
4
5
6
7
8
9
10
11
SqlNodeList keywordList;
SqlNodeList selectList;
SqlNode from;
SqlNode where;
SqlNodeList groupBy;
SqlNode having;
SqlNodeList windowDecls;
SqlNodeList orderBy;
SqlNode offset;
SqlNode fetch;
SqlMatchRecognize matchRecognize;

rel

1
2
3
4
5
6
7
val rexBuilder: RexBuilder = createRexBuilder
val cluster: RelOptCluster = FlinkRelOptClusterFactory.create(planner, rexBuilder)
val config = SqlToRelConverter.configBuilder()
.withTrimUnusedFields(false).withConvertTableAccess(false).build()
val sqlToRelConverter: SqlToRelConverter = new SqlToRelConverter(
new ViewExpanderImpl, validator, createCatalogReader, cluster, convertletTable, config)
root = sqlToRelConverter.convertQuery(validatedSqlNode, false, true)

这个过程是将SqlNode转化为RelNode的过程,RelNode表示关系型表达式,代表的是对数据的一个操作常见的有Project,Scan,Filter,Join等。通过explain可以看到相应的逻辑执行计划,以下还包括优化后的物理执行计划的一部分

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
 == Abstract Syntax Tree ==
LogicalProject(c=[$1], EXPR$1=[$2], EXPR$2=[$2], EXPR$3=[$3])
LogicalAggregate(group=[{0, 1}], EXPR$1=[COUNT()], EXPR$3=[COUNT($3)])
LogicalProject($f0=[TUMBLE($3, 5000)], c=[$2], $f2=[1], b=[$1])
LogicalTableScan(table=[[T1]])

== Optimized Logical Plan ==
DataStreamCalc(select=[c, EXPR$1, EXPR$1 AS EXPR$2, EXPR$3])
DataStreamGroupWindowAggregate(groupBy=[c], window=[TumblingGroupWindow('w$, 'rowtime, 5000.millis)], select=[c, COUNT(*) AS EXPR$1, COUNT(b) AS EXPR$3])
DataStreamCalc(select=[rowtime, c, 1 AS $f2, b])
DataStreamScan(table=[[_DataStreamTable_0]])

== Physical Execution Plan ==
Stage 1 : Data Source
content : collect elements with CollectionInputFormat

Stage 2 : Operator
content : Timestamps/Watermarks
ship_strategy : FORWARD

Stage 3 : Operator
content : from: (a, b, c, rowtime)
ship_strategy : FORWARD

Stage 4 : Operator
content : select: (rowtime, c, 1 AS $f2, b)
ship_strategy : FORWARD

Stage 5 : Operator
content : time attribute: (rowtime)
ship_strategy : FORWARD

Stage 7 : Operator
content : groupBy: (c), window: (TumblingGroupWindow('w$, 'rowtime, 5000.millis)), select: (c, COUNT(*) AS EXPR$1, COUNT(b) AS EXPR$3)
ship_strategy : HASH

Stage 8 : Operator
content : select: (c, EXPR$1, EXPR$1 AS EXPR$2, EXPR$3)
ship_strategy : FORWARD

Flink sql查询的时候只做到这里的LogicalNode生成之后就完成了,等待sink才会触发下一步优化和转化逻辑。

RelNode,RexNode

RelNode的实现类有LogicalProject,LogicalScan等表示的是数据处理方式,rexnode表示的是行表达式,是包含在一个RelNode中的,RexNode类中的exps字段就存储了相应的数据操作所需要的行表达式

参考以下讨论:

Difference between sqlnode and relnode and rexnode
https://www.mail-archive.com/dev@calcite.apache.org/msg01674.html

RexTraits, RelTraitDef

这个表示的是一个RelNode的物理特性,用于在convertRule中使用

在一个ConverterRule中的convert

1
2
val scan: FlinkLogicalNativeTableScan = rel.asInstanceOf[FlinkLogicalNativeTableScan]
val traitSet: RelTraitSet = rel.getTraitSet.replace(FlinkConventions.DATASTREAM)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public RelTraitSet replace(
RelTrait trait) {
// Quick check for common case
if (containsShallow(traits, trait)) {
return this;
}
final RelTraitDef traitDef = trait.getTraitDef();
int index = findIndex(traitDef);
if (index < 0) {
// Trait is not present. Ignore it.
return this;
}

return replace(index, trait);
}

实际上是把某一类relTraitsDef的trait实现更换掉,相当于变更了RelNode的物理实现,planner转化过程应该只是trait的转化过程以及相应的RelNode的物理化的过程,按照论文中的解释:

Traits. Calcite does not use different entities to represent logical
and physical operators. Instead, it describes the physical properties
associated with an operator using traits. These traits help the optimizer evaluate the cost of different alternative plans. Changing a
trait value does not change the logical expression being evaluated,
i.e., the rows produced by the given operator will still be the same

因此在Flink中的转化StreamTableEnvironment#optimize,match之后根据HepPlanner和VocanoPlanner进行convert转化成物理算子,例如将

LogicalJoin(RelNode) -> DataStreamJoin(FlinkRelNode) -> translateToPlan -> NonWindowJoin (runtime)

总结

Flink SQL具体在flink中的实现分为LogicalPlan层,经过应用rule optimize之后的RelNode层,例如: DataStreamJoin, 再通过translate的时候code generator以及调用相应的runtime层的具体算子实现(这个对应的是physical plan的翻译),以上就完成了从SQL到Flink执行计划的翻译。从整个流程看calcite全程参与,使用方式非常的方便,足见整个calcite框架的扩展性做的很好。
Flink SQL中还有许多其他的细节:SQL中的回撤消息,join,distinct的具体通用算子的实现,还有sql优化,分流,ddl的实现,antlr的实现,窗口聚合等等这些实现细节后文再具体分析。

参考

介绍calcite与flink sql比较好的几篇文章:

https://zhuanlan.zhihu.com/p/48735419
https://arxiv.org/pdf/1802.10233.pdf calcite的论文
https://zhuanlan.zhihu.com/p/51221350
https://zhuanlan.zhihu.com/p/58249033
https://zhuanlan.zhihu.com/p/59643962


http://matt33.com/2019/03/17/apache-calcite-planner/
http://matt33.com/2019/03/07/apache-calcite-process-flow/
https://www.slideshare.net/julianhyde/costbased-query-optimization-in-apache-phoenix-using-apache-calcite?qid=b7a1ca0f-e7bf-49ad-bc51-0615ec8a4971&v=&b=&from_search=4


https://issues.apache.org/jira/browse/FLINK-7146 Flink SQL DDL支持
https://docs.google.com/document/d/1TTP-GCC8wSsibJaSUyFZ_5NBAHYEB1FVmPpP7RgDGBA/edit#heading=h.wpsqidkaaoil doc

https://github.com/TatianaJin/calcite_playground/wiki/Query-Planning-&-Optimization-II.a:-VolcanoPlanner-Basics

谢谢支持