0%

Flink SQL原理之SQL执行流程

Flink SQL原理之SQL执行流程

Flink使用Calcite实现了SQL的解析、转换、执行计划优化和转换,那么,Flink SQL是如何执行的呢?

Calcite处理SQL的流程

  1. SQL解析(SQL -> SqlNode): 将SQL解析为AST(抽象语法树, Calcite中用SqlNode表示)
  2. SqlNode验证(SqlNode -> SqlNode) : 根据元数据信息(表名、字段名、函数名和数据类型等)进行语法验证
  3. 语义分析(SqlNode -> RelNode/RexNode: relational expression): 根据SqlNode与元数据信息构建RelNode树,也就是逻辑计划(Logical Plan)
  4. 逻辑计划优化(RelNode->RelNode): 优化器的核心,Calcite提供了两种Planner(HepPlanner和VolcanoPlanner),按照相应的规则(Rule)进行优化
    • HepPlanner: 启发式优化器,RBO, 按照规则匹配,直到最大次数或遍历后不再match rule
    • VolcanoPlanner: CBO,一直迭代,直到找到cost最小的plan
  5. 生成物理执行计划: 将Plan映射为Flink Graph

测试样例

MySQL实体表

1
2
3
4
5
6
7
8
9
10
CREATE TABLE IF NOT EXISTS t_rt_agg_result(
ks int(11) not null primary key auto_increment ,
biz int NOT NULL COMMENT 'biz code',
ei varchar(10) NOT NULL COMMENT 'first level event name',
sei varchar(10) NULL COMMENT 'second level event name',
uv long NULL COMMENT 'user count',
pv long NULL,
time_id varchar(10) NOT NULL,
update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
)

数据源

1
2
3
4
5
6
7
8
9
10
11
12
CREATE TABLE orders (
biz int,
ei STRING,
sei STRING,
ui STRing,
etime timestamp(3),
watermark for etime as etime - interval '5' second
) WITH (
'connector.type' = 'filesystem',
'connector.path'='/Users/zhangzuofeng1/orders.csv',
'format.type'='csv'
)

假设选用Blink planner运行两个SQL:

  • 创建MySQL表的DDL语句
  • 从Source消费数据写入MySQL的DML语句
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
CREATE TABLE agg_result (
biz int,
ei STRING,
sei STRING,
uv bigint,
pv bigint,
time_id STRING
) WITH (
'connector.type' = 'jdbc',
'connector.driver' = 'com.mysql.cj.jdbc.Driver',
'connector.url' = 'jdbc:mysql://localhost:3306/fdata',
'connector.table' = 't_agg_result',
'connector.username' = 'root',
'connector.password' = 'abc.ABC.123'
)
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
INSERT INTO MyUserTable
SELECT
biz,
ei,
sei,
count(DISTINCT ui) AS uv,
count(*) AS pv,
DATE_FORMAT(TUMBLE_START(etime, INTERVAL '1' MINUTE), 'yyyyMMddHHmm') AS timeId
FROM orders
WHERE ui IS NOT NULL
GROUP BY biz,
ei,
sei,
TUMBLE(etime, INTERVAL '1' MINUTE)

insert into agg_result
SELECT
biz,
ei,
sei,
count(DISTINCT ui) AS uv,
count(*) AS pv,
DATE_FORMAT(TUMBLE_START(etime, INTERVAL '1' MINUTE), 'yyyyMMddHHmm') AS timeId
FROM orders
WHERE ui IS NOT NULL
GROUP BY biz,
ei,
sei,
TUMBLE(etime, INTERVAL '1' MINUTE)

执行下面MySQL

1
2
tEnv.sqlUpdate(ddlSql);
tEnv.sqlUpdate(dmlSql);
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
== Abstract Syntax Tree ==
LogicalProject(ei=[$0], sei=[$1], ui=[$2], etime=[$3])
+- LogicalFilter(condition=[<>($0, _UTF-16LE'a')])
+- LogicalTableScan(table=[[default_catalog, default_database, Unregistered_TableSource_1298380324, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]])

== Optimized Logical Plan ==
Calc(select=[ei, sei, ui, etime], where=[<>(ei, _UTF-16LE'a':VARCHAR(10) CHARACTER SET "UTF-16LE")])
+- TableSourceScan(table=[[default_catalog, default_database, Unregistered_TableSource_1298380324, source: [CsvTableSource(read fields: ei, sei, ui, etime)]]], fields=[ei, sei, ui, etime])

== Physical Execution Plan ==
Stage 1 : Data Source
content : Source: Custom File source

Stage 2 : Operator
content : CsvTableSource(read fields: ei, sei, ui, etime)
ship_strategy : REBALANCE

Stage 3 : Operator
content : SourceConversion(table=[default_catalog.default_database.Unregistered_TableSource_1298380324, source: [CsvTableSource(read fields: ei, sei, ui, etime)]], fields=[ei, sei, ui, etime])
ship_strategy : FORWARD

Stage 4 : Operator
content : Calc(select=[ei, sei, ui, etime], where=[(ei <> _UTF-16LE'a':VARCHAR(10) CHARACTER SET "UTF-16LE")])
ship_strategy : FORWARD
1
{"nodes":[{"id":10,"type":"Source: Custom File source","pact":"Data Source","contents":"Source: Custom File source","parallelism":1},{"id":11,"type":"CsvTableSource(read fields: biz, ei, sei, ui, etime)","pact":"Operator","contents":"CsvTableSource(read fields: biz, ei, sei, ui, etime)","parallelism":12,"predecessors":[{"id":10,"ship_strategy":"REBALANCE","side":"second"}]},{"id":12,"type":"SourceConversion(table=[default_catalog.default_database.orders, source: [CsvTableSource(read fields: biz, ei, sei, ui, etime)]], fields=[biz, ei, sei, ui, etime])","pact":"Operator","contents":"SourceConversion(table=[default_catalog.default_database.orders, source: [CsvTableSource(read fields: biz, ei, sei, ui, etime)]], fields=[biz, ei, sei, ui, etime])","parallelism":12,"predecessors":[{"id":11,"ship_strategy":"FORWARD","side":"second"}]},{"id":13,"type":"WatermarkAssigner(rowtime=[etime], watermark=[(etime - 5000:INTERVAL SECOND)])","pact":"Operator","contents":"WatermarkAssigner(rowtime=[etime], watermark=[(etime - 5000:INTERVAL SECOND)])","parallelism":12,"predecessors":[{"id":12,"ship_strategy":"FORWARD","side":"second"}]},{"id":14,"type":"Calc(select=[biz, ei, sei, ui, etime], where=[ui IS NOT NULL])","pact":"Operator","contents":"Calc(select=[biz, ei, sei, ui, etime], where=[ui IS NOT NULL])","parallelism":12,"predecessors":[{"id":13,"ship_strategy":"FORWARD","side":"second"}]},{"id":16,"type":"GroupWindowAggregate(groupBy=[biz, ei, sei], window=[TumblingGroupWindow('w$, etime, 60000)], properties=[w$start, w$end, w$rowtime, w$proctime], select=[biz, ei, sei, COUNT(DISTINCT ui) AS uv, COUNT(*) AS pv, start('w$) AS w$start, end('w$) AS w$end, rowtime('w$) AS w$rowtime, proctime('w$) AS w$proctime])","pact":"Operator","contents":"GroupWindowAggregate(groupBy=[biz, ei, sei], window=[TumblingGroupWindow('w$, etime, 60000)], properties=[w$start, w$end, w$rowtime, w$proctime], select=[biz, ei, sei, COUNT(DISTINCT ui) AS uv, COUNT(*) AS pv, start('w$) AS w$start, end('w$) AS w$end, rowtime('w$) AS w$rowtime, proctime('w$) AS w$proctime])","parallelism":12,"predecessors":[{"id":14,"ship_strategy":"HASH","side":"second"}]},{"id":17,"type":"Calc(select=[biz, ei, sei, uv, pv, (w$start DATE_FORMAT _UTF-16LE'yyyyMMddHHmm') AS timeId])","pact":"Operator","contents":"Calc(select=[biz, ei, sei, uv, pv, (w$start DATE_FORMAT _UTF-16LE'yyyyMMddHHmm') AS timeId])","parallelism":12,"predecessors":[{"id":16,"ship_strategy":"FORWARD","side":"second"}]},{"id":18,"type":"SinkConversionToTuple2","pact":"Operator","contents":"SinkConversionToTuple2","parallelism":12,"predecessors":[{"id":17,"ship_strategy":"FORWARD","side":"second"}]},{"id":19,"type":"Sink: JDBCUpsertTableSink(biz, ei, sei, uv, pv, time_id)","pact":"Data Sink","contents":"Sink: JDBCUpsertTableSink(biz, ei, sei, uv, pv, time_id)","parallelism":12,"predecessors":[{"id":18,"ship_strategy":"FORWARD","side":"second"}]}]}

Parser: Provides methods for parsing SQL objects from a SQL string.

org.apache.flink.table.planner.delegation.PlannerBase

org.apache.flink.table.planner.operations.SqlToOperationConverter#convert

SqlNode 转换为 Operation

  1. 结合元数据验证

https://cloud.tencent.com/developer/article/1803116

https://cloud.tencent.com/developer/article/1697401