Flink SQL原理之SQL执行流程
Flink使用Calcite实现了SQL的解析、转换、执行计划优化和转换,那么,Flink SQL是如何执行的呢?
Calcite处理SQL的流程
- SQL解析(SQL -> SqlNode): 将SQL解析为AST(抽象语法树, Calcite中用SqlNode表示)
- SqlNode验证(SqlNode -> SqlNode) : 根据元数据信息(表名、字段名、函数名和数据类型等)进行语法验证
- 语义分析(SqlNode -> RelNode/RexNode: relational expression): 根据SqlNode与元数据信息构建RelNode树,也就是逻辑计划(Logical Plan)
- 逻辑计划优化(RelNode->RelNode): 优化器的核心,Calcite提供了两种Planner(HepPlanner和VolcanoPlanner),按照相应的规则(Rule)进行优化
- HepPlanner: 启发式优化器,RBO, 按照规则匹配,直到最大次数或遍历后不再match rule
- VolcanoPlanner: CBO,一直迭代,直到找到cost最小的plan
- 生成物理执行计划: 将Plan映射为Flink Graph
Flink SQL的执行流程
测试样例
MySQL实体表
1 | CREATE TABLE IF NOT EXISTS t_rt_agg_result( |
数据源
1 | CREATE TABLE orders ( |
假设选用Blink planner运行两个SQL:
- 创建MySQL表的DDL语句
- 从Source消费数据写入MySQL的DML语句
1 | CREATE TABLE agg_result ( |
1 | INSERT INTO MyUserTable |
执行下面MySQL
1 | tEnv.sqlUpdate(ddlSql); |
1 | == Abstract Syntax Tree == |
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
- 结合元数据验证