- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
导读
本文是 Apache Flink 中 Table API 与 SQL 的核心入门与深度指南,基于仓库中的官方文档 docs/content.zh/docs/dev/table/common.md 展开,并结合flink-table模块的真实源码进行佐证。Table API 和 SQL 集成在同一套 API 中,其核心概念是Table,它同时充当查询的输入与输出。读完本文,你将掌握 Table API/SQL 程序的通用结构、如何创建TableEnvironment、如何在 Catalog 中注册临时表与永久表、如何用 Table API 和 SQL 两种方式查询表、如何将结果输出到 TableSink,以及查询的翻译、优化与解释(explain)机制,从而能够独立搭建从数据源到结果输出的完整 Flink 表处理程序。
Table API 和 SQL 程序的结构
所有用于批处理和流处理的 Table API 和 SQL 程序都遵循相同的模式,可概括为五个步骤:创建TableEnvironment→ 创建源表(source table)→ 创建输出表(sink table)→ 通过 Table API 或 SQL 构建Table查询 → 将结果表输出到 sink。下面的 Java 示例展示了这一通用结构:
import org.apache.flink.table.api.*; import org.apache.flink.connector.datagen.table.DataGenConnectorOptions; // Create a TableEnvironment for batch or streaming execution. // See the "Create a TableEnvironment" section for details. TableEnvironment tableEnv = TableEnvironment.create(/*…*/); // Create a source table tableEnv.createTemporaryTable("SourceTable", TableDescriptor.forConnector("datagen") .schema(Schema.newBuilder() .column("f0", DataTypes.STRING()) .build()) .option(DataGenConnectorOptions.ROWS_PER_SECOND, 100L) .build()); // Create a sink table (using SQL DDL) tableEnv.executeSql("CREATE TEMPORARY TABLE SinkTable WITH ('connector' = 'blackhole') LIKE SourceTable (EXCLUDING OPTIONS) "); // Create a Table object from a Table API query Table table1 = tableEnv.from("SourceTable"); // Create a Table object from a SQL query Table table2 = tableEnv.sqlQuery("SELECT * FROM SourceTable"); // Emit a Table API result Table to a TableSink, same for SQL result TableResult tableResult = table1.insertInto("SinkTable").execute();对应的 Python 版本使用pyflink.table包,通过executeSql创建表、from_path读取表、sql_query执行 SQL、execute_insert输出结果:
from pyflink.table import * # Create a TableEnvironment for batch or streaming execution table_env = ... # see "Create a TableEnvironment" section # Create a source table table_env.executeSql("""CREATE TEMPORARY TABLE SourceTable ( f0 STRING ) WITH ( 'connector' = 'datagen', 'rows-per-second' = '100' ) """) # Create a sink table table_env.executeSql("CREATE TEMPORARY TABLE SinkTable WITH ('connector' = 'blackhole') LIKE SourceTable (EXCLUDING OPTIONS) ") # Create a Table from a Table API query table1 = table_env.from_path("SourceTable").select(...) # Create a Table from a SQL query table2 = table_env.sql_query("SELECT ... FROM SourceTable ...") # Emit a Table API result Table to a TableSink, same for SQL result table_result = table1.execute_insert("SinkTable")注意:示例中
datagen连接器用于无界生成测试数据。其rows-per-second选项在源码中定义于 DataGenConnectorOptions.java,默认值为10000行/秒(见 DataGenConnectorOptionsUtil.java),并支持number-of-rows限定总行数(默认无限、无界生成)。blackhole连接器则是一个只接收数据、不产生任何输出的“黑洞” sink,非常适合快速验证管线正确性。关于 Scala:所有 Flink Scala API 均已弃用(deprecated),并将在未来的 Flink 版本中移除。你仍然可以用 Scala 构建应用,但建议迁移到 DataStream 和/或 Table API 的 Java 版本。
与 DataStream 的集成:Table API 和 SQL 查询可以很容易地集成并嵌入到 DataStream 程序中。请参阅与 DataStream API 集成章节了解如何将 DataStream 与表之间的相互转化。
创建 TableEnvironment
TableEnvironment是 Table API 和 SQL 的核心概念,它负责:
- 在内部的 catalog 中注册
Table - 注册外部的 catalog
- 加载可插拔模块(modules)
- 执行 SQL 查询
- 注册自定义函数(scalar、table 或 aggregation 函数)
DataStream和Table之间的转换(面向StreamTableEnvironment)
Table总是与特定的TableEnvironment绑定,不能在同一条查询中使用不同TableEnvironment中的表,例如对它们进行 join 或 union 操作。
TableEnvironment通过静态方法TableEnvironment.create()创建,该方法在源码 TableEnvironment.java 中提供两种重载:接收EnvironmentSettings或直接接收Configuration。最常用的方式是基于EnvironmentSettings指定执行模式:
import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.TableEnvironment; EnvironmentSettings settings = EnvironmentSettings .newInstance() .inStreamingMode() // 流处理模式 //.inBatchMode() // 批处理模式 .build(); TableEnvironment tEnv = TableEnvironment.create(settings);Python 中通过EnvironmentSettings.in_streaming_mode()和EnvironmentSettings.in_batch_mode()分别创建流式与批式环境:
from pyflink.table import EnvironmentSettings, TableEnvironment # create a streaming TableEnvironment env_settings = EnvironmentSettings.in_streaming_mode() table_env = TableEnvironment.create(env_settings) # create a batch TableEnvironment env_settings = EnvironmentSettings.in_batch_mode() table_env = TableEnvironment.create(env_settings)两种模式的选择决定了查询语义:流式模式以无界数据流为输入、支持窗口与状态化聚合;批式模式则以有界数据集为输入,执行与传统数据库类似的批处理语义。
从 StreamExecutionEnvironment 创建 StreamTableEnvironment
当需要与 DataStream API 互操作时,可以从现有的StreamExecutionEnvironment创建StreamTableEnvironment:
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);Python 版本同样支持:
from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import StreamTableEnvironment s_env = StreamExecutionEnvironment.get_execution_environment() t_env = StreamTableEnvironment.create(s_env)这样创建的StreamTableEnvironment既可以通过 Table API/SQL 编写关系查询,又可以随时将Table与DataStream互相转换,实现关系查询与流式算子(如 keyBy、map、process)的混合编程。
在 Catalog 中创建表
TableEnvironment维护着一个由标识符(identifier)创建的表 catalog 映射。标识符由三个部分组成:catalog 名称、数据库名称以及对象名称。如果 catalog 或数据库没有指明,就会使用当前默认值(参见下文扩展表标识符)。
Table可以是虚拟的(视图VIEWS)也可以是常规的(表TABLES):
- 视图
VIEWS可以从已经存在的Table中创建,一般是 Table API 或 SQL 的查询结果; - 表
TABLES描述的是外部数据,例如文件、数据库表或消息队列。
临时表(Temporary Table)和永久表(Permanent Table)
表可以是临时的,与单个 Flink 会话(session)的生命周期相关;也可以是永久的,在多个 Flink 会话和集群(cluster)中可见。
- 永久表需要 catalog(例如 Hive Metastore)来维护表的元数据。一旦永久表被创建,它将对任何连接到该 catalog 的 Flink 会话可见且持续存在,直至被明确删除。
- 临时表通常保存于内存中,仅在创建它们的 Flink 会话持续期间存在,对其它会话不可见。它们不与任何 catalog 或数据库绑定,但可以在一个命名空间(namespace)中创建。即使它们对应的数据库被删除,临时表也不会被删除。
屏蔽(Shadowing)
可以使用与已存在的永久表相同的标识符去注册临时表。此时临时表会屏蔽永久表,并且只要临时表存在,永久表就无法访问——所有使用该标识符的查询都将作用于临时表。
屏蔽机制对实验(experimentation)非常有用:可以先对一个临时表执行完全相同的查询,例如只包含一个子集的数据,或者数据是不确定的;一旦验证了查询的正确性,就可以对实际的生产表进行查询,无需修改任何 SQL。
创建表
虚拟表(Virtual Tables)
在 SQL 的术语中,Table API 的对象对应于视图(虚拟表)。它封装了一个逻辑查询计划,可以通过以下方式在 catalog 中创建:
// get a TableEnvironment TableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section // table is the result of a simple projection query Table projTable = tableEnv.from("X").select(...); // register the Table projTable as table "projectedTable" tableEnv.createTemporaryView("projectedTable", projTable);Python 中对应的注册方法为register_table:
proj_table = table_env.from_path("X").select(...) table_env.register_table("projectedTable", proj_table)注意:从传统数据库系统的角度来看,Table对象与VIEW视图非常像。也就是说,定义了Table的查询没有被优化,而且会被内嵌到另一个引用了这个注册表的查询中。如果多个查询都引用了同一个注册表,那么它会被内嵌到每个查询中并执行多次——注册表的结果不会被共享。
Connector Tables
另一种创建TABLE的方式是通过 connector 声明。Connector 描述了存储表数据的外部系统,例如 Apache Kafka 或常规的文件系统都可以通过这种方式来声明。这类表既可以通过 Table API 的TableDescriptor直接创建,也可以切换到 SQL DDL 创建:
// Using table descriptors final TableDescriptor sourceDescriptor = TableDescriptor.forConnector("datagen") .schema(Schema.newBuilder() .column("f0", DataTypes.STRING()) .build()) .option(DataGenConnectorOptions.ROWS_PER_SECOND, 100L) .build(); tableEnv.createTable("SourceTableA", sourceDescriptor); // 永久表(需要 catalog 支持) tableEnv.createTemporaryTable("SourceTableB", sourceDescriptor); // 临时表 // Using SQL DDL tableEnv.executeSql("CREATE [TEMPORARY] TABLE MyTable (...) WITH (...)");TableDescriptor.forConnector(...)以编程方式构建表的连接器类型、schema 与连接选项,DataGenConnectorOptions.ROWS_PER_SECOND即源码中定义的ConfigOption(配置键rows-per-second),使用强类型常量可以避免手写字符串拼写错误。SQL DDL 则提供了与CREATE TABLE一致的声明式语法,二者创建的表可以互换使用。
扩展表标识符
表总是通过三元标识符注册,包括 catalog 名、数据库名和表名。用户可以指定一个 catalog 和数据库作为“当前 catalog”和“当前数据库”,这样三元标识符的前两个部分就可以省略;未指定时使用当前的 catalog 和当前数据库。用户也可以通过 Table API 或 SQL 切换当前的 catalog 和当前的数据库。
标识符遵循 SQL 标准,因此使用时需要用反引号(`)进行转义。以下 Java 示例演示了不同标识符写法:
TableEnvironment tEnv = ...; tEnv.useCatalog("custom_catalog"); tEnv.useDatabase("custom_database"); Table table = ...; // register the view named 'exampleView' in the catalog named 'custom_catalog' // in the database named 'custom_database' tableEnv.createTemporaryView("exampleView", table); // register the view named 'exampleView' in the catalog named 'custom_catalog' // in the database named 'other_database' tableEnv.createTemporaryView("other_database.exampleView", table); // register the view named 'example.View' in the catalog named 'custom_catalog' // in the database named 'custom_database' tableEnv.createTemporaryView("`example.View`", table); // register the view named 'exampleView' in the catalog named 'other_catalog' // in the database named 'other_database' tableEnv.createTemporaryView("other_catalog.other_database.exampleView", table);从源码接口看,useCatalog(String)与useDatabase(String)均定义于 TableEnvironment.java 的TableEnvironment接口中,是切换当前命名空间的官方 API。Python 中对应为use_catalog(...)与use_database(...),注册方法为create_temporary_view(...)。
查询表
Table API
Table API 是关于 Java 和 Scala 的集成语言式查询 API。与 SQL 相反,Table API 的查询不是由字符串指定,而是在宿主语言中逐步构建。
Table API 基于Table类,该类表示一个表(流或批处理),并提供使用关系操作的方法。这些方法返回一个新的Table对象,该对象表示对输入 Table 进行关系操作的结果。一些关系操作由多个方法调用组成,例如table.groupBy(...).select(...),其中groupBy(...)指定table的分组,而select(...)是在分组上的投影。
文档 Table API 说明了所有流处理和批处理表支持的 Table API 算子。以下示例展示了一个简单的 Table API 聚合查询——计算法国所有客户的收入:
// get a TableEnvironment TableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section // register Orders table // scan registered Orders table Table orders = tableEnv.from("Orders"); // compute revenue for all customers from France Table revenue = orders .filter($("cCountry").isEqual("FRANCE")) .groupBy($("cID"), $("cName")) .select($("cID"), $("cName"), $("revenue").sum().as("revSum")); // emit or convert Table // execute queryPython 版本使用col(...)引用列:
orders = table_env.from_path("Orders") revenue = orders \ .filter(col('cCountry') == 'FRANCE') \ .group_by(col('cID'), col('cName')) \ .select(col('cID'), col('cName'), col('revenue').sum.alias('revSum'))注意$("...")是基于字符串的表达式引用(fromDataStream转换或已有表均可使用),filter/groupBy/select等算子都返回新的Table,因而可以无限链式组合。
SQL
Flink SQL 是基于实现了 SQL 标准的 Apache Calcite 描述了 Flink 对流处理和批处理表的 SQL 支持。
下面的示例演示了如何指定查询并将结果作为Table对象返回:
// get a TableEnvironment TableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section // register Orders table // compute revenue for all customers from France Table revenue = tableEnv.sqlQuery( "SELECT cID, cName, SUM(revenue) AS revSum " + "FROM Orders " + "WHERE cCountry = 'FRANCE' " + "GROUP BY cID, cName" ); // emit or convert Table // execute query如下示例展示了如何指定一个更新查询(insert query),将查询的结果插入到已注册的表中:
// get a TableEnvironment TableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section // register "Orders" table // register "RevenueFrance" output table // compute revenue for all customers from France and emit to "RevenueFrance" tableEnv.executeSql( "INSERT INTO RevenueFrance " + "SELECT cID, cName, SUM(revenue) AS revSum " + "FROM Orders " + "WHERE cCountry = 'FRANCE' " + "GROUP BY cID, cName" );Python 对应为sql_query(...)与execute_sql(...)。可以看到,sqlQuery用于读取(返回Table对象),而executeSql用于执行 DDL/DML(如CREATE TABLE、INSERT INTO),这是两者在用法上的关键区别。
混用 Table API 和 SQL
Table API 和 SQL 查询的混用非常简单,因为它们都返回Table对象:
- 可以在 SQL 查询返回的
Table对象上定义 Table API 查询; - 在
TableEnvironment中注册的结果表可以在 SQL 查询的FROM子句中引用,通过这种方法就可以在 Table API 查询的结果上定义 SQL 查询。
输出表
Table通过写入TableSink输出。TableSink是一个通用接口,用于支持多种文件格式(如 CSV、Apache Parquet、Apache Avro)、存储系统(如 JDBC、Apache HBase、Apache Cassandra、Elasticsearch)或消息队列系统(如 Apache Kafka、RabbitMQ)。
- 批处理
Table只能写入BatchTableSink; - 流处理
Table需要指定写入AppendStreamTableSink、RetractStreamTableSink或UpsertStreamTableSink之一(取决于结果流的更新模式)。
请参考文档 Table Sources & Sinks 获取更多关于可用 Sink 的信息以及如何自定义DynamicTableSink。
方法Table.insertInto(String tableName)定义了一个完整的端到端管道,将源表中的数据传输到一个被注册的输出表中。该方法通过名称在 catalog 中查找输出表,并确认Tableschema 与输出表 schema 一致。从源码看,insertInto返回TablePipeline对象(见 Table.java),可以通过TablePipeline.explain()和TablePipeline.execute()分别解释和执行一个数据流管道。
下面的示例演示如何输出Table,其中输出表使用filesystem连接器 + CSV 格式,并以|作为字段分隔符:
// get a TableEnvironment TableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section // create an output Table final Schema schema = Schema.newBuilder() .column("a", DataTypes.INT()) .column("b", DataTypes.STRING()) .column("c", DataTypes.BIGINT()) .build(); tableEnv.createTemporaryTable("CsvSinkTable", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "/path/to/file") .format(FormatDescriptor.forFormat("csv") .option("field-delimiter", "|") .build()) .build()); // compute a result Table using Table API operators and/or SQL queries Table result = ...; // Prepare the insert into pipeline TablePipeline pipeline = result.insertInto("CsvSinkTable"); // Print explain details pipeline.printExplain(); // emit the result Table to the registered TableSink pipeline.execute();Python 中直接通过result.execute_insert("CsvSinkTable")一步完成插入与执行。
这里体现了 sink 管道的两个阶段:insertInto只是准备(构建TablePipeline,可先行printExplain()检查执行计划),真正触发执行的是pipeline.execute()。
翻译与执行查询
不论输入数据源是流式的还是批式的,Table API 和 SQL 查询都会被转换成 DataStream 程序。查询在内部表示为逻辑查询计划,并被翻译成两个阶段:
- 优化逻辑执行计划
- 翻译成 DataStream 程序
Table API 或 SQL 查询在下列情况下会被翻译:
TableEnvironment.executeSql()被调用时:用于执行一条 SQL 语句,一旦被调用,SQL 语句立即被翻译;TablePipeline.execute()被调用时:用于执行一个源表到输出表的数据流,一旦被调用,Table API 程序立即被翻译;Table.execute()被调用时:用于将一个表的内容收集到本地,一旦被调用,Table API 程序立即被翻译;StatementSet.execute()被调用时:TablePipeline(通过StatementSet.add()输出给某个 Sink)和 INSERT 语句(通过调用StatementSet.addInsertSql())会先被缓存到StatementSet中,当StatementSet.execute()被调用时,所有的 sink 会被优化成一张有向无环图(DAG),从而共享公共子计划、避免重复计算;Table被转换成DataStream时(参阅与 DataStream 集成):转换完成后,它就成为一个普通的 DataStream 程序,并会在调用StreamExecutionEnvironment.execute()时被执行。
StatementSet由TableEnvironment.createStatementSet()创建(接口定义见 TableEnvironment.java),适合在一条作业中批量提交多个 sink 的场景。
查询优化
Apache Flink 使用并扩展了 Apache Calcite 来执行复杂的查询优化,包括一系列基于规则和基于成本的优化,例如:
- 基于 Apache Calcite 的子查询解相关(subquery decorrelation)
- 投影剪裁(projection pruning)
- 分区剪裁(partition pruning)
- 过滤器下推(filter push-down)
- 子计划消除重复数据以避免重复计算
- 特殊子查询重写,包括两部分:
- 将
IN和EXISTS转换为 left semi-joins - 将
NOT IN和NOT EXISTS转换为 left anti-join
- 将
- 可选 join 重新排序:通过
table.optimizer.join-reorder-enabled配置启用
注意:当前仅在子查询重写的结合条件下支持IN/EXISTS/NOT IN/NOT EXISTS。
优化器不仅基于计划,还基于可从数据源获得的丰富统计信息,以及每个算子(例如 io、cpu、网络和内存)的细粒度成本来做出明智的决策。
高级用户可以通过CalciteConfig对象提供自定义优化,通过调用TableEnvironment#getConfig#setPlannerConfig将其提供给 TableEnvironment。
解释表(Explain)
Table API 提供了一种机制来解释计算Table的逻辑和优化查询计划。这是通过Table.explain()方法或者StatementSet.explain()方法完成的:
Table.explain()返回一个Table的计划;StatementSet.explain()返回多 sink 计划的结果。
它们返回一个描述三种计划的字符串:
- 关系查询的抽象语法树(the Abstract Syntax Tree),即未优化的逻辑查询计划;
- 优化的逻辑查询计划;
- 物理执行计划。
此外,可以用TableEnvironment.explainSql()方法和TableEnvironment.executeSql()方法支持执行一个EXPLAIN语句获取逻辑和优化查询计划,请参阅 EXPLAIN 页面。
以下代码展示了给定Table使用Table.explain()的示例:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tEnv = StreamTableEnvironment.create(env); DataStream<Tuple2<Integer, String>> stream1 = env.fromElements(new Tuple2<>(1, "hello")); DataStream<Tuple2<Integer, String>> stream2 = env.fromElements(new Tuple2<>(1, "hello")); // explain Table API Table table1 = tEnv.fromDataStream(stream1, $("count"), $("word")); Table table2 = tEnv.fromDataStream(stream2, $("count"), $("word")); Table table = table1 .where($("word").like("F%")) .unionAll(table2); System.out.println(table.explain());上述例子的输出:
== Abstract Syntax Tree == LogicalUnion(all=[true]) :- LogicalFilter(condition=[LIKE($1, _UTF-16LE'F%')]) : +- LogicalTableScan(table=[[Unregistered_DataStream_1]]) +- LogicalTableScan(table=[[Unregistered_DataStream_2]]) == Optimized Physical Plan == Union(all=[true], union=[count, word]) :- Calc(select=[count, word], where=[LIKE(word, _UTF-16LE'F%')]) : +- DataStreamScan(table=[[Unregistered_DataStream_1]], fields=[count, word]) +- DataStreamScan(table=[[Unregistered_DataStream_2]], fields=[count, word]) == Optimized Execution Plan == Union(all=[true], union=[count, word]) :- Calc(select=[count, word], where=[LIKE(word, _UTF-16LE'F%')]) : +- DataStreamScan(table=[[Unregistered_DataStream_1]], fields=[count, word]) +- DataStreamScan(table=[[Unregistered_DataStream_2]], fields=[count, word])可以看到:AST 阶段展示的是未经优化的逻辑算子(LogicalUnion、LogicalFilter、LogicalTableScan);优化后的物理计划中LIKE过滤条件已经被下推进Calc算子;执行计划则给出了最终可执行的算子拓扑。explain是排查查询语义与优化效果最直接的工具。
多 sink 计划的解释
当使用StatementSet提交多个 sink 时,可以用StatementSet.explain()观察多 sink 计划——所有 sink 被优化成一张有向无环图,公共子计划会被复用。以下 Java 示例定义了两个文件源和两个文件 sink:
EnvironmentSettings settings = EnvironmentSettings.inStreamingMode(); TableEnvironment tEnv = TableEnvironment.create(settings); final Schema schema = Schema.newBuilder() .column("count", DataTypes.INT()) .column("word", DataTypes.STRING()) .build(); tEnv.createTemporaryTable("MySource1", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "/source/path1") .format("csv") .build()); tEnv.createTemporaryTable("MySource2", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "/source/path2") .format("csv") .build()); tEnv.createTemporaryTable("MySink1", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "/sink/path1") .format("csv") .build()); tEnv.createTemporaryTable("MySink2", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "/sink/path2") .format("csv") .build()); StatementSet stmtSet = tEnv.createStatementSet(); Table table1 = tEnv.from("MySource1").where($("word").like("F%")); stmtSet.add(table1.insertInto("MySink1")); Table table2 = table1.unionAll(tEnv.from("MySource2")); stmtSet.add(table2.insertInto("MySink2")); String explanation = stmtSet.explain(); System.out.println(explanation);Python 中对应为create_statement_set()、stmt_set.add_insert("MySink1", table1)与stmt_set.explain()。
多 sink 计划的输出(节选):
== Abstract Syntax Tree == LogicalLegacySink(name=[`default_catalog`.`default_database`.`MySink1`], fields=[count, word]) +- LogicalFilter(condition=[LIKE($1, _UTF-16LE'F%')]) +- LogicalTableScan(table=[[default_catalog, default_database, MySource1, source: [CsvTableSource(read fields: count, word)]]]) LogicalLegacySink(name=[`default_catalog`.`default_database`.`MySink2`], fields=[count, word]) +- LogicalUnion(all=[true]) :- LogicalFilter(condition=[LIKE($1, _UTF-16LE'F%')]) : +- LogicalTableScan(table=[[default_catalog, default_database, MySource1, source: [CsvTableSource(read fields: count, word)]]]) +- LogicalTableScan(table=[[default_catalog, default_database, MySource2, source: [CsvTableSource(read fields: count, word)]]]) == Optimized Physical Plan == LegacySink(name=[`default_catalog`.`default_database`.`MySink1`], fields=[count, word]) +- Calc(select=[count, word], where=[LIKE(word, _UTF-16LE'F%')]) +- LegacyTableSourceScan(table=[[default_catalog, default_database, MySource1, source: [CsvTableSource(read fields: count, word)]]], fields=[count, word]) LegacySink(name=[`default_catalog`.`default_database`.`MySink2`], fields=[count, word]) +- Union(all=[true], union=[count, word]) :- Calc(select=[count, word], where=[LIKE(word, _UTF-16LE'F%')]) : +- LegacyTableSourceScan(table=[[default_catalog, default_database, MySource1, source: [CsvTableSource(read fields: count, word)]]], fields=[count, word]) +- LegacyTableSourceScan(table=[[default_catalog, default_database, MySource2, source: [CsvTableSource(read fields: count, word)]]], fields=[count, word]) == Optimized Execution Plan == Calc(select=[count, word], where=[LIKE(word, _UTF-16LE'F%')])(reuse_id=[1]) +- LegacyTableSourceScan(table=[[default_catalog, default_database, MySource1, source: [CsvTableSource(read fields: count, word)]]], fields=[count, word]) LegacySink(name=[`default_catalog`.`default_database`.`MySink1`], fields=[count, word]) +- Reused(reference_id=[1]) LegacySink(name=[`default_catalog`.`default_database`.`MySink2`], fields=[count, word]) +- Union(all=[true], union=[count, word]) :- Reused(reference_id=[1]) +- LegacyTableSourceScan(table=[[default_catalog, default_database, MySource2, source: [CsvTableSource(read fields: count, word)]]], fields=[count, word])注意执行计划中的reuse_id=[1]与Reused(reference_id=[1]):由于MySink1与MySink2都引用了对MySource1的相同过滤查询,StatementSet优化器将其识别为公共子计划并复用,避免了同一份Calc计算被重复执行——这正是StatementSet相比多次单独executeSql的核心优势之一。
总结
围绕Table这一核心概念,Flink Table API 与 SQL 形成了完整一致的编程模型:通过TableEnvironment统一管理 catalog、模块与函数注册;通过临时/永久表、虚拟视图与 connector 表三种形态描述数据;以 Table API 或 SQL 两种等价方式构建查询;经insertInto/TablePipeline将结果输出到 sink;最终由基于 Calcite 的优化器完成两阶段翻译与优化,并在Table.execute()、TablePipeline.execute()、StatementSet.execute()等触发点真正执行。explain机制则为调试与优化提供了从 AST 到物理执行计划的全程可视化。本文所涉及的源码证据均可在仓库flink-table与flink-connectors模块中找到,读者可结合 docs/content.zh/docs/dev/table/tableApi.md、docs/content.zh/docs/dev/table/sql/overview.md 与 docs/content.zh/docs/dev/table/catalogs.md 继续深入。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
终极百度贴吧个性化体验指南:TiebaTS模块完全解析
终极百度贴吧个性化体验指南:TiebaTS模块完全解析 TiebaTS是一款基于Xposed框架的开源百度贴吧增强模块,专为追求纯净、高效贴吧浏览体验的用户设计
后端运维观测告警可观测性人工智能AI Agent3个核心策略:深度优化MediaPipe GPU性能的完整指南
3个核心策略:深度优化MediaPipe GPU性能的完整指南 MediaPipe作为跨平台的机器学习框架,为实时媒体处理提供了强大的GPU加速能力。在追求极致
人工智能机器学习计算机视觉多模态本地部署Flink SQL JSON完全指南:从解析到嵌套查询实战
Flink SQL JSON完全指南:从解析到嵌套查询实战 你是否还在为JSON数据解析头疼?面对多层嵌套结构无从下手?本文将带你掌握Apache Flink
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考