개념 및 공통 API
개념 및 공통 API (Concepts & Common API)
Table API와 SQL은 통합된 공동 API로 구성됩니다. 이 API의 중심 개념은 쿼리의 입력과 출력으로 사용되는 Table입니다. 이 문서는 Table API와 SQL 쿼리가 있는 프로그램의 공통 구조, Table을 등록하는 방법, Table을 쿼리하는 방법, Table을 방출하는 방법을 보여줍니다.
출처: 문서
본문
Table API 및 SQL 프로그램의 구조
다음 코드 예시는 Table API와 SQL 프로그램의 공통 구조를 보여줍니다.
모든 Flink Scala API는 deprecated이며 향후 Flink 버전에서 제거될 예정입니다. 여전히 Scala로 애플리케이션을 빌드할 수 있지만 DataStream 및/또는 Table API의 Java 버전으로 이동해야 합니다.
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();
Scala
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.
val 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
val table1 = tableEnv.from("SourceTable")
// Create a Table object from a SQL query
val table2 = tableEnv.sqlQuery("SELECT * FROM SourceTable")
// Emit a Table API result Table to a TableSink, same for SQL result
val tableResult = table1.insertInto("SinkTable").execute()
Python
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")
Table API와 SQL 쿼리는 DataStream 프로그램에 쉽게 통합되고 내장될 수 있습니다. DataStream을 Table로 또는 그 반대로 변환하는 방법을 배우려면 DataStream API 통합 페이지를 살펴보세요.
TableEnvironment 생성
TableEnvironment는 Table API와 SQL 통합의 진입점이며 다음을 담당합니다.
- 내부 카탈로그에
Table등록 - 카탈로그 등록
- 플러그형 모듈 로드
- SQL 쿼리 실행
- 사용자 정의(스칼라, 테이블 또는 집계) 함수 등록
DataStream과Table사이 변환(StreamTableEnvironment의 경우)
모든 TableEnvironment 메서드의 완전한 참조는 TableEnvironment API 페이지를 참조하세요.
Table은 항상 특정 TableEnvironment에 바인딩됩니다. 같은 쿼리에서 다른 TableEnvironment의 테이블을 결합하는 것(예: 조인이나 union)은 불가능합니다. TableEnvironment는 정적 TableEnvironment.create() 메서드를 호출하여 생성합니다.
Java
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);
Scala
import org.apache.flink.table.api.{EnvironmentSettings, TableEnvironment}
val settings = EnvironmentSettings
.newInstance()
.inStreamingMode()
//.inBatchMode()
.build()
val tEnv = TableEnvironment.create(settings)
Python
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와 상호 운용할 수 있습니다.
Java
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);
Scala
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.EnvironmentSettings
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
val env = StreamExecutionEnvironment.getExecutionEnvironment
val 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)
카탈로그에서 테이블 생성
TableEnvironment는 식별자로 생성되는 테이블의 카탈로그 맵을 유지합니다. 각 식별자는 3부분으로 구성됩니다: 카탈로그 이름, 데이터베이스 이름, 객체 이름. 카탈로그나 데이터베이스가 지정되지 않으면 현재 기본값이 사용됩니다(Table identifier expanding 섹션의 예시 참조).
테이블은 가상(VIEWS) 또는 일반(TABLES)일 수 있습니다. VIEWS는 기존 Table 객체(보통 Table API 또는 SQL 쿼리의 결과)에서 만들 수 있습니다. TABLES는 파일, 데이터베이스 테이블 또는 메시지 큐 같은 외부 데이터를 설명합니다.
임시 vs 영구 테이블
테이블은 단일 Flink 세션의 수명주기에 묶인 임시(temporary)일 수도 있고, 여러 Flink 세션과 클러스터에 걸쳐 보이는 영구(permanent)일 수도 있습니다.
영구 테이블은 테이블에 대한 메타데이터를 유지하기 위해 카탈로그(예: Hive Metastore)가 필요합니다. 영구 테이블이 생성되면 카탈로그에 연결된 모든 Flink 세션에 보이며, 테이블이 명시적으로 드롭될 때까지 계속 존재합니다.
반면 임시 테이블은 항상 메모리에 저장되며 생성된 Flink 세션의 기간 동안만 존재합니다. 이 테이블은 다른 세션에는 보이지 않습니다. 어떤 카탈로그나 데이터베이스에도 바인딩되지 않지만 그 네임스페이스에 생성될 수 있습니다. 임시 테이블은 해당 데이터베이스가 제거되어도 삭제되지 않습니다.
Shadowing
기존 영구 테이블과 같은 식별자로 임시 테이블을 등록하는 것이 가능합니다. 임시 테이블은 영구 테이블을 가려(shadow) 임시 테이블이 존재하는 동안 영구 테이블에 접근할 수 없게 만듭니다. 그 식별자를 가진 모든 쿼리는 임시 테이블에 대해 실행됩니다.
이는 실험에 유용할 수 있습니다. 데이터의 하위 집합만 있거나 난독화된 임시 테이블에 대해 동일한 쿼리를 먼저 실행할 수 있게 합니다. 쿼리가 올바른지 확인되면 실제 프로덕션 테이블에 대해 실행할 수 있습니다.
테이블 생성
가상 테이블
Table API 객체는 SQL 용어의 VIEW(가상 테이블)에 해당합니다. 논리적 쿼리 계획을 캡슐화합니다. 다음과 같이 카탈로그에 생성할 수 있습니다.
Java
// 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);
Scala
// get a TableEnvironment
val tableEnv = ... // see "Create a TableEnvironment" section
// table is the result of a simple projection query
val projTable: Table = tableEnv.from("X").select(...)
// register the Table projTable as table "projectedTable"
tableEnv.createTemporaryView("projectedTable", projTable)
Python
# get a TableEnvironment
table_env = ... # see "Create a TableEnvironment" section
# table is the result of a simple projection query
proj_table = table_env.from_path("X").select(...)
# register the Table projTable as table "projectedTable"
table_env.register_table("projectedTable", proj_table)
참고: Table 객체는 관계형 데이터베이스 시스템의 VIEW와 유사합니다. 즉 Table을 정의하는 쿼리는 최적화되지 않지만 다른 쿼리가 등록된 Table을 참조할 때 인라인됩니다. 여러 쿼리가 같은 등록된 Table을 참조하면 각 참조 쿼리마다 인라인되고 여러 번 실행됩니다. 즉 등록된 Table의 결과는 공유되지 않습니다.
커넥터 테이블
커넥터 선언에서 관계형 데이터베이스에서 알려진 것과 같은 TABLE을 생성하는 것도 가능합니다. 커넥터는 테이블의 데이터를 저장하는 외부 시스템을 설명합니다. Apache Kafka나 일반 파일 시스템 같은 저장 시스템을 여기 선언할 수 있습니다.
이러한 테이블은 Table API를 직접 사용하거나 SQL DDL로 전환하여 생성할 수 있습니다.
Java
// 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);
tableEnv.createTemporaryTable("SourceTableB", sourceDescriptor);
// Using SQL DDL
tableEnv.executeSql("CREATE [TEMPORARY] TABLE MyTable (...) WITH (...)");
Python
# Using table descriptors
source_descriptor = TableDescriptor.for_connector("datagen") \
.schema(Schema.new_builder()
.column("f0", DataTypes.STRING())
.build()) \
.option("rows-per-second", "100") \
.build()
t_env.create_table("SourceTableA", source_descriptor)
t_env.create_temporary_table("SourceTableB", source_descriptor)
# Using SQL DDL
t_env.execute_sql("CREATE [TEMPORARY] TABLE MyTable (...) WITH (...)")
Table 식별자 확장
테이블은 항상 카탈로그, 데이터베이스, 테이블 이름으로 구성된 3부분 식별자로 등록됩니다.
사용자는 그 안에 하나의 카탈로그와 하나의 데이터베이스를 "현재 카탈로그(current catalog)"와 "현재 데이터베이스(current database)"로 설정할 수 있습니다. 이를 통해 위에서 언급한 3부분 식별자의 처음 두 부분은 선택 사항이 될 수 있습니다 - 제공되지 않으면 현재 카탈로그와 현재 데이터베이스를 참조합니다. 사용자는 table API 또는 SQL로 현재 카탈로그와 현재 데이터베이스를 전환할 수 있습니다.
식별자는 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);
Scala
// get a TableEnvironment
val tEnv: TableEnvironment = ...
tEnv.useCatalog("custom_catalog")
tEnv.useDatabase("custom_database")
val 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)
Python
# get a TableEnvironment
t_env = TableEnvironment.create(...)
t_env.use_catalog("custom_catalog")
t_env.use_database("custom_database")
table = ...
# register the view named 'exampleView' in the catalog named 'custom_catalog'
# in the database named 'custom_database'
t_env.create_temporary_view("other_database.exampleView", table)
# register the view named 'example.View' in the catalog named 'custom_catalog'
# in the database named 'custom_database'
t_env.create_temporary_view("`example.View`", table)
# register the view named 'exampleView' in the catalog named 'other_catalog'
# in the database named 'other_database'
t_env.create_temporary_view("other_catalog.other_database.exampleView", table)
테이블 쿼리
Table API
Table API는 Scala와 Java를 위한 언어 통합 쿼리 API입니다. SQL과 달리 쿼리는 String으로 지정되지 않고 호스트 언어에서 단계별로 구성됩니다.
API는 (스트리밍 또는 배치) 테이블을 나타내고 관계형 연산을 적용하는 메서드를 제공하는 Table 클래스에 기반합니다. 이 메서드는 입력 Table에 관계형 연산을 적용한 결과를 나타내는 새 Table 객체를 반환합니다. 일부 관계형 연산은 table.groupBy(...).select() 같은 다중 메서드 호출로 구성되며, 여기서 groupBy(...)는 table의 그룹화를 지정하고 select(...)는 table의 그룹화에 대한 투영을 지정합니다.
Table API 문서는 스트리밍 및 배치 테이블에서 지원되는 모든 Table API 연산을 설명합니다.
다음 예시는 간단한 Table API 집계 쿼리를 보여줍니다.
Java
// 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 query
Scala
// get a TableEnvironment
val tableEnv = ... // see "Create a TableEnvironment" section
// register Orders table
// scan registered Orders table
val orders = tableEnv.from("Orders")
// compute revenue for all customers from France
val revenue = orders
.filter($"cCountry" === "FRANCE")
.groupBy($"cID", $"cName")
.select($"cID", $"cName", $"revenue".sum AS "revSum")
// emit or convert Table
// execute query
참고: Scala Table API는 달러 기호($)로 시작하는 Scala String interpolation을 사용하여 Table의 속성을 참조합니다. Table API는 Scala implicits를 사용합니다. 다음을 반드시 import하세요.
org.apache.flink.table.api._- 암시적 표현식 변환용- DataStream과 변환하려면
org.apache.flink.api.scala._와org.apache.flink.table.api.bridge.scala._
Python
# get a TableEnvironment
table_env = # see "Create a TableEnvironment" section
# register Orders table
# scan registered Orders table
orders = table_env.from_path("Orders")
# compute revenue for all customers from France
revenue = orders \
.filter(col('cCountry') == 'FRANCE') \
.group_by(col('cID'), col('cName')) \
.select(col('cID'), col('cName'), col('revenue').sum.alias('revSum'))
# emit or convert Table
# execute query
SQL
Flink의 SQL 통합은 SQL 표준을 구현하는 Apache Calcite에 기반합니다. SQL 쿼리는 일반 String으로 지정됩니다.
SQL 문서는 스트리밍 및 배치 테이블에 대한 Flink의 SQL 지원을 설명합니다.
다음 예시는 쿼리를 지정하고 결과를 Table로 반환하는 방법을 보여줍니다.
Java
// 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
Scala
// get a TableEnvironment
val tableEnv = ... // see "Create a TableEnvironment" section
// register Orders table
// compute revenue for all customers from France
val revenue = tableEnv.sqlQuery("""
|SELECT cID, cName, SUM(revenue) AS revSum
|FROM Orders
|WHERE cCountry = 'FRANCE'
|GROUP BY cID, cName
""".stripMargin)
// emit or convert Table
// execute query
Python
# get a TableEnvironment
table_env = ... # see "Create a TableEnvironment" section
# register Orders table
# compute revenue for all customers from France
revenue = table_env.sql_query(
"SELECT cID, cName, SUM(revenue) AS revSum "
"FROM Orders "
"WHERE cCountry = 'FRANCE' "
"GROUP BY cID, cName"
)
# emit or convert Table
# execute query
다음 예시는 등록된 테이블에 결과를 삽입하는 갱신 쿼리를 지정하는 방법을 보여줍니다.
Java
// 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"
);
Scala
// get a TableEnvironment
val 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
""".stripMargin)
Python
# get a TableEnvironment
table_env = ... # see "Create a TableEnvironment" section
# register "Orders" table
# register "RevenueFrance" output table
# compute revenue for all customers from France and emit to "RevenueFrance"
table_env.execute_sql(
"INSERT INTO RevenueFrance "
"SELECT cID, cName, SUM(revenue) AS revSum "
"FROM Orders "
"WHERE cCountry = 'FRANCE' "
"GROUP BY cID, cName"
)
Table API와 SQL 혼합
Table API와 SQL 쿼리는 둘 다 Table 객체를 반환하므로 쉽게 혼합할 수 있습니다.
- Table API 쿼리는 SQL 쿼리가 반환한
Table객체에 정의할 수 있습니다. - SQL 쿼리는 결과 Table 등록 후
TableEnvironment에 등록하고 SQL 쿼리의FROM절에서 참조하여 Table API 쿼리의 결과에 정의할 수 있습니다.
테이블 방출
Table은 TableSink에 쓰는 것으로 방출됩니다. TableSink는 다양한 파일 형식(예: CSV, Apache Parquet, Apache Avro), 저장 시스템(예: JDBC, Apache HBase, Apache Cassandra, Elasticsearch) 또는 메시징 시스템(예: Apache Kafka, RabbitMQ)을 지원하는 범용 인터페이스입니다.
배치 Table은 BatchTableSink에만 쓸 수 있고, 스트리밍 Table은 AppendStreamTableSink, RetractStreamTableSink 또는 UpsertStreamTableSink 중 하나를 요구합니다.
사용 가능한 sink와 사용자 지정 DynamicTableSink 구현 방법에 대한 자세한 내용은 Table Sources & Sinks 문서를 참조하세요.
Table.insertInto(String tableName) 메서드는 소스 테이블을 등록된 sink 테이블로 방출하는 완전한 end-to-end 파이프라인을 정의합니다. 이 메서드는 이름으로 카탈로그에서 테이블 sink를 찾고 Table의 스키마가 sink의 스키마와 동일한지 검증합니다. 파이프라인은 TablePipeline.explain()으로 설명할 수 있고 TablePipeline.execute()를 호출하여 실행할 수 있습니다.
다음 예시는 Table을 방출하는 방법을 보여줍니다.
Java
// 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();
Scala
// get a TableEnvironment
val tableEnv = ... // see "Create a TableEnvironment" section
// create an output Table
val 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
val result: Table = ...
// Prepare the insert into pipeline
val pipeline = result.insertInto("CsvSinkTable")
// Print explain details
pipeline.printExplain()
// emit the result Table to the registered TableSink
pipeline.execute()
Python
# get a TableEnvironment
table_env = ... # see "Create a TableEnvironment" section
# create a TableSink
schema = Schema.new_builder()
.column("a", DataTypes.INT())
.column("b", DataTypes.STRING())
.column("c", DataTypes.BIGINT())
.build()
table_env.create_temporary_table("CsvSinkTable", TableDescriptor.for_connector("filesystem")
.schema(schema)
.option("path", "/path/to/file")
.format(FormatDescriptor.for_format("csv")
.option("field-delimiter", "|")
.build())
.build())
# compute a result Table using Table API operators and/or SQL queries
result = ...
# emit the result Table to the registered TableSink
result.execute_insert("CsvSinkTable")
테이블 출력
디버깅 및 개발 목적으로 Table의 내용을 콘솔에 출력할 수 있습니다. Table.execute().print()를 호출하면 됩니다.
Java
Table table = tableEnv.fromValues(1, 2, 3);
table.execute().print();
Scala
val table = tableEnv.fromValues(1, 2, 3)
table.execute().print()
Python
table = table_env.from_elements([(1, 'Hi'), (2, 'Hello')], ['id', 'data'])
table.execute().print()
이는 테이블의 구체화를 트리거하고 클라이언트의 메모리에 내용을 수집합니다. 큰 테이블의 경우
Table.limit()으로 행 수를 제한하는 것이 좋습니다.
결과를 클라이언트로 수집
TableResult.collect()를 사용해 Table의 결과를 클라이언트로 수집할 수 있습니다. 이는 결과를 프로그래밍 방식으로 처리하는 데 사용할 수 있는 반복자를 반환합니다.
Java
Table table = tableEnv.fromValues(1, 2, 3);
try (CloseableIterator<Row> it = table.execute().collect()) {
while (it.hasNext()) {
Row row = it.next();
// process row
}
}
Scala
val table = tableEnv.fromValues(1, 2, 3)
val it = table.execute().collect()
try {
while (it.hasNext) {
val row = it.next()
// process row
}
} finally {
it.close()
}
Python
table = table_env.from_elements([(1, 'Hi'), (2, 'Hello')], ['id', 'data'])
with table.execute().collect() as results:
for row in results:
print(row)
이는 테이블의 구체화를 트리거하고 클라이언트의 메모리에 내용을 수집합니다. 큰 테이블의 경우
Table.limit()으로 행 수를 제한하는 것이 좋습니다.
결과를 여러 Sink 테이블로 방출
StatementSet을 사용하여 여러 Table을 단일 작업에서 여러 sink 테이블로 방출할 수 있습니다. 이는 여러 개별 작업을 실행하는 것보다 효율적입니다.
Java
// create source table
Table sourceTable = tableEnv.from("SourceTable");
// create sink tables
tableEnv.executeSql("CREATE TABLE SinkTable1 (...) WITH (...)");
tableEnv.executeSql("CREATE TABLE SinkTable2 (...) WITH (...)");
// create a statement set
StatementSet stmtSet = tableEnv.createStatementSet();
// add insert statements
stmtSet.add(sourceTable.insertInto("SinkTable1"));
stmtSet.addInsertSql("INSERT INTO SinkTable2 SELECT * FROM SourceTable");
// execute all statements together
stmtSet.execute();
Scala
// create source table
val sourceTable = tableEnv.from("SourceTable")
// create sink tables
tableEnv.executeSql("CREATE TABLE SinkTable1 (...) WITH (...)")
tableEnv.executeSql("CREATE TABLE SinkTable2 (...) WITH (...)")
// create a statement set
val stmtSet = tableEnv.createStatementSet()
// add insert statements
stmtSet.add(sourceTable.insertInto("SinkTable1"))
stmtSet.addInsertSql("INSERT INTO SinkTable2 SELECT * FROM SourceTable")
// execute all statements together
stmtSet.execute()
Python
# create source table
table = table_env.from_elements([(1, 'Hi'), (2, 'Hello')], ['id', 'data'])
table_env.create_temporary_view("source_table", table)
# create sink tables
table_env.execute_sql("""
CREATE TABLE first_sink_table (
id BIGINT,
data VARCHAR
) WITH (
'connector' = 'print'
)
""")
table_env.execute_sql("""
CREATE TABLE second_sink_table (
id BIGINT,
data VARCHAR
) WITH (
'connector' = 'print'
)
""")
# create a statement set
statement_set = table_env.create_statement_set()
# add insert statements
statement_set.add_insert("first_sink_table", table)
statement_set.add_insert_sql("INSERT INTO second_sink_table SELECT * FROM source_table")
# execute all statements together
statement_set.execute().wait()
쿼리 변환 및 실행
Table API와 SQL 쿼리는 입력이 스트리밍이든 배치든 DataStream 프로그램으로 변환됩니다. 쿼리는 내부적으로 논리적 쿼리 계획으로 표현되며 두 단계로 변환됩니다.
- 논리적 계획의 최적화,
- DataStream 프로그램으로의 변환.
Table API 또는 SQL 쿼리는 다음 경우에 변환됩니다.
TableEnvironment.executeSql()을 호출할 때. 이 메서드는 주어진 문을 실행하는 데 사용되며, 이 메서드가 호출되는 즉시 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으로 변환될 때 변환됩니다(Integration with DataStream 참조). 변환되면 일반 DataStream 프로그램이 되며StreamExecutionEnvironment.execute()가 호출될 때 실행됩니다.
쿼리 최적화
Apache Flink는 정교한 쿼리 최적화를 수행하기 위해 Apache Calcite를 활용하고 확장합니다. 여기에는 다음과 같은 일련의 규칙 및 비용 기반 최적화가 포함됩니다.
- Apache Calcite에 기반한 하위 쿼리 비상관화(Subquery decorrelation)
- 프로젝트 프루닝(Project pruning)
- 파티션 프루닝(Partition pruning)
- 필터 푸시다운(Filter push-down)
- 중복 계산 방지를 위한 하위 계획 중복 제거(Sub-plan deduplication)
- 두 부분을 포함하는 특수 하위 쿼리 재작성:
- IN과 EXISTS를 왼쪽 semi-join으로 변환
- NOT IN과 NOT EXISTS를 왼쪽 anti-join으로 변환
- 선택적 조인 재정렬(Join reordering)
table.optimizer.join-reorder-enabled로 활성화
참고: IN/EXISTS/NOT IN/NOT EXISTS는 현재 하위 쿼리 재작성에서 접속 조건(conjunctive condition)에서만 지원됩니다.
최적화기는 계획뿐 아니라 데이터 소스에서 얻을 수 있는 풍부한 통계와 각 연산자에 대한 io, cpu, 네트워크, 메모리 같은 세밀한 비용을 기반으로 지능적인 결정을 내립니다.
고급 사용자는 TableEnvironment#getConfig#setPlannerConfig를 호출해 테이블 환경에 제공할 수 있는 CalciteConfig 객체를 통해 사용자 지정 최적화를 제공할 수 있습니다.
테이블 설명 (Explaining)
Table API는 Table을 계산하는 논리적 및 최적화된 쿼리 계획을 설명하는 메커니즘을 제공합니다. 이는 Table.explain() 메서드 또는 StatementSet.explain() 메서드를 통해 수행됩니다. Table.explain()은 Table의 계획을 반환합니다. StatementSet.explain()은 여러 sink의 계획을 반환합니다. 세 가지 계획을 설명하는 String을 반환합니다.
- 관계형 쿼리의 Abstract Syntax Tree, 즉 최적화되지 않은 논리적 쿼리 계획,
- 최적화된 논리적 쿼리 계획,
- 물리적 실행 계획.
TableEnvironment.explainSql()과 TableEnvironment.executeSql()은 EXPLAIN 문을 실행하여 계획을 얻는 것을 지원합니다. EXPLAIN 페이지를 참조하세요.
다음 코드는 주어진 Table에 대해 Table.explain() 메서드를 사용한 예시와 해당 출력을 보여줍니다.
Java
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());
Scala
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tEnv = StreamTableEnvironment.create(env)
val table1 = env.fromElements((1, "hello")).toTable(tEnv, $"count", $"word")
val table2 = env.fromElements((1, "hello")).toTable(tEnv, $"count", $"word")
val table = table1
.where($"word".like("F%"))
.unionAll(table2)
println(table.explain())
Python
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
table1 = t_env.from_elements([(1, "hello")], ["count", "word"])
table2 = t_env.from_elements([(1, "hello")], ["count", "word"])
table = table1 \
.where(col('word').like('F%')) \
.union_all(table2)
print(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])
다음 코드는 StatementSet.explain() 메서드를 사용한 다중 sink 계획의 예시와 해당 출력을 보여줍니다.
Java
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);
Scala
val settings = EnvironmentSettings.inStreamingMode()
val tEnv = TableEnvironment.create(settings)
val 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())
val stmtSet = tEnv.createStatementSet()
val table1 = tEnv.from("MySource1").where($"word".like("F%"))
stmtSet.add(table1.insertInto("MySink1"))
val table2 = table1.unionAll(tEnv.from("MySource2"))
stmtSet.add(table2.insertInto("MySink2"))
val explanation = stmtSet.explain()
println(explanation)
Python
settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(environment_settings=settings)
schema = Schema.new_builder()
.column("count", DataTypes.INT())
.column("word", DataTypes.STRING())
.build()
t_env.create_temporary_table("MySource1", TableDescriptor.for_connector("filesystem")
.schema(schema)
.option("path", "/source/path1")
.format("csv")
.build())
t_env.create_temporary_table("MySource2", TableDescriptor.for_connector("filesystem")
.schema(schema)
.option("path", "/source/path2")
.format("csv")
.build())
t_env.create_temporary_table("MySink1", TableDescriptor.for_connector("filesystem")
.schema(schema)
.option("path", "/sink/path1")
.format("csv")
.build())
t_env.create_temporary_table("MySink2", TableDescriptor.for_connector("filesystem")
.schema(schema)
.option("path", "/sink/path2")
.format("csv")
.build())
stmt_set = t_env.create_statement_set()
table1 = t_env.from_path("MySource1").where(col('word').like('F%'))
stmt_set.add_insert("MySink1", table1)
table2 = table1.union_all(t_env.from_path("MySource2"))
stmt_set.add_insert("MySink2", table2)
explanation = stmt_set.explain()
print(explanation)
다중 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])