DataStream API 통합
DataStream API 통합 (DataStream API Integration)
Table API와 DataStream API는 데이터 처리 파이프라인을 정의할 때 똑같이 중요한 API입니다. DataStream API는 시간, 상태, 데이터플로 관리를 저수준 명령형 방식으로 제공하고, Table API는 구조화된 선언형 API를 제공합니다. 두 API는 유계와 무계 스트림 모두를 처리할 수 있으며, Flink는 두 API를 매끄럽게 통합하기 위한 특별한 브리징 기능을 제공합니다. 이 문서는 StreamTableEnvironment를 사용해 DataStream과 Table 사이를 변환하는 방법, 배치 런타임 모드, changelog 스트림 처리, 그리고 데이터 타입 매핑을 설명합니다.
출처: 문서
본문
데이터 처리 파이프라인을 정의할 때 Table API와 DataStream API는 모두 똑같이 중요합니다.
DataStream API는 스트림 처리의 기본 요소(즉 시간(time), 상태(state), 데이터플로 관리)를 비교적 저수준의 명령형 프로그래밍 API로 제공합니다. Table API는 많은 내부 구현을 추상화하고 구조화된 선언형 API를 제공합니다.
두 API 모두 유계(bounded) 스트림과 무계(unbounded) 스트림으로 작업할 수 있습니다.
유계 스트림은 과거 데이터를 처리할 때 관리해야 합니다. 무계 스트림은 먼저 과거 데이터로 초기화될 수 있는 실시간 처리 시나리오에서 발생합니다.
효율적인 실행을 위해 두 API는 유계 스트림을 최적화된 배치 실행 모드에서 처리하는 방법을 제공합니다. 그러나 배치는 스트리밍의 특수한 경우에 불과하므로, 유계 스트림으로 이루어진 파이프라인을 일반 스트리밍 실행 모드에서도 실행할 수 있습니다.
한쪽 API의 파이프라인은 다른 API에 의존하지 않고도 end-to-end로 정의할 수 있습니다. 하지만 다음과 같은 여러 이유로 두 API를 혼합하는 것이 유용할 수 있습니다:
- DataStream API로 메인 파이프라인을 구현하기 전에, 카탈로그에 접근하거나 외부 시스템에 쉽게 연결하기 위해 테이블 생태계를 활용합니다.
- DataStream API로 메인 파이프라인을 구현하기 전에, 상태 없는 데이터 정규화와 정제를 위해 일부 SQL 함수에 접근합니다.
- Table API에 저수준 연산(예: 커스텀 타이머 처리)이 없는 경우 수시로 DataStream API로 전환합니다.
Flink는 DataStream API와의 통합을 최대한 매끄럽게 만들기 위한 특별한 브리징 기능을 제공합니다.
DataStream API와 Table API 사이를 전환하면 약간의 변환 오버헤드가 발생합니다. 예를 들어, 부분적으로 바이너리 데이터를 다루는 테이블 런타임의 내부 데이터 구조(즉
RowData)를 더 사용자 친화적인 데이터 구조(즉Row)로 변환해야 합니다. 보통 이 오버헤드는 무시할 수 있지만, 완전성을 위해 여기에 언급합니다.
DataStream과 Table 간 변환
Flink는 DataStream API와 통합하기 위한 전용 StreamTableEnvironment를 제공합니다. 이 환경(environment)은 일반적인 TableEnvironment를 확장해 추가 메서드를 제공하며, DataStream API에서 사용하는 StreamExecutionEnvironment를 매개변수로 받습니다.
다음 코드는 두 API 사이를 오가며 변환하는 예시를 보여줍니다. Table의 컬럼 이름과 타입은 DataStream의 TypeInformation에서 자동으로 파생됩니다. DataStream API는 기본적으로 changelog 처리를 지원하지 않으므로, 이 코드는 스트림→테이블 및 테이블→스트림 변환 동안 append-only/insert-only 의미론을 가정합니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.types.Row;
// create environments of both APIs
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// create a DataStream
DataStream<String> dataStream = env.fromElements("Alice", "Bob", "John");
// interpret the insert-only DataStream as a Table
Table inputTable = tableEnv.fromDataStream(dataStream);
// register the Table object as a view and query it
tableEnv.createTemporaryView("InputTable", inputTable);
Table resultTable = tableEnv.sqlQuery("SELECT UPPER(f0) FROM InputTable");
// interpret the insert-only Table as a DataStream again
DataStream<Row> resultStream = tableEnv.toDataStream(resultTable);
// add a printing sink and execute in DataStream API
resultStream.print();
env.execute();
// prints:
// +I[ALICE]
// +I[BOB]
// +I[JOHN]
Scala
import org.apache.flink.api.scala._
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
// create environments of both APIs
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env)
// create a DataStream
val dataStream = env.fromElements("Alice", "Bob", "John")
// interpret the insert-only DataStream as a Table
val inputTable = tableEnv.fromDataStream(dataStream)
// register the Table object as a view and query it
tableEnv.createTemporaryView("InputTable", inputTable)
val resultTable = tableEnv.sqlQuery("SELECT UPPER(f0) FROM InputTable")
// interpret the insert-only Table as a DataStream again
val resultStream = tableEnv.toDataStream(resultTable)
// add a printing sink and execute in DataStream API
resultStream.print()
env.execute()
// prints:
// +I[ALICE]
// +I[BOB]
// +I[JOHN]
Python
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
from pyflink.common.typeinfo import Types
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# create a DataStream
ds = env.from_collection(["Alice", "Bob", "John"], Types.STRING())
# interpret the insert-only DataStream as a Table
t = t_env.from_data_stream(ds)
# register the Table object as a view and query it
t_env.create_temporary_view("InputTable", t)
res_table = t_env.sql_query("SELECT UPPER(f0) FROM InputTable")
# interpret the insert-only Table as a DataStream again
res_ds = t_env.to_data_stream(res_table)
# add a printing sink and execute in DataStream API
res_ds.print()
env.execute()
# prints:
# +I[ALICE]
# +I[BOB]
# +I[JOHN]
fromDataStream과 toDataStream의 완전한 의미론은 아래의 전용 섹션에서 확인할 수 있습니다. 특히 이 섹션에서는 더 복잡하고 중첩된 타입으로 스키마 파생에 영향을 주는 방법을 다룹니다. 또한 event-time과 워터마크를 사용하는 방법도 다룹니다.
쿼리 종류에 따라, 대부분의 경우 결과 동적 테이블은 Table을 DataStream으로 변환할 때 insert-only 변경만이 아니라 retraction과 다른 종류의 업데이트도 생성하는 파이프라인입니다. 테이블→스트림 변환 중에는 다음과 유사한 예외가 발생할 수 있습니다:
Table sink 'Unregistered_DataStream_Sink_1' doesn't support consuming update changes [...].
이런 경우 쿼리를 다시 수정하거나 toChangelogStream으로 전환해야 합니다.
다음 예시는 업데이트되는 테이블을 어떻게 변환하는지 보여줍니다. 모든 결과 행은 row.getKind()를 호출해 조회할 수 있는 변경 플래그가 있는 changelog 항목을 나타냅니다. 예시에서 Alice의 두 번째 점수는 update before(-U)와 update after(+U) 변경을 만듭니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.types.Row;
// create environments of both APIs
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// create a DataStream
DataStream<Row> dataStream = env.fromElements(
Row.of("Alice", 12),
Row.of("Bob", 10),
Row.of("Alice", 100));
// interpret the insert-only DataStream as a Table
Table inputTable = tableEnv.fromDataStream(dataStream).as("name", "score");
// register the Table object as a view and query it
// the query contains an aggregation that produces updates
tableEnv.createTemporaryView("InputTable", inputTable);
Table resultTable = tableEnv.sqlQuery(
"SELECT name, SUM(score) FROM InputTable GROUP BY name");
// interpret the updating Table as a changelog DataStream
DataStream<Row> resultStream = tableEnv.toChangelogStream(resultTable);
// add a printing sink and execute in DataStream API
resultStream.print();
env.execute();
// prints:
// +I[Alice, 12]
// +I[Bob, 10]
// -U[Alice, 12]
// +U[Alice, 112]
Scala
import org.apache.flink.api.scala.typeutils.Types
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
import org.apache.flink.types.Row
// create environments of both APIs
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env)
// create a DataStream
val dataStream = env.fromElements(
Row.of("Alice", Int.box(12)),
Row.of("Bob", Int.box(10)),
Row.of("Alice", Int.box(100))
)(Types.ROW(Types.STRING, Types.INT))
// interpret the insert-only DataStream as a Table
val inputTable = tableEnv.fromDataStream(dataStream).as("name", "score")
// register the Table object as a view and query it
// the query contains an aggregation that produces updates
tableEnv.createTemporaryView("InputTable", inputTable)
val resultTable = tableEnv.sqlQuery("SELECT name, SUM(score) FROM InputTable GROUP BY name")
// interpret the updating Table as a changelog DataStream
val resultStream = tableEnv.toChangelogStream(resultTable)
// add a printing sink and execute in DataStream API
resultStream.print()
env.execute()
// prints:
// +I[Alice, 12]
// +I[Bob, 10]
// -U[Alice, 12]
// +U[Alice, 112]
Python
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
from pyflink.common.typeinfo import Types
# create environments of both APIs
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# create a DataStream
ds = env.from_collection([("Alice", 12), ("Bob", 10), ("Alice", 100)],
type_info=Types.ROW_NAMED(
["a", "b"],
[Types.STRING(), Types.INT()]))
input_table = t_env.from_data_stream(ds).alias("name", "score")
# register the Table object as a view and query it
# the query contains an aggregation that produces updates
t_env.create_temporary_view("InputTable", input_table)
res_table = t_env.sql_query("SELECT name, SUM(score) FROM InputTable GROUP BY name")
# interpret the updating Table as a changelog DataStream
res_stream = t_env.to_changelog_stream(res_table)
# add a printing sink and execute in DataStream API
res_stream.print()
env.execute()
# prints:
# +I[Alice, 12]
# +I[Bob, 10]
# -U[Alice, 12]
# +U[Alice, 112]
fromChangelogStream과 toChangelogStream의 완전한 의미론은 아래의 전용 섹션에서 확인할 수 있습니다. 특히 이 섹션은 더 복잡하고 중첩된 타입으로 스키마 파생에 영향을 주는 방법을 다룹니다. event-time과 워터마크로 작업하는 방법을 다루며, 입력 및 출력 스트림에 대한 기본 키(primary key)와 changelog 모드를 선언하는 방법도 논의합니다.
위의 예시는 들어오는 각 레코드마다 행 단위 업데이트를 계속 방출하여 최종 결과를 증분 방식으로 계산하는 방법을 보여줍니다. 그러나 입력 스트림이 유한한(즉 유계인) 경우에는 배치 처리 원리를 활용해 결과를 더 효율적으로 계산할 수 있습니다.
배치 처리에서는 연산자가 결과를 방출하기 전에 전체 입력 테이블을 소비하는 연속 단계로 실행될 수 있습니다. 예를 들어, join 연산자는 실제 조인을 수행하기 전에 두 유계 입력을 정렬하거나(즉 sort-merge join 알고리즘), 다른 입력을 소비하기 전에 한 입력으로 해시 테이블을 구축할 수 있습니다(즉 hash join 알고리즘의 build/probe 단계).
DataStream API와 Table API 모두 전용 *배치 런타임 모드(batch runtime mode)*를 제공합니다.
다음 예시는 단일 플래그를 전환하는 것만으로 통합 파이프라인이 배치 데이터와 스트리밍 데이터를 모두 처리할 수 있음을 보여줍니다.
Java
import org.apache.flink.api.common.RuntimeExecutionMode;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
// setup DataStream API
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// set the batch runtime mode
env.setRuntimeMode(RuntimeExecutionMode.BATCH);
// uncomment this for streaming mode
// env.setRuntimeMode(RuntimeExecutionMode.STREAMING);
// setup Table API
// the table environment adopts the runtime mode during initialization
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// define the same pipeline as above
// prints in BATCH mode:
// +I[Bob, 10]
// +I[Alice, 112]
// prints in STREAMING mode:
// +I[Alice, 12]
// +I[Bob, 10]
// -U[Alice, 12]
// +U[Alice, 112]
Scala
import org.apache.flink.api.common.RuntimeExecutionMode
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
// setup DataStream API
val env = StreamExecutionEnvironment.getExecutionEnvironment()
// set the batch runtime mode
env.setRuntimeMode(RuntimeExecutionMode.BATCH)
// uncomment this for streaming mode
// env.setRuntimeMode(RuntimeExecutionMode.STREAMING)
// setup Table API
// the table environment adopts the runtime mode during initialization
val tableEnv = StreamTableEnvironment.create(env)
// define the same pipeline as above
// prints in BATCH mode:
// +I[Bob, 10]
// +I[Alice, 112]
// prints in STREAMING mode:
// +I[Alice, 12]
// +I[Bob, 10]
// -U[Alice, 12]
// +U[Alice, 112]
Python
from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode
from pyflink.table import StreamTableEnvironment
# setup DataStream API
env = StreamExecutionEnvironment.get_execution_environment()
# set the batch runtime mode
env.set_runtime_mode(RuntimeExecutionMode.BATCH)
# uncomment this for streaming mode
# env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
# setup Table API
# the table environment adopts the runtime mode during initialization
table_env = StreamTableEnvironment.create(env)
# define the same pipeline as above
# prints in BATCH mode:
# +I[Bob, 10]
# +I[Alice, 112]
# prints in STREAMING mode:
# +I[Alice, 12]
# +I[Bob, 10]
# -U[Alice, 12]
# +U[Alice, 112]
changelog를 외부 시스템(예: 키-값 저장소)에 적용하면 두 모드 모두 정확히 동일한 출력 테이블을 생성할 수 있음을 확인할 수 있습니다. 결과를 방출하기 전에 모든 입력 데이터를 소비함으로써, 배치 모드의 changelog는 insert-only 변경으로만 구성됩니다. 더 자세한 내용은 아래의 전용 배치 모드 섹션을 참조하세요.
의존성과 임포트(Imports)
Table API와 DataStream API를 결합하는 프로젝트는 다음 브리징 모듈 중 하나를 추가해야 합니다. 이들은 flink-table-api-java 또는 flink-table-api-scala에 대한 전이 의존성과 해당 언어별 DataStream API 모듈을 포함합니다.
Java
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge_2.12</artifactId>
<version>2.3.0</version>
<scope>provided</scope>
</dependency>
Scala
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-scala-bridge_2.12</artifactId>
<version>2.3.0</version>
<scope>provided</scope>
</dependency>
DataStream API와 Table API의 Java 또는 Scala 버전을 사용해 공통 파이프라인을 선언하려면 다음 임포트가 필요합니다.
Java
// imports for Java DataStream API
import org.apache.flink.streaming.api.*;
import org.apache.flink.streaming.api.environment.*;
// imports for Table API with bridging to Java DataStream API
import org.apache.flink.table.api.*;
import org.apache.flink.table.api.bridge.java.*;
Scala
// imports for Scala DataStream API
import org.apache.flink.api.scala._
import org.apache.flink.streaming.api.scala._
// imports for Table API with bridging to Scala DataStream API
import org.apache.flink.table.api._
import org.apache.flink.table.api.bridge.scala._
Python
# imports for Python DataStream API
from pyflink.datastream import *
# imports for Table API to Python DataStream API
from pyflink.table import *
자세한 내용은 configuration 섹션을 참조하세요.
구성(Configuration)
TableEnvironment는 전달받은 StreamExecutionEnvironment의 모든 구성 옵션을 채택합니다. 그러나 StreamExecutionEnvironment의 구성에 대한 이후 변경사항이 인스턴스화 후 StreamTableEnvironment로 전파된다는 보장은 없습니다. Table API에서 DataStream API로의 옵션 전파는 플래닝(planning) 중에 발생합니다.
Table API로 전환하기 전에 DataStream API에서 모든 구성 옵션을 미리 설정하는 것을 권장합니다.
Java
import java.time.ZoneId;
import org.apache.flink.core.execution.CheckpointingMode;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
// create Java DataStream API
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// set various configuration early
env.setMaxParallelism(256);
env.getConfig().addDefaultKryoSerializer(MyCustomType.class, CustomKryoSerializer.class);
env.getCheckpointConfig().setCheckpointingConsistencyMode(CheckpointingMode.EXACTLY_ONCE);
// then switch to Java Table API
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// set configuration early
tableEnv.getConfig().setLocalTimeZone(ZoneId.of("Europe/Berlin"));
// start defining your pipelines in both APIs...
Scala
import java.time.ZoneId
import org.apache.flink.core.execution.CheckpointingMode
import org.apache.flink.api.scala._
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala._
// create Scala DataStream API
val env = StreamExecutionEnvironment.getExecutionEnvironment
// set various configuration early
env.setMaxParallelism(256)
env.getConfig.addDefaultKryoSerializer(classOf[MyCustomType], classOf[CustomKryoSerializer])
env.getCheckpointConfig.setCheckpointingConsistencyMode(CheckpointingMode.EXACTLY_ONCE)
// then switch to Scala Table API
val tableEnv = StreamTableEnvironment.create(env)
// set configuration early
tableEnv.getConfig.setLocalTimeZone(ZoneId.of("Europe/Berlin"))
// start defining your pipelines in both APIs...
Python
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
from pyflink.datastream.checkpointing_mode import CheckpointingMode
# create Python DataStream API
env = StreamExecutionEnvironment.get_execution_environment()
# set various configuration early
env.set_max_parallelism(256)
env.get_checkpoint_config().set_checkpointing_mode(CheckpointingMode.EXACTLY_ONCE)
# then switch to Python Table API
t_env = StreamTableEnvironment.create(env)
# set configuration early
t_env.get_config().set_local_timezone("Europe/Berlin")
# start defining your pipelines in both APIs...
실행 동작(Execution Behavior)
두 API 모두 파이프라인을 실행하는 메서드를 제공합니다. 즉, 요청하면 클러스터에 제출되어 실행이 트리거될 작업 그래프를 컴파일합니다. 결과는 선언된 싱크(sink)로 스트리밍됩니다.
보통 두 API는 메서드 이름에 execute라는 용어로 이러한 동작을 나타냅니다. 그러나 실행 동작은 Table API와 DataStream API 사이에 약간 다릅니다.
DataStream API
DataStream API의 StreamExecutionEnvironment는 *빌더 패턴(builder pattern)*을 사용해 복잡한 파이프라인을 구성합니다. 파이프라인은 싱크로 끝날 수도 있고 아닐 수도 있는 여러 분기(branch)로 나뉠 수 있습니다. 환경은 작업 제출 시점까지 정의된 모든 분기를 버퍼링합니다.
StreamExecutionEnvironment.execute()는 구성된 전체 파이프라인을 제출하고 이후 빌더를 비웁니다. 즉, 더 이상 소스와 싱크가 선언되지 않은 상태가 되어 새 파이프라인을 빌더에 추가할 수 있습니다. 따라서 모든 DataStream 프로그램은 보통 StreamExecutionEnvironment.execute() 호출로 끝납니다. 또는 DataStream.executeAndCollect()가 결과를 로컬 클라이언트로 스트리밍하기 위한 싱크를 암시적으로 정의합니다.
Table API
Table API에서 분기 파이프라인은 각 분기가 최종 싱크를 선언해야 하는 StatementSet 내에서만 지원됩니다. TableEnvironment와 StreamTableEnvironment 모두 전용 일반 execute() 메서드를 제공하지 않습니다. 대신 단일 소스-투-싱크 파이프라인 또는 명령문 집합(statement set)을 제출하는 메서드를 제공합니다:
Java
// execute with explicit sink
tableEnv.from("InputTable").insertInto("OutputTable").execute();
tableEnv.executeSql("INSERT INTO OutputTable SELECT * FROM InputTable");
tableEnv.createStatementSet()
.add(tableEnv.from("InputTable").insertInto("OutputTable"))
.add(tableEnv.from("InputTable").insertInto("OutputTable2"))
.execute();
tableEnv.createStatementSet()
.addInsertSql("INSERT INTO OutputTable SELECT * FROM InputTable")
.addInsertSql("INSERT INTO OutputTable2 SELECT * FROM InputTable")
.execute();
// execute with implicit local sink
tableEnv.from("InputTable").execute().print();
tableEnv.executeSql("SELECT * FROM InputTable").print();
Python
# execute with explicit sink
table_env.from_path("input_table").execute_insert("output_table")
table_env.execute_sql("INSERT INTO output_table SELECT * FROM input_table")
table_env.create_statement_set() \
.add_insert("output_table", input_table) \
.add_insert("output_table2", input_table) \
.execute()
table_env.create_statement_set() \
.add_insert_sql("INSERT INTO output_table SELECT * FROM input_table") \
.add_insert_sql("INSERT INTO output_table2 SELECT * FROM input_table") \
.execute()
# execute with implicit local sink
table_env.from_path("input_table").execute().print()
table_env.execute_sql("SELECT * FROM input_table").print()
두 실행 동작을 결합하려면 StreamTableEnvironment.toDataStream 또는 StreamTableEnvironment.toChangelogStream 호출마다 Table API 하위 파이프라인이 구체화(즉 컴파일)되어 DataStream API 파이프라인 빌더에 삽입됩니다. 즉, 이후에 StreamExecutionEnvironment.execute() 또는 DataStream.executeAndCollect를 호출해야 합니다. Table API에서의 실행은 이러한 "외부 부분"을 트리거하지 않습니다.
Java
// (1)
// adds a branch with a printing sink to the StreamExecutionEnvironment
tableEnv.toDataStream(table).print();
// (2)
// executes a Table API end-to-end pipeline as a Flink job and prints locally,
// thus (1) has still not been executed
table.execute().print();
// executes the DataStream API pipeline with the sink defined in (1) as a
// Flink job, (2) was already running before
env.execute();
Python
# (1)
# adds a branch with a printing sink to the StreamExecutionEnvironment
table_env.to_data_stream(table).print()
# (2)
# executes a Table API end-to-end pipeline as a Flink job and prints locally,
# thus (1) has still not been executed
table.execute().print()
# executes the DataStream API pipeline with the sink defined in (1) as a
# Flink job, (2) was already running before
env.execute()
배치 런타임 모드(Batch Runtime Mode)
배치 런타임 모드는 유계(bounded) Flink 프로그램을 위한 특수 실행 모드입니다.
일반적으로 *유계성(boundedness)*은 해당 소스에서 오는 모든 레코드가 실행 전에 알려져 있는지, 아니면 잠재적으로 무기한 새로운 데이터가 나타나는지 알려주는 데이터 소스의 속성입니다. 반대로 작업(job)은 모든 소스가 유계이면 유계이고, 그렇지 않으면 무계입니다.
반면 스트리밍 런타임 모드는 유계와 무계 작업 모두에 사용할 수 있습니다.
서로 다른 실행 모드에 대한 자세한 내용은 해당 DataStream API 섹션을 참조하세요.
Table API & SQL 플래너는 두 모드 각각에 대해 전문화된 최적화 규칙과 런타임 연산자 집합을 제공합니다.
현재 런타임 모드는 소스로부터 자동으로 파생되지 않으므로, 명시적으로 설정해야 하며 그렇지 않으면 StreamTableEnvironment를 인스턴스화할 때 StreamExecutionEnvironment에서 채택됩니다:
Java
import org.apache.flink.api.common.RuntimeExecutionMode;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
// adopt mode from StreamExecutionEnvironment
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setRuntimeMode(RuntimeExecutionMode.BATCH);
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// or
// set mode explicitly for StreamTableEnvironment
// it will be propagated to StreamExecutionEnvironment during planning
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, EnvironmentSettings.inBatchMode());
Scala
import org.apache.flink.api.common.RuntimeExecutionMode
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
import org.apache.flink.table.api.EnvironmentSettings
// adopt mode from StreamExecutionEnvironment
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setRuntimeMode(RuntimeExecutionMode.BATCH)
val tableEnv = StreamTableEnvironment.create(env)
// or
// set mode explicitly for StreamTableEnvironment
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env, EnvironmentSettings.inBatchMode)
Python
from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode
from pyflink.table import EnvironmentSettings, StreamTableEnvironment
# adopt mode from StreamExecutionEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
env.set_runtime_mode(RuntimeExecutionMode.BATCH)
table_env = StreamTableEnvironment.create(env)
# or
# set mode explicitly for StreamTableEnvironment
# it will be propagated to StreamExecutionEnvironment during planning
env = StreamExecutionEnvironment.get_execution_environment()
table_env = StreamTableEnvironment.create(env, EnvironmentSettings.in_batch_mode())
런타임 모드를 BATCH로 설정하기 전에 다음 전제 조건을 충족해야 합니다:
- 모든 소스가 스스로 유계라고 선언해야 합니다.
- 현재 테이블 소스는 insert-only 변경만 방출해야 합니다.
- 연산자는 정렬 및 기타 중간 결과를 위해 충분한 off-heap 메모리가 필요합니다.
- 모든 테이블 연산이 배치 모드에서 사용 가능해야 합니다. 현재 일부 연산은 스트리밍 모드에서만 사용 가능합니다. 해당 Table API & SQL 페이지를 확인하세요.
배치 실행은 (다른 것들 중에서도) 다음과 같은 영향을 미칩니다:
- 점진적 워터마크(progressive watermark)는 생성되지도, 연산자에서 사용되지도 않습니다. 그러나 소스는 종료 전에 최대 워터마크를 방출합니다.
- 태스크 간 교환은
execution.batch-shuffle-mode에 따라 블로킹될 수 있습니다. 이는 동일한 파이프라인을 스트리밍 모드로 실행하는 것과 비교해 리소스 요구량이 더 적어질 수 있음을 의미합니다. - 체크포인트는 비활성화됩니다. 인위적인 상태 백엔드가 삽입됩니다.
- 테이블 연산은 증분 업데이트를 생성하지 않고, insert-only changelog 스트림으로 변환되는 완전한 최종 결과만 생성합니다.
배치 처리는 스트림 처리의 특수한 경우로 간주될 수 있으므로, 유계와 무계 데이터 모두에 가장 일반적인 구현인 스트리밍 파이프라인을 먼저 구현하는 것을 권장합니다.
이론적으로 스트리밍 파이프라인은 모든 연산자를 실행할 수 있습니다. 그러나 실제로는 일부 연산이 계속 커지는 상태를 초래해 의미가 없으므로 지원되지 않을 수 있습니다. 전역 정렬(global sort)은 배치 모드에서만 사용할 수 있는 예입니다. 간단히 말해: 동작하는 스트리밍 파이프라인을 배치 모드에서 실행하는 것은 가능해야 하지만, 그 반대가 반드시 가능한 것은 아닙니다.
다음 예시는 DataGen 테이블 소스를 사용해 배치 모드를 실험하는 방법을 보여줍니다. 많은 소스는 예를 들어 종료 오프셋이나 타임스탬프를 정의하는 방식으로 커넥터를 암시적으로 유계로 만드는 옵션을 제공합니다. 우리 예시에서는 number-of-rows 옵션으로 행 수를 제한합니다.
Java
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.TableDescriptor;
Table table =
tableEnv.from(
TableDescriptor.forConnector("datagen")
.option("number-of-rows", "10") // make the source bounded
.schema(
Schema.newBuilder()
.column("uid", DataTypes.TINYINT())
.column("payload", DataTypes.STRING())
.build())
.build());
// convert the Table to a DataStream and further transform the pipeline
tableEnv.toDataStream(table)
.keyBy(r -> r.<Byte>getFieldAs("uid"))
.map(r -> "My custom operator: " + r.<String>getFieldAs("payload"))
.executeAndCollect()
.forEachRemaining(System.out::println);
// prints:
// My custom operator: 9660912d30a43c7b035e15bd...
// My custom operator: 29f5f706d2144f4a4f9f52a0...
// ...
Scala
import org.apache.flink.api.scala._
import org.apache.flink.table.api._
val table =
tableEnv.from(
TableDescriptor.forConnector("datagen")
.option("number-of-rows", "10") // make the source bounded
.schema(
Schema.newBuilder()
.column("uid", DataTypes.TINYINT())
.column("payload", DataTypes.STRING())
.build())
.build())
// convert the Table to a DataStream and further transform the pipeline
tableEnv.toDataStream(table)
.keyBy(r => r.getFieldAs[Byte]("uid"))
.map(r => "My custom operator: " + r.getFieldAs[String]("payload"))
.executeAndCollect()
.foreach(println)
// prints:
// My custom operator: 9660912d30a43c7b035e15bd...
// My custom operator: 29f5f706d2144f4a4f9f52a0...
// ...
Python
from pyflink.table import TableDescriptor, Schema, DataTypes
table = table_env.from_descriptor(
TableDescriptor.for_connector("datagen")
.option("number-of-rows", "10")
.schema(
Schema.new_builder()
.column("uid", DataTypes.TINYINT())
.column("payload", DataTypes.STRING())
.build())
.build())
# convert the Table to a DataStream and further transform the pipeline
collect = table_env.to_data_stream(table) \
.key_by(lambda r: r[0]) \
.map(lambda r: "My custom operator: " + r[1]) \
.execute_and_collect()
for c in collect:
print(c)
# prints:
# My custom operator: 9660912d30a43c7b035e15bd...
# My custom operator: 29f5f706d2144f4a4f9f52a0...
# ...
Changelog 통합(Changelog Unification)
대부분의 경우 스트리밍에서 배치 모드로 전환하거나 그 반대로 전환할 때 파이프라인 정의 자체는 Table API와 DataStream API 모두에서 동일하게 유지될 수 있습니다. 그러나 앞서 언급했듯이, 배치 모드에서 증분 연산을 피하기 때문에 결과 changelog 스트림은 다를 수 있습니다.
event-time에 의존하고 완전성 마커로 워터마크를 활용하는 시간 기반 연산은 런타임 모드와 무관한 insert-only changelog 스트림을 생성할 수 있습니다.
다음 Java 예시는 API 수준뿐만 아니라 결과 changelog 스트림에서도 통합된 Flink 프로그램을 보여줍니다. 이 예시는 두 테이블의 시간 속성(ts)을 기반으로 하는 interval join을 사용해 SQL에서 두 테이블(UserTable과 OrderTable)을 조인합니다. 그리고 KeyedProcessFunction과 value state를 사용해 사용자 이름을 중복 제거하는 커스텀 연산자를 DataStream API로 구현합니다.
Java
import org.apache.flink.api.common.functions.OpenContext;
import org.apache.flink.api.common.RuntimeExecutionMode;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.types.Row;
import org.apache.flink.util.Collector;
import java.time.LocalDateTime;
// setup DataStream API
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// use BATCH or STREAMING mode
env.setRuntimeMode(RuntimeExecutionMode.BATCH);
// setup Table API
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// create a user stream
DataStream<Row> userStream = env
.fromElements(
Row.of(LocalDateTime.parse("2021-08-21T13:00:00"), 1, "Alice"),
Row.of(LocalDateTime.parse("2021-08-21T13:05:00"), 2, "Bob"),
Row.of(LocalDateTime.parse("2021-08-21T13:10:00"), 2, "Bob"))
.returns(
Types.ROW_NAMED(
new String[] {"ts", "uid", "name"},
Types.LOCAL_DATE_TIME, Types.INT, Types.STRING));
// create an order stream
DataStream<Row> orderStream = env
.fromElements(
Row.of(LocalDateTime.parse("2021-08-21T13:02:00"), 1, 122),
Row.of(LocalDateTime.parse("2021-08-21T13:07:00"), 2, 239),
Row.of(LocalDateTime.parse("2021-08-21T13:11:00"), 2, 999))
.returns(
Types.ROW_NAMED(
new String[] {"ts", "uid", "amount"},
Types.LOCAL_DATE_TIME, Types.INT, Types.INT));
// create corresponding tables
tableEnv.createTemporaryView(
"UserTable",
userStream,
Schema.newBuilder()
.column("ts", DataTypes.TIMESTAMP(3))
.column("uid", DataTypes.INT())
.column("name", DataTypes.STRING())
.watermark("ts", "ts - INTERVAL '1' SECOND")
.build());
tableEnv.createTemporaryView(
"OrderTable",
orderStream,
Schema.newBuilder()
.column("ts", DataTypes.TIMESTAMP(3))
.column("uid", DataTypes.INT())
.column("amount", DataTypes.INT())
.watermark("ts", "ts - INTERVAL '1' SECOND")
.build());
// perform interval join
Table joinedTable =
tableEnv.sqlQuery(
"SELECT U.name, O.amount " +
"FROM UserTable U, OrderTable O " +
"WHERE U.uid = O.uid AND O.ts BETWEEN U.ts AND U.ts + INTERVAL '5' MINUTES");
DataStream<Row> joinedStream = tableEnv.toDataStream(joinedTable);
joinedStream.print();
// implement a custom operator using ProcessFunction and value state
joinedStream
.keyBy(r -> r.<String>getFieldAs("name"))
.process(
new KeyedProcessFunction<String, Row, String>() {
ValueState<String> seen;
@Override
public void open(OpenContext openContext) {
seen = getRuntimeContext().getState(
new ValueStateDescriptor<>("seen", String.class));
}
@Override
public void processElement(Row row, Context ctx, Collector<String> out)
throws Exception {
String name = row.getFieldAs("name");
if (seen.value() == null) {
seen.update(name);
out.collect(name);
}
}
})
.print();
// execute unified pipeline
env.execute();
// prints (in both BATCH and STREAMING mode):
// +I[Bob, 239]
// +I[Alice, 122]
// +I[Bob, 999]
//
// Bob
// Alice
Python
from datetime import datetime
from pyflink.common import Row, Types
from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode,
KeyedProcessFunction, RuntimeContext
from pyflink.datastream.state import ValueStateDescriptor
from pyflink.table import StreamTableEnvironment, Schema, DataTypes
# setup DataStream API
env = StreamExecutionEnvironment.get_execution_environment()
# use BATCH or STREAMING mode
env.set_runtime_mode(RuntimeExecutionMode.BATCH)
# setup Table API
table_env = StreamTableEnvironment.create(env)
# create a user stream
t_format = "%Y-%m-%dT%H:%M:%S"
user_stream = env.from_collection(
[Row(datetime.strptime("2021-08-21T13:00:00", t_format), 1, "Alice"),
Row(datetime.strptime("2021-08-21T13:05:00", t_format), 2, "Bob"),
Row(datetime.strptime("2021-08-21T13:10:00", t_format), 2, "Bob")],
type_info=Types.ROW_NAMED(["ts1", "uid", "name"],
[Types.SQL_TIMESTAMP(), Types.INT(), Types.STRING()]))
# create an order stream
order_stream = env.from_collection(
[Row(datetime.strptime("2021-08-21T13:02:00", t_format), 1, 122),
Row(datetime.strptime("2021-08-21T13:07:00", t_format), 2, 239),
Row(datetime.strptime("2021-08-21T13:11:00", t_format), 2, 999)],
type_info=Types.ROW_NAMED(["ts1", "uid", "amount"],
[Types.SQL_TIMESTAMP(), Types.INT(), Types.INT()]))
# # create corresponding tables
table_env.create_temporary_view(
"user_table",
user_stream,
Schema.new_builder()
.column_by_expression("ts", "CAST(ts1 AS TIMESTAMP(3))")
.column("uid", DataTypes.INT())
.column("name", DataTypes.STRING())
.watermark("ts", "ts - INTERVAL '1' SECOND")
.build())
table_env.create_temporary_view(
"order_table",
order_stream,
Schema.new_builder()
.column_by_expression("ts", "CAST(ts1 AS TIMESTAMP(3))")
.column("uid", DataTypes.INT())
.column("amount", DataTypes.INT())
.watermark("ts", "ts - INTERVAL '1' SECOND")
.build())
# perform interval join
joined_table = table_env.sql_query(
"SELECT U.name, O.amount " +
"FROM user_table U, order_table O " +
"WHERE U.uid = O.uid AND O.ts BETWEEN U.ts AND U.ts + INTERVAL '5' MINUTES")
joined_stream = table_env.to_data_stream(joined_table)
joined_stream.print()
# implement a custom operator using ProcessFunction and value state
class MyProcessFunction(KeyedProcessFunction):
def __init__(self):
self.seen = None
def open(self, runtime_context: RuntimeContext):
state_descriptor = ValueStateDescriptor("seen", Types.STRING())
self.seen = runtime_context.get_state(state_descriptor)
def process_element(self, value, ctx):
name = value[0]
if self.seen.value() is None:
self.seen.update(name)
yield name
joined_stream \
.key_by(lambda r: r[0]) \
.process(MyProcessFunction()) \
.print()
# execute unified pipeline
env.execute()
# prints (in both BATCH and STREAMING mode):
# +I[Bob, 239]
# +I[Alice, 122]
# +I[Bob, 999]
#
# Bob
# Alice
(Insert-Only) 스트림 처리
StreamTableEnvironment는 DataStream API로/로부터 변환하기 위해 다음 메서드를 제공합니다:
fromDataStream(DataStream): insert-only 변경과 임의 타입의 스트림을 테이블로 해석합니다. event-time과 워터마크는 기본적으로 전파되지 않습니다.fromDataStream(DataStream, Schema): insert-only 변경과 임의 타입의 스트림을 테이블로 해석합니다. 선택적 스키마를 사용하면 컬럼 데이터 타입을 보강하고 시간 속성, 워터마크 전략, 기타 계산 컬럼, 또는 기본 키를 추가할 수 있습니다.createTemporaryView(String, DataStream): SQL에서 접근할 수 있도록 스트림을 이름으로 등록합니다.createTemporaryView(String, fromDataStream(DataStream))의 단축 형태입니다.createTemporaryView(String, DataStream, Schema): SQL에서 접근할 수 있도록 스트림을 이름으로 등록합니다.createTemporaryView(String, fromDataStream(DataStream, Schema))의 단축 형태입니다.toDataStream(Table): 테이블을 insert-only 변경 스트림으로 변환합니다. 기본 스트림 레코드 타입은org.apache.flink.types.Row입니다. 단일 rowtime 속성 컬럼이 DataStream API의 레코드로 다시 기록됩니다. 워터마크도 전파됩니다.toDataStream(Table, AbstractDataType): 테이블을 insert-only 변경 스트림으로 변환합니다. 이 메서드는 원하는 스트림 레코드 타입을 표현하기 위해 데이터 타입을 받습니다. 플래너는 (아마도 중첩된) 데이터 타입의 필드에 컬럼을 매핑하기 위해 암시적 캐스트를 삽입하고 컬럼을 재정렬할 수 있습니다.toDataStream(Table, Class): 반사적으로 원하는 데이터 타입을 빠르게 만들기 위한toDataStream(Table, DataTypes.of(Class))의 단축 형태입니다.
Table API 관점에서 DataStream API로/로부터 변환하는 것은 SQL에서 CREATE TABLE DDL로 정의된 가상 테이블 커넥터에서 읽거나 쓰는 것과 유사합니다.
가상 CREATE TABLE name (schema) WITH (options) 문의 스키마 부분은 DataStream의 타입 정보에서 자동으로 파생되거나, 보강되거나, org.apache.flink.table.api.Schema를 사용해 완전히 수동으로 정의될 수 있습니다.
가상 DataStream 테이블 커넥터는 모든 행에 대해 다음 메타데이터를 노출합니다:
| Key | 데이터 타입 | 설명 | R/W |
|---|---|---|---|
rowtime |
TIMESTAMP_LTZ(3) NOT NULL |
스트림 레코드의 타임스탬프. | R/W |
가상 DataStream 테이블 소스는 SupportsSourceWatermark를 구현하므로, DataStream API에서 워터마크를 채택하기 위한 워터마크 전략으로 SOURCE_WATERMARK() 내장 함수를 호출할 수 있습니다.
fromDataStream 예시
다음 코드는 다양한 시나리오에서 fromDataStream을 사용하는 방법을 보여줍니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
import java.time.Instant;
// some example POJO
public static class User {
public String name;
public Integer score;
public Instant event_time;
// default constructor for DataStream API
public User() {}
// fully assigning constructor for Table API
public User(String name, Integer score, Instant event_time) {
this.name = name;
this.score = score;
this.event_time = event_time;
}
}
// create a DataStream
DataStream<User> dataStream =
env.fromElements(
new User("Alice", 4, Instant.ofEpochMilli(1000)),
new User("Bob", 6, Instant.ofEpochMilli(1001)),
new User("Alice", 10, Instant.ofEpochMilli(1002)));
// === EXAMPLE 1 ===
// derive all physical columns automatically
Table table = tableEnv.fromDataStream(dataStream);
table.printSchema();
// prints:
// (
// `name` STRING,
// `score` INT,
// `event_time` TIMESTAMP_LTZ(9)
// )
// === EXAMPLE 2 ===
// derive all physical columns automatically
// but add computed columns (in this case for creating a proctime attribute column)
Table table = tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.columnByExpression("proc_time", "PROCTIME()")
.build());
table.printSchema();
// prints:
// (
// `name` STRING,
// `score` INT NOT NULL,
// `event_time` TIMESTAMP_LTZ(9),
// `proc_time` TIMESTAMP_LTZ(3) NOT NULL *PROCTIME* AS PROCTIME()
//)
// === EXAMPLE 3 ===
// derive all physical columns automatically
// but add computed columns (in this case for creating a rowtime attribute column)
// and a custom watermark strategy
Table table =
tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.columnByExpression("rowtime", "CAST(event_time AS TIMESTAMP_LTZ(3))")
.watermark("rowtime", "rowtime - INTERVAL '10' SECOND")
.build());
table.printSchema();
// prints:
// (
// `name` STRING,
// `score` INT,
// `event_time` TIMESTAMP_LTZ(9),
// `rowtime` TIMESTAMP_LTZ(3) *ROWTIME* AS CAST(event_time AS TIMESTAMP_LTZ(3)),
// WATERMARK FOR `rowtime`: TIMESTAMP_LTZ(3) AS rowtime - INTERVAL '10' SECOND
// )
// === EXAMPLE 4 ===
// derive all physical columns automatically
// but access the stream record's timestamp for creating a rowtime attribute column
// also rely on the watermarks generated in the DataStream API
// we assume that a watermark strategy has been defined for `dataStream` before
// (not part of this example)
Table table =
tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)")
.watermark("rowtime", "SOURCE_WATERMARK()")
.build());
table.printSchema();
// prints:
// (
// `name` STRING,
// `score` INT,
// `event_time` TIMESTAMP_LTZ(9),
// `rowtime` TIMESTAMP_LTZ(3) *ROWTIME* METADATA,
// WATERMARK FOR `rowtime`: TIMESTAMP_LTZ(3) AS SOURCE_WATERMARK()
// )
// === EXAMPLE 5 ===
// define physical columns manually
// in this example,
// - we can reduce the default precision of timestamps from 9 to 3
// - we also project the columns and put `event_time` to the beginning
Table table =
tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.column("event_time", "TIMESTAMP_LTZ(3)")
.column("name", "STRING")
.column("score", "INT")
.watermark("event_time", "SOURCE_WATERMARK()")
.build());
table.printSchema();
// prints:
// (
// `event_time` TIMESTAMP_LTZ(3) *ROWTIME*,
// `name` VARCHAR(200),
// `score` INT
// )
// note: the watermark strategy is not shown due to the inserted column reordering projection
Scala
import org.apache.flink.api.scala._
import java.time.Instant
// some example case class
case class User(name: String, score: java.lang.Integer, event_time: java.time.Instant)
// create a DataStream
val dataStream = env.fromElements(
User("Alice", 4, Instant.ofEpochMilli(1000)),
User("Bob", 6, Instant.ofEpochMilli(1001)),
User("Alice", 10, Instant.ofEpochMilli(1002)))
// === EXAMPLE 1 ===
// derive all physical columns automatically
val table = tableEnv.fromDataStream(dataStream)
table.printSchema()
// prints:
// (
// `name` STRING,
// `score` INT,
// `event_time` TIMESTAMP_LTZ(9)
// )
// === EXAMPLE 2 ===
// derive all physical columns automatically
// but add computed columns (in this case for creating a proctime attribute column)
val table = tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.columnByExpression("proc_time", "PROCTIME()")
.build())
table.printSchema()
// prints:
// (
// `name` STRING,
// `score` INT NOT NULL,
// `event_time` TIMESTAMP_LTZ(9),
// `proc_time` TIMESTAMP_LTZ(3) NOT NULL *PROCTIME* AS PROCTIME()
//)
// === EXAMPLE 3 ===
// derive all physical columns automatically
// but add computed columns (in this case for creating a rowtime attribute column)
// and a custom watermark strategy
val table =
tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.columnByExpression("rowtime", "CAST(event_time AS TIMESTAMP_LTZ(3))")
.watermark("rowtime", "rowtime - INTERVAL '10' SECOND")
.build())
table.printSchema()
// prints:
// (
// `name` STRING,
// `score` INT,
// `event_time` TIMESTAMP_LTZ(9),
// `rowtime` TIMESTAMP_LTZ(3) *ROWTIME* AS CAST(event_time AS TIMESTAMP_LTZ(3)),
// WATERMARK FOR `rowtime`: TIMESTAMP_LTZ(3) AS rowtime - INTERVAL '10' SECOND
// )
// === EXAMPLE 4 ===
// derive all physical columns automatically
// but access the stream record's timestamp for creating a rowtime attribute column
// also rely on the watermarks generated in the DataStream API
// we assume that a watermark strategy has been defined for `dataStream` before
// (not part of this example)
val table =
tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)")
.watermark("rowtime", "SOURCE_WATERMARK()")
.build())
table.printSchema()
// prints:
// (
// `name` STRING,
// `score` INT,
// `event_time` TIMESTAMP_LTZ(9),
// `rowtime` TIMESTAMP_LTZ(3) *ROWTIME* METADATA,
// WATERMARK FOR `rowtime`: TIMESTAMP_LTZ(3) AS SOURCE_WATERMARK()
// )
// === EXAMPLE 5 ===
// define physical columns manually
// in this example,
// - we can reduce the default precision of timestamps from 9 to 3
// - we also project the columns and put `event_time` to the beginning
val table =
tableEnv.fromDataStream(
dataStream,
Schema.newBuilder()
.column("event_time", "TIMESTAMP_LTZ(3)")
.column("name", "STRING")
.column("score", "INT")
.watermark("event_time", "SOURCE_WATERMARK()")
.build())
table.printSchema()
// prints:
// (
// `event_time` TIMESTAMP_LTZ(3) *ROWTIME*,
// `name` VARCHAR(200),
// `score` INT
// )
// note: the watermark strategy is not shown due to the inserted column reordering projection
Python
from pyflink.common.time import Instant
from pyflink.common.types import Row
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, Schema
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
ds = env.from_collection([
Row("Alice", 12, Instant.of_epoch_milli(1000)),
Row("Bob", 5, Instant.of_epoch_milli(1001)),
Row("Alice", 10, Instant.of_epoch_milli(1002))],
type_info=Types.ROW_NAMED(['name', 'score', 'event_time'], [Types.STRING(), Types.INT(), Types.INSTANT()]))
# === EXAMPLE 1 ===
# derive all physical columns automatically
table = t_env.from_data_stream(ds)
table.print_schema()
# prints:
# (
# `name` STRING,
# `score` INT,
# `event_time` TIMESTAMP_LTZ(9)
# )
# === EXAMPLE 2 ===
# derive all physical columns automatically
# but add computed columns (in this case for creating a proctime attribute column)
table = t_env.from_data_stream(
ds,
Schema.new_builder()
.column_by_expression("proc_time", "PROCTIME()")
.build())
table.print_schema()
# prints:
# (
# `name` STRING,
# `score` INT,
# `event_time` TIMESTAMP_LTZ(9),
# `proc_time` TIMESTAMP_LTZ(3) NOT NULL *PROCTIME* AS PROCTIME()
# )
# === EXAMPLE 3 ===
# derive all physical columns automatically
# but add computed columns (in this case for creating a rowtime attribute column)
# and a custom watermark strategy
table = t_env.from_data_stream(
ds,
Schema.new_builder()
.column_by_expression("rowtime", "CAST(event_time AS TIMESTAMP_LTZ(3))")
.watermark("rowtime", "rowtime - INTERVAL '10' SECOND")
.build())
table.print_schema()
# prints:
# (
# `name` STRING,
# `score` INT,
# `event_time` TIMESTAMP_LTZ(9),
# `rowtime` TIMESTAMP_LTZ(3) *ROWTIME* AS CAST(event_time AS TIMESTAMP_LTZ(3)),
# WATERMARK FOR `rowtime`: TIMESTAMP_LTZ(3) AS rowtime - INTERVAL '10' SECOND
# )
# === EXAMPLE 4 ===
# derive all physical columns automatically
# but access the stream record's timestamp for creating a rowtime attribute column
# also rely on the watermarks generated in the DataStream API
# we assume that a watermark strategy has been defined for `dataStream` before
# (not part of this example)
table = t_env.from_data_stream(
ds,
Schema.new_builder()
.column_by_metadata("rowtime", "TIMESTAMP_LTZ(3)")
.watermark("rowtime", "SOURCE_WATERMARK()")
.build())
table.print_schema()
# prints:
# (
# `name` STRING,
# `score` INT,
# `event_time` TIMESTAMP_LTZ(9),
# `rowtime` TIMESTAMP_LTZ(3) *ROWTIME* METADATA,
# WATERMARK FOR `rowtime`: TIMESTAMP_LTZ(3) AS SOURCE_WATERMARK()
# )
# === EXAMPLE 5 ===
# define physical columns manually
# in this example,
# - we can reduce the default precision of timestamps from 9 to 3
# - we also project the columns and put `event_time` to the beginning
table = t_env.from_data_stream(
ds,
Schema.new_builder()
.column("event_time", "TIMESTAMP_LTZ(3)")
.column("name", "STRING")
.column("score", "INT")
.watermark("event_time", "SOURCE_WATERMARK()")
.build())
table.print_schema()
# prints:
# (
# `event_time` TIMESTAMP_LTZ(3) *ROWTIME*,
# `name` STRING,
# `score` INT
# )
# note: the watermark strategy is not shown due to the inserted column reordering projection
예시 1은 시간 기반 연산이 필요 없는 간단한 사용 사례를 보여줍니다.
예시 4는 윈도우나 interval join 같은 시간 기반 연산이 파이프라인의 일부여야 하는 가장 일반적인 사용 사례입니다. 예시 2는 이러한 시간 기반 연산이 processing time에서 동작해야 하는 가장 일반적인 사용 사례입니다.
예시 5는 전적으로 사용자의 선언에 의존합니다. 이는 DataStream API의 제네릭 타입(Table API에서는 RAW가 됨)을 적절한 데이터 타입으로 대체하는 데 유용할 수 있습니다.
DataType는 TypeInformation보다 풍부하므로, 불변(immutable) POJO 및 기타 복잡한 데이터 구조를 쉽게 활성화할 수 있습니다. 다음 Java 예시는 가능한 것을 보여줍니다. 지원되는 타입에 대한 자세한 내용은 DataStream API의 Data Types & Serialization 페이지도 확인하세요.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
// the DataStream API does not support immutable POJOs yet,
// the class will result in a generic type that is a RAW type in Table API by default
public static class User {
public final String name;
public final Integer score;
public User(String name, Integer score) {
this.name = name;
this.score = score;
}
}
// create a DataStream
DataStream<User> dataStream = env.fromElements(
new User("Alice", 4),
new User("Bob", 6),
new User("Alice", 10));
// since fields of a RAW type cannot be accessed, every stream record is treated as an atomic type
// leading to a table with a single column `f0`
Table table = tableEnv.fromDataStream(dataStream);
table.printSchema();
// prints:
// (
// `f0` RAW('User', '...')
// )
// instead, declare a more useful data type for columns using the Table API's type system
// in a custom schema and rename the columns in a following `as` projection
Table table = tableEnv
.fromDataStream(
dataStream,
Schema.newBuilder()
.column("f0", DataTypes.of(User.class))
.build())
.as("user");
table.printSchema();
// prints:
// (
// `user` *User<`name` STRING,`score` INT>*
// )
// data types can be extracted reflectively as above or explicitly defined
Table table = tableEnv
.fromDataStream(
dataStream,
Schema.newBuilder()
.column(
"f0",
DataTypes.STRUCTURED(
User.class,
DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("score", DataTypes.INT())))
.build())
.as("user");
table.printSchema();
// prints:
// (
// `user` *User<`name` STRING,`score` INT>*
// )
Python
Custom PoJo Class is unsupported in PyFlink now.
createTemporaryView 예시
DataStream은 (스키마로 보강될 수 있는) 뷰로 직접 등록될 수 있습니다.
DataStream에서 생성된 뷰는 임시 뷰로만 등록할 수 있습니다. 인라인/익명 특성 때문에 영구 카탈로그에는 등록할 수 없습니다.
다음 코드는 다양한 시나리오에서 createTemporaryView를 사용하는 방법을 보여줍니다.
Java
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.Schema;
// create some DataStream
DataStream<Tuple2<Long, String>> dataStream = env.fromElements(
Tuple2.of(12L, "Alice"),
Tuple2.of(0L, "Bob"));
// === EXAMPLE 1 ===
// register the DataStream as view "MyView" in the current session
// all columns are derived automatically
tableEnv.createTemporaryView("MyView", dataStream);
tableEnv.from("MyView").printSchema();
// prints:
// (
// `f0` BIGINT NOT NULL,
// `f1` STRING
// )
// === EXAMPLE 2 ===
// register the DataStream as view "MyView" in the current session,
// provide a schema to adjust the columns similar to `fromDataStream`
// in this example, the derived NOT NULL information has been removed
tableEnv.createTemporaryView(
"MyView",
dataStream,
Schema.newBuilder()
.column("f0", "BIGINT")
.column("f1", "STRING")
.build());
tableEnv.from("MyView").printSchema();
// prints:
// (
// `f0` BIGINT,
// `f1` STRING
// )
// === EXAMPLE 3 ===
// use the Table API before creating the view if it is only about renaming columns
tableEnv.createTemporaryView(
"MyView",
tableEnv.fromDataStream(dataStream).as("id", "name"));
tableEnv.from("MyView").printSchema();
// prints:
// (
// `id` BIGINT NOT NULL,
// `name` STRING
// )
Scala
// create some DataStream
val dataStream: DataStream[(Long, String)] = env.fromElements(
(12L, "Alice"),
(0L, "Bob"))
// === EXAMPLE 1 ===
// register the DataStream as view "MyView" in the current session
// all columns are derived automatically
tableEnv.createTemporaryView("MyView", dataStream)
tableEnv.from("MyView").printSchema()
// prints:
// (
// `_1` BIGINT NOT NULL,
// `_2` STRING
// )
// === EXAMPLE 2 ===
// register the DataStream as view "MyView" in the current session,
// provide a schema to adjust the columns similar to `fromDataStream`
// in this example, the derived NOT NULL information has been removed
tableEnv.createTemporaryView(
"MyView",
dataStream,
Schema.newBuilder()
.column("_1", "BIGINT")
.column("_2", "STRING")
.build())
tableEnv.from("MyView").printSchema()
// prints:
// (
// `_1` BIGINT,
// `_2` STRING
// )
// === EXAMPLE 3 ===
// use the Table API before creating the view if it is only about renaming columns
tableEnv.createTemporaryView(
"MyView",
tableEnv.fromDataStream(dataStream).as("id", "name"))
tableEnv.from("MyView").printSchema()
// prints:
// (
// `id` BIGINT NOT NULL,
// `name` STRING
// )
Python
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import DataTypes, StreamTableEnvironment, Schema
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
ds = env.from_collection([(12, "Alice"), (0, "Bob")], type_info=Types.TUPLE([Types.LONG(), Types.STRING()]))
# === EXAMPLE 1 ===
# register the DataStream as view "MyView" in the current session
# all columns are derived automatically
t_env.create_temporary_view("MyView", ds)
t_env.from_path("MyView").print_schema()
# prints:
# (
# `f0` BIGINT NOT NULL,
# `f1` STRING
# )
# === EXAMPLE 2 ===
# register the DataStream as view "MyView" in the current session,
# provide a schema to adjust the columns similar to `fromDataStream`
# in this example, the derived NOT NULL information has been removed
t_env.create_temporary_view(
"MyView",
ds,
Schema.new_builder()
.column("f0", "BIGINT")
.column("f1", "STRING")
.build())
t_env.from_path("MyView").print_schema()
# prints:
# (
# `f0` BIGINT,
# `f1` STRING
# )
# === EXAMPLE 3 ===
# use the Table API before creating the view if it is only about renaming columns
t_env.create_temporary_view(
"MyView",
t_env.from_data_stream(ds).alias("id", "name"))
t_env.from_path("MyView").print_schema()
# prints:
# (
# `id` BIGINT NOT NULL,
# `name` STRING
# )
toDataStream 예시
다음 코드는 다양한 시나리오에서 toDataStream을 사용하는 방법을 보여줍니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.Table;
import org.apache.flink.types.Row;
import java.time.Instant;
// POJO with mutable fields
// since no fully assigning constructor is defined, the field order
// is alphabetical [event_time, name, score]
public static class User {
public String name;
public Integer score;
public Instant event_time;
}
tableEnv.executeSql(
"CREATE TABLE GeneratedTable "
+ "("
+ " name STRING,"
+ " score INT,"
+ " event_time TIMESTAMP_LTZ(3),"
+ " WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND"
+ ")"
+ "WITH ('connector'='datagen')");
Table table = tableEnv.from("GeneratedTable");
// === EXAMPLE 1 ===
// use the default conversion to instances of Row
// since `event_time` is a single rowtime attribute, it is inserted into the DataStream
// metadata and watermarks are propagated
DataStream<Row> dataStream = tableEnv.toDataStream(table);
// === EXAMPLE 2 ===
// a data type is extracted from class `User`,
// the planner reorders fields and inserts implicit casts where possible to convert internal
// data structures to the desired structured type
// since `event_time` is a single rowtime attribute, it is inserted into the DataStream
// metadata and watermarks are propagated
DataStream<User> dataStream = tableEnv.toDataStream(table, User.class);
// === EXAMPLE 3 ===
// data types can be extracted reflectively as above or explicitly defined
DataStream<User> dataStream =
tableEnv.toDataStream(
table,
DataTypes.STRUCTURED(
User.class,
DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("score", DataTypes.INT()),
DataTypes.FIELD("event_time", DataTypes.TIMESTAMP_LTZ(3))));
Scala
import org.apache.flink.streaming.api.scala.DataStream
import org.apache.flink.table.api.DataTypes
case class User(name: String, score: java.lang.Integer, event_time: java.time.Instant)
tableEnv.executeSql(
"""
CREATE TABLE GeneratedTable (
name STRING,
score INT,
event_time TIMESTAMP_LTZ(3),
WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
)
WITH ('connector'='datagen')
"""
)
val table = tableEnv.from("GeneratedTable")
// === EXAMPLE 1 ===
// use the default conversion to instances of Row
// since `event_time` is a single rowtime attribute, it is inserted into the DataStream
// metadata and watermarks are propagated
val dataStream: DataStream[Row] = tableEnv.toDataStream(table)
// === EXAMPLE 2 ===
// a data type is extracted from class `User`,
// the planner reorders fields and inserts implicit casts where possible to convert internal
// data structures to the desired structured type
// since `event_time` is a single rowtime attribute, it is inserted into the DataStream
// metadata and watermarks are propagated
val dataStream: DataStream[User] = tableEnv.toDataStream(table, classOf[User])
// === EXAMPLE 3 ===
// data types can be extracted reflectively as above or explicitly defined
val dataStream: DataStream[User] =
tableEnv.toDataStream(
table,
DataTypes.STRUCTURED(
classOf[User],
DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("score", DataTypes.INT()),
DataTypes.FIELD("event_time", DataTypes.TIMESTAMP_LTZ(3))))
Python
t_env.execute_sql(
"CREATE TABLE GeneratedTable "
+ "("
+ " name STRING,"
+ " score INT,"
+ " event_time TIMESTAMP_LTZ(3),"
+ " WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND"
+ ")"
+ "WITH ('connector'='datagen')");
table = t_env.from_path("GeneratedTable");
# === EXAMPLE 1 ===
# use the default conversion to instances of Row
# since `event_time` is a single rowtime attribute, it is inserted into the DataStream
# metadata and watermarks are propagated
ds = t_env.to_data_stream(table)
toDataStream은 업데이트되지 않는(non-updating) 테이블만 지원한다는 점에 유의하세요. 보통 윈도우, interval join, MATCH_RECOGNIZE 절 같은 시간 기반 연산은 프로젝션과 필터 같은 단순 연산 옆에서 insert-only 파이프라인에 잘 맞습니다.
업데이트를 생성하는 연산이 있는 파이프라인은 toChangelogStream을 사용할 수 있습니다.
Changelog 스트림 처리
내부적으로 Flink의 테이블 런타임은 changelog 프로세서입니다. concepts 페이지는 동적 테이블과 스트림이 서로 어떻게 관련되는지를 설명합니다.
StreamTableEnvironment는 이러한 변경 데이터 캡처(change data capture, CDC) 기능을 노출하기 위해 다음 메서드를 제공합니다:
fromChangelogStream(DataStream): changelog 항목의 스트림을 테이블로 해석합니다. 런타임에RowKind플래그가 평가되므로 스트림 레코드 타입은org.apache.flink.types.Row여야 합니다. event-time과 워터마크는 기본적으로 전파되지 않습니다. 이 메서드는 기본ChangelogMode로 (즉org.apache.flink.types.RowKind에 나열된) 모든 종류의 변경을 포함하는 changelog를 기대합니다.fromChangelogStream(DataStream, Schema):fromDataStream(DataStream, Schema)와 유사하게DataStream에 대한 스키마를 정의할 수 있습니다. 그 외의 의미론은fromChangelogStream(DataStream)과 동일합니다.fromChangelogStream(DataStream, Schema, ChangelogMode): 스트림을 changelog로 해석하는 방법을 완전히 제어합니다. 전달된ChangelogMode는 플래너가 insert-only, upsert, 또는 retract 동작을 구분하는 데 도움을 줍니다.toChangelogStream(Table):fromChangelogStream(DataStream)의 역연산입니다.org.apache.flink.types.Row인스턴스로 이루어진 스트림을 생성하고 런타임에 모든 레코드에RowKind플래그를 설정합니다. 이 메서드는 모든 종류의 업데이트 테이블을 지원합니다. 입력 테이블에 단일 rowtime 컬럼이 있으면 스트림 레코드의 타임스탬프로 전파됩니다. 워터마크도 전파됩니다.toChangelogStream(Table, Schema):fromChangelogStream(DataStream, Schema)의 역연산입니다. 이 메서드는 생성되는 컬럼 데이터 타입을 보강할 수 있습니다. 필요한 경우 플래너는 암시적 캐스트를 삽입할 수 있습니다. rowtime을 메타데이터 컬럼으로 기록하는 것도 가능합니다.toChangelogStream(Table, Schema, ChangelogMode): 테이블을 changelog 스트림으로 변환하는 방법을 완전히 제어합니다. 전달된ChangelogMode는 플래너가 insert-only, upsert, 또는 retract 동작을 구분하는 데 도움을 줍니다.
Table API 관점에서 DataStream API로/로부터 변환하는 것은 SQL에서 CREATE TABLE DDL로 정의된 가상 테이블 커넥터에서 읽거나 쓰는 것과 유사합니다.
fromChangelogStream은 fromDataStream과 유사하게 동작하므로, 여기를 계속하기 전에 이전 섹션을 읽는 것을 권장합니다.
이 가상 커넥터는 스트림 레코드의 rowtime 메타데이터 읽기와 쓰기도 지원합니다.
가상 테이블 소스는 SupportsSourceWatermark를 구현합니다.
fromChangelogStream 예시
다음 코드는 다양한 시나리오에서 fromChangelogStream을 사용하는 방법을 보여줍니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.connector.ChangelogMode;
import org.apache.flink.types.Row;
import org.apache.flink.types.RowKind;
// === EXAMPLE 1 ===
// interpret the stream as a retract stream
// create a changelog DataStream
DataStream<Row> dataStream =
env.fromElements(
Row.ofKind(RowKind.INSERT, "Alice", 12),
Row.ofKind(RowKind.INSERT, "Bob", 5),
Row.ofKind(RowKind.UPDATE_BEFORE, "Alice", 12),
Row.ofKind(RowKind.UPDATE_AFTER, "Alice", 100));
// interpret the DataStream as a Table
Table table = tableEnv.fromChangelogStream(dataStream);
// register the table under a name and perform an aggregation
tableEnv.createTemporaryView("InputTable", table);
tableEnv
.executeSql("SELECT f0 AS name, SUM(f1) AS score FROM InputTable GROUP BY f0")
.print();
// prints:
// +----+--------------------------------+-------------+
// | op | name | score |
// +----+--------------------------------+-------------+
// | +I | Bob | 5 |
// | +I | Alice | 12 |
// | -D | Alice | 12 |
// | +I | Alice | 100 |
// +----+--------------------------------+-------------+
// === EXAMPLE 2 ===
// interpret the stream as an upsert stream (without a need for UPDATE_BEFORE)
// create a changelog DataStream
DataStream<Row> dataStream =
env.fromElements(
Row.ofKind(RowKind.INSERT, "Alice", 12),
Row.ofKind(RowKind.INSERT, "Bob", 5),
Row.ofKind(RowKind.UPDATE_AFTER, "Alice", 100));
// interpret the DataStream as a Table
Table table =
tableEnv.fromChangelogStream(
dataStream,
Schema.newBuilder().primaryKey("f0").build(),
ChangelogMode.upsert());
// register the table under a name and perform an aggregation
tableEnv.createTemporaryView("InputTable", table);
tableEnv
.executeSql("SELECT f0 AS name, SUM(f1) AS score FROM InputTable GROUP BY f0")
.print();
// prints:
// +----+--------------------------------+-------------+
// | op | name | score |
// +----+--------------------------------+-------------+
// | +I | Bob | 5 |
// | +I | Alice | 12 |
// | -U | Alice | 12 |
// | +U | Alice | 100 |
// +----+--------------------------------+-------------+
Scala
import org.apache.flink.api.scala.typeutils.Types
import org.apache.flink.table.api.Schema
import org.apache.flink.table.connector.ChangelogMode
import org.apache.flink.types.{Row, RowKind}
// === EXAMPLE 1 ===
// interpret the stream as a retract stream
// create a changelog DataStream
val dataStream = env.fromElements(
Row.ofKind(RowKind.INSERT, "Alice", Int.box(12)),
Row.ofKind(RowKind.INSERT, "Bob", Int.box(5)),
Row.ofKind(RowKind.UPDATE_BEFORE, "Alice", Int.box(12)),
Row.ofKind(RowKind.UPDATE_AFTER, "Alice", Int.box(100))
)(Types.ROW(Types.STRING, Types.INT))
// interpret the DataStream as a Table
val table = tableEnv.fromChangelogStream(dataStream)
// register the table under a name and perform an aggregation
tableEnv.createTemporaryView("InputTable", table)
tableEnv
.executeSql("SELECT f0 AS name, SUM(f1) AS score FROM InputTable GROUP BY f0")
.print()
// prints:
// +----+--------------------------------+-------------+
// | op | name | score |
// +----+--------------------------------+-------------+
// | +I | Bob | 5 |
// | +I | Alice | 12 |
// | -D | Alice | 12 |
// | +I | Alice | 100 |
// +----+--------------------------------+-------------+
// === EXAMPLE 2 ===
// interpret the stream as an upsert stream (without a need for UPDATE_BEFORE)
// create a changelog DataStream
val dataStream = env.fromElements(
Row.ofKind(RowKind.INSERT, "Alice", Int.box(12)),
Row.ofKind(RowKind.INSERT, "Bob", Int.box(5)),
Row.ofKind(RowKind.UPDATE_AFTER, "Alice", Int.box(100))
)(Types.ROW(Types.STRING, Types.INT))
// interpret the DataStream as a Table
val table =
tableEnv.fromChangelogStream(
dataStream,
Schema.newBuilder().primaryKey("f0").build(),
ChangelogMode.upsert())
// register the table under a name and perform an aggregation
tableEnv.createTemporaryView("InputTable", table)
tableEnv
.executeSql("SELECT f0 AS name, SUM(f1) AS score FROM InputTable GROUP BY f0")
.print()
// prints:
// +----+--------------------------------+-------------+
// | op | name | score |
// +----+--------------------------------+-------------+
// | +I | Bob | 5 |
// | +I | Alice | 12 |
// | -U | Alice | 12 |
// | +U | Alice | 100 |
// +----+--------------------------------+-------------+
Python
from pyflink.common import Row, RowKind
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import DataTypes, StreamTableEnvironment, Schema
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# === EXAMPLE 1 ===
# create a changelog DataStream
ds = env.from_collection([
Row.of_kind(RowKind.INSERT, "Alice", 12),
Row.of_kind(RowKind.INSERT, "Bob", 5),
Row.of_kind(RowKind.UPDATE_BEFORE, "Alice", 12),
Row.of_kind(RowKind.UPDATE_AFTER, "Alice", 100)],
type_info=Types.ROW([Types.STRING(),Types.INT()]))
# interpret the DataStream as a Table
table = t_env.from_changelog_stream(ds)
# register the table under a name and perform an aggregation
t_env.create_temporary_view("InputTable", table)
t_env.execute_sql("SELECT f0 AS name, SUM(f1) AS score FROM InputTable GROUP BY f0").print()
# prints:
# +----+--------------------------------+-------------+
# | op | name | score |
# +----+--------------------------------+-------------+
# | +I | Bob | 5 |
# | +I | Alice | 12 |
# | -D | Alice | 12 |
# | +I | Alice | 100 |
# +----+--------------------------------+-------------+
# === EXAMPLE 2 ===
# interpret the stream as an upsert stream (without a need for UPDATE_BEFORE)
# create a changelog DataStream
ds = env.from_collection([
Row.of_kind(RowKind.INSERT, "Alice", 12),
Row.of_kind(RowKind.INSERT, "Bob", 5),
Row.of_kind(RowKind.UPDATE_AFTER, "Alice", 100)],
type_info=Types.ROW([Types.STRING(),Types.INT()]))
# interpret the DataStream as a Table
table = t_env.from_changelog_stream(
ds,
Schema.new_builder().primary_key("f0").build(),
ChangelogMode.upsert())
# register the table under a name and perform an aggregation
t_env.create_temporary_view("InputTable", table)
t_env.execute_sql("SELECT f0 AS name, SUM(f1) AS score FROM InputTable GROUP BY f0").print()
# prints:
# +----+--------------------------------+-------------+
# | op | name | score |
# +----+--------------------------------+-------------+
# | +I | Bob | 5 |
# | +I | Alice | 12 |
# | -U | Alice | 12 |
# | +U | Alice | 100 |
# +----+--------------------------------+-------------+
예시 1에 나온 기본 ChangelogMode는 모든 종류의 변경을 수용하므로 대부분의 사용 사례에서 충분해야 합니다.
그러나 예시 2는 upsert 모드를 사용해 업데이트 메시지 수를 50% 줄임으로써 효율성을 위해 들어오는 변경 종류를 제한하는 방법을 보여줍니다. toChangelogStream에 기본 키와 upsert changelog 모드를 정의하면 결과 메시지 수를 줄일 수 있습니다.
toChangelogStream 예시
다음 코드는 다양한 시나리오에서 toChangelogStream을 사용하는 방법을 보여줍니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.data.StringData;
import org.apache.flink.types.Row;
import org.apache.flink.util.Collector;
import static org.apache.flink.table.api.Expressions.*;
// create Table with event-time
tableEnv.executeSql(
"CREATE TABLE GeneratedTable "
+ "("
+ " name STRING,"
+ " score INT,"
+ " event_time TIMESTAMP_LTZ(3),"
+ " WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND"
+ ")"
+ "WITH ('connector'='datagen')");
Table table = tableEnv.from("GeneratedTable");
// === EXAMPLE 1 ===
// convert to DataStream in the simplest and most general way possible (no event-time)
Table simpleTable = tableEnv
.fromValues(row("Alice", 12), row("Alice", 2), row("Bob", 12))
.as("name", "score")
.groupBy($("name"))
.select($("name"), $("score").sum());
tableEnv
.toChangelogStream(simpleTable)
.executeAndCollect()
.forEachRemaining(System.out::println);
// prints:
// +I[Bob, 12]
// +I[Alice, 12]
// -U[Alice, 12]
// +U[Alice, 14]
// === EXAMPLE 2 ===
// convert to DataStream in the simplest and most general way possible (with event-time)
DataStream<Row> dataStream = tableEnv.toChangelogStream(table);
// since `event_time` is a single time attribute in the schema, it is set as the
// stream record's timestamp by default; however, at the same time, it remains part of the Row
dataStream.process(
new ProcessFunction<Row, Void>() {
@Override
public void processElement(Row row, Context ctx, Collector<Void> out) {
// prints: [name, score, event_time]
System.out.println(row.getFieldNames(true));
// timestamp exists twice
assert ctx.timestamp() == row.<Instant>getFieldAs("event_time").toEpochMilli();
}
});
env.execute();
// === EXAMPLE 3 ===
// convert to DataStream but write out the time attribute as a metadata column which means
// it is not part of the physical schema anymore
DataStream<Row> dataStream = tableEnv.toChangelogStream(
table,
Schema.newBuilder()
.column("name", "STRING")
.column("score", "INT")
.columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)")
.build());
// the stream record's timestamp is defined by the metadata; it is not part of the Row
dataStream.process(
new ProcessFunction<Row, Void>() {
@Override
public void processElement(Row row, Context ctx, Collector<Void> out) {
// prints: [name, score]
System.out.println(row.getFieldNames(true));
// timestamp exists once
System.out.println(ctx.timestamp());
}
});
env.execute();
// === EXAMPLE 4 ===
// for advanced users, it is also possible to use more internal data structures for efficiency
// note that this is only mentioned here for completeness because using internal data structures
// adds complexity and additional type handling
// however, converting a TIMESTAMP_LTZ column to `Long` or STRING to `byte[]` might be convenient,
// also structured types can be represented as `Row` if needed
DataStream<Row> dataStream = tableEnv.toChangelogStream(
table,
Schema.newBuilder()
.column(
"name",
DataTypes.STRING().bridgedTo(StringData.class))
.column(
"score",
DataTypes.INT())
.column(
"event_time",
DataTypes.TIMESTAMP_LTZ(3).bridgedTo(Long.class))
.build());
// leads to a stream of Row(name: StringData, score: Integer, event_time: Long)
Scala
import org.apache.flink.api.scala._
import org.apache.flink.streaming.api.functions.ProcessFunction
import org.apache.flink.streaming.api.scala.DataStream
import org.apache.flink.table.api._
import org.apache.flink.types.Row
import org.apache.flink.util.Collector
import java.time.Instant
// create Table with event-time
tableEnv.executeSql(
"""
CREATE TABLE GeneratedTable (
name STRING,
score INT,
event_time TIMESTAMP_LTZ(3),
WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
)
WITH ('connector'='datagen')
"""
)
val table = tableEnv.from("GeneratedTable")
// === EXAMPLE 1 ===
// convert to DataStream in the simplest and most general way possible (no event-time)
val simpleTable = tableEnv
.fromValues(row("Alice", 12), row("Alice", 2), row("Bob", 12))
.as("name", "score")
.groupBy($"name")
.select($"name", $"score".sum())
tableEnv
.toChangelogStream(simpleTable)
.executeAndCollect()
.foreach(println)
// prints:
// +I[Bob, 12]
// +I[Alice, 12]
// -U[Alice, 12]
// +U[Alice, 14]
// === EXAMPLE 2 ===
// convert to DataStream in the simplest and most general way possible (with event-time)
val dataStream: DataStream[Row] = tableEnv.toChangelogStream(table)
// since `event_time` is a single time attribute in the schema, it is set as the
// stream record's timestamp by default; however, at the same time, it remains part of the Row
dataStream.process(new ProcessFunction[Row, Unit] {
override def processElement(
row: Row,
ctx: ProcessFunction[Row, Unit]#Context,
out: Collector[Unit]): Unit = {
// prints: [name, score, event_time]
println(row.getFieldNames(true))
// timestamp exists twice
assert(ctx.timestamp() == row.getFieldAs[Instant]("event_time").toEpochMilli)
}
})
env.execute()
// === EXAMPLE 3 ===
// convert to DataStream but write out the time attribute as a metadata column which means
// it is not part of the physical schema anymore
val dataStream: DataStream[Row] = tableEnv.toChangelogStream(
table,
Schema.newBuilder()
.column("name", "STRING")
.column("score", "INT")
.columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)")
.build())
// the stream record's timestamp is defined by the metadata; it is not part of the Row
dataStream.process(new ProcessFunction[Row, Unit] {
override def processElement(
row: Row,
ctx: ProcessFunction[Row, Unit]#Context,
out: Collector[Unit]): Unit = {
// prints: [name, score]
println(row.getFieldNames(true))
// timestamp exists once
println(ctx.timestamp())
}
})
env.execute()
// === EXAMPLE 4 ===
// for advanced users, it is also possible to use more internal data structures for better
// efficiency
// note that this is only mentioned here for completeness because using internal data structures
// adds complexity and additional type handling
// however, converting a TIMESTAMP_LTZ column to `Long` or STRING to `byte[]` might be convenient,
// also structured types can be represented as `Row` if needed
val dataStream: DataStream[Row] = tableEnv.toChangelogStream(
table,
Schema.newBuilder()
.column(
"name",
DataTypes.STRING().bridgedTo(classOf[StringData]))
.column(
"score",
DataTypes.INT())
.column(
"event_time",
DataTypes.TIMESTAMP_LTZ(3).bridgedTo(class[Long]))
.build())
// leads to a stream of Row(name: StringData, score: Integer, event_time: Long)
Python
from pyflink.common import Row
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.functions import ProcessFunction
from pyflink.table import DataTypes, StreamTableEnvironment, Schema
from pyflink.table.expressions import col
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# create Table with event-time
t_env.execute_sql(
"CREATE TABLE GeneratedTable "
+ "("
+ " name STRING,"
+ " score INT,"
+ " event_time TIMESTAMP_LTZ(3),"
+ " WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND"
+ ")"
+ "WITH ('connector'='datagen')")
table = t_env.from_path("GeneratedTable")
# === EXAMPLE 1 ===
# convert to DataStream in the simplest and most general way possible (no event-time)
simple_table = t_env.from_elements([Row("Alice", 12), Row("Alice", 2), Row("Bob", 12)],
DataTypes.ROW([DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("score", DataTypes.INT())]))
simple_table = simple_table.group_by(col('name')).select(col('name'), col('score').sum)
t_env.to_changelog_stream(simple_table).print()
env.execute()
# prints:
# +I[Bob, 12]
# +I[Alice, 12]
# -U[Alice, 12]
# +U[Alice, 14]
# === EXAMPLE 2 ===
# convert to DataStream in the simplest and most general way possible (with event-time)
ds = t_env.to_changelog_stream(table)
# since `event_time` is a single time attribute in the schema, it is set as the
# stream record's timestamp by default; however, at the same time, it remains part of the Row
class MyProcessFunction(ProcessFunction):
def process_element(self, row, ctx):
print(row)
assert ctx.timestamp() == row.event_time.to_epoch_milli()
ds.process(MyProcessFunction())
env.execute()
# === EXAMPLE 3 ===
# convert to DataStream but write out the time attribute as a metadata column which means
# it is not part of the physical schema anymore
ds = t_env.to_changelog_stream(
table,
Schema.new_builder()
.column("name", "STRING")
.column("score", "INT")
.column_by_metadata("rowtime", "TIMESTAMP_LTZ(3)")
.build())
class MyProcessFunction(ProcessFunction):
def process_element(self, row, ctx):
print(row)
print(ctx.timestamp())
ds.process(MyProcessFunction())
env.execute()
예시 4에서 데이터 타입에 대해 지원되는 변환에 대한 자세한 내용은 Table API의 Data Types 페이지를 참조하세요.
toChangelogStream(Table).executeAndCollect()의 동작은 Table.execute().collect()를 호출하는 것과 같습니다. 그러나 toChangelogStream(Table)은 DataStream API의 후속 ProcessFunction에서 생성된 워터마크에 접근할 수 있게 해주므로 테스트에 더 유용할 수 있습니다.
DataStream API에 Table API 파이프라인 추가하기
단일 Flink 작업은 서로 나란히 실행되는 여러 개의 연결되지 않은 파이프라인으로 구성될 수 있습니다.
Table API에서 정의된 소스-투-싱크 파이프라인은 전체로서 StreamExecutionEnvironment에 첨부될 수 있으며, DataStream API에서 execute 메서드 중 하나를 호출할 때 제출됩니다.
그러나 소스가 반드시 테이블 소스일 필요는 없으며, 이전에 Table API로 변환된 다른 DataStream 파이프라인일 수도 있습니다. 따라서 DataStream API 프로그램에 테이블 싱크를 사용할 수 있습니다.
이 기능은 StreamTableEnvironment.createStatementSet()로 생성된 전용 StreamStatementSet 인스턴스를 통해 사용할 수 있습니다. statement set을 사용하면 플래너가 추가된 모든 명령문을 함께 최적화하고, StreamStatementSet.attachAsDataStream()을 호출할 때 StreamExecutionEnvironment에 추가되는 하나 이상의 end-to-end 파이프라인을 만들 수 있습니다.
다음 예시는 단일 작업 내에서 DataStream API 프로그램에 테이블 프로그램을 추가하는 방법을 보여줍니다.
Java
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink;
import org.apache.flink.table.api.*;
import org.apache.flink.table.api.bridge.java.*;
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
StreamStatementSet statementSet = tableEnv.createStatementSet();
// create some source
TableDescriptor sourceDescriptor =
TableDescriptor.forConnector("datagen")
.option("number-of-rows", "3")
.schema(
Schema.newBuilder()
.column("myCol", DataTypes.INT())
.column("myOtherCol", DataTypes.BOOLEAN())
.build())
.build();
// create some sink
TableDescriptor sinkDescriptor = TableDescriptor.forConnector("print").build();
// add a pure Table API pipeline
Table tableFromSource = tableEnv.from(sourceDescriptor);
statementSet.add(tableFromSource.insertInto(sinkDescriptor));
// use table sinks for the DataStream API pipeline
DataStream<Integer> dataStream = env.fromElements(1, 2, 3);
Table tableFromStream = tableEnv.fromDataStream(dataStream);
statementSet.add(tableFromStream.insertInto(sinkDescriptor));
// attach both pipelines to StreamExecutionEnvironment
// (the statement set will be cleared after calling this method)
statementSet.attachAsDataStream();
// define other DataStream API parts
env.fromElements(4, 5, 6).sinkTo(new DiscardingSink<>());
// use DataStream API to submit the pipelines
env.execute();
// prints similar to:
// +I[1618440447, false]
// +I[1259693645, true]
// +I[158588930, false]
// +I[1]
// +I[2]
// +I[3]
Scala
import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink
import org.apache.flink.streaming.api.scala._
import org.apache.flink.table.api._
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env)
val statementSet = tableEnv.createStatementSet()
// create some source
val sourceDescriptor = TableDescriptor.forConnector("datagen")
.option("number-of-rows", "3")
.schema(Schema.newBuilder
.column("myCol", DataTypes.INT)
.column("myOtherCol", DataTypes.BOOLEAN).build)
.build
// create some sink
val sinkDescriptor = TableDescriptor.forConnector("print").build
// add a pure Table API pipeline
val tableFromSource = tableEnv.from(sourceDescriptor)
statementSet.add(tableFromSource.insertInto(sinkDescriptor))
// use table sinks for the DataStream API pipeline
val dataStream = env.fromElements(1, 2, 3)
val tableFromStream = tableEnv.fromDataStream(dataStream)
statementSet.add(tableFromStream.insertInto(sinkDescriptor))
// attach both pipelines to StreamExecutionEnvironment
// (the statement set will be cleared calling this method)
statementSet.attachAsDataStream()
// define other DataStream API parts
env.fromElements(4, 5, 6).sinkTo(new DiscardingSink[Int]())
// now use DataStream API to submit the pipelines
env.execute()
// prints similar to:
// +I[1618440447, false]
// +I[1259693645, true]
// +I[158588930, false]
// +I[1]
// +I[2]
// +I[3]
Python
from pyflink.common import Encoder
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors.file_system import FileSink
from pyflink.table import StreamTableEnvironment, TableDescriptor, Schema, DataTypes
env = StreamExecutionEnvironment.get_execution_environment()
table_env = StreamTableEnvironment.create(env)
statement_set = table_env.create_statement_set()
# create some source
source_descriptor = TableDescriptor.for_connector("datagen") \
.option("number-of-rows", "3") \
.schema(
Schema.new_builder()
.column("my_col", DataTypes.INT())
.column("my_other_col", DataTypes.BOOLEAN())
.build()) \
.build()
# create some sink
sink_descriptor = TableDescriptor.for_connector("print").build()
# add a pure Table API pipeline
table_from_source = table_env.from_descriptor(source_descriptor)
statement_set.add_insert(sink_descriptor, table_from_source)
# use table sinks for the DataStream API pipeline
data_stream = env.from_collection([1, 2, 3])
table_from_stream = table_env.from_data_stream(data_stream)
statement_set.add_insert(sink_descriptor, table_from_stream)
# attach both pipelines to StreamExecutionEnvironment
# (the statement set will be cleared after calling this method)
statement_set.attach_as_datastream()
# define other DataStream API parts
env.from_collection([4, 5, 6]) \
.add_sink(FileSink
.for_row_format('/tmp/output', Encoder.simple_string_encoder())
.build())
# use DataStream API to submit the pipelines
env.execute()
# prints similar to:
# +I[1618440447, false]
# +I[1259693645, true]
# +I[158588930, false]
# +I[1]
# +I[2]
# +I[3]
Scala의 암시적 변환(Implicit Conversions)
Scala API 사용자는 Scala의 implicit 기능을 활용해 위의 모든 변환 메서드를 더 유창하게 사용할 수 있습니다.
org.apache.flink.table.api.bridge.scala._를 통해 패키지 객체를 임포트하면 이러한 implicit을 API에서 사용할 수 있습니다.
활성화하면 toTable이나 toChangelogTable 같은 메서드를 DataStream 객체에서 직접 호출할 수 있습니다. 마찬가지로 toDataStream과 toChangelogStream은 Table 객체에서 사용할 수 있습니다. 또한 DataStream[Row]에 특화된 DataStream API 메서드를 요청하면 Table 객체가 changelog 스트림으로 변환됩니다.
implicit 변환의 사용은 항상 의식적인 결정이어야 합니다. IDE가 실제 Table API 메서드를 제안하는지, 아니면 implicit을 통한 DataStream API 메서드를 제안하는지 주의해야 합니다.
예를 들어,
table.execute().collect()는 Table API에 머무르지만table.executeAndCollect()는 implicit으로 DataStream API의executeAndCollect()메서드를 사용하므로 API 변환을 강제합니다.
import org.apache.flink.streaming.api.scala._
import org.apache.flink.table.api.bridge.scala._
import org.apache.flink.types.Row
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env)
val dataStream: DataStream[(Int, String)] = env.fromElements((42, "hello"))
// call toChangelogTable() implicitly on the DataStream object
val table: Table = dataStream.toChangelogTable(tableEnv)
// force implicit conversion
val dataStreamAgain1: DataStream[Row] = table
// call toChangelogStream() implicitly on the Table object
val dataStreamAgain2: DataStream[Row] = table.toChangelogStream
TypeInformation과 DataType 간 매핑
DataStream API는 스트림에서 이동하는 레코드 타입을 설명하기 위해 org.apache.flink.api.common.typeinfo.TypeInformation 인스턴스를 사용합니다. 특히, 한 DataStream 연산자에서 다른 연산자로 레코드를 직렬화하고 역직렬화하는 방법을 정의합니다. 또한 상태를 savepoint와 checkpoint로 직렬화하는 데도 도움이 됩니다.
Table API는 레코드를 내부적으로 표현하기 위해 커스텀 데이터 구조를 사용하고, 사용자가 소스, 싱크, UDF, 또는 DataStream API에서 더 쉽게 사용할 수 있도록 데이터 구조가 변환되는 외부 형식을 선언하기 위해 org.apache.flink.table.types.DataType를 노출합니다.
DataType는 논리적 SQL 타입에 대한 세부 정보도 포함하므로 TypeInformation보다 풍부합니다. 따라서 변환 중에 일부 세부 사항이 암시적으로 추가됩니다.
Table의 컬럼 이름과 타입은 DataStream의 TypeInformation에서 자동으로 파생됩니다. DataStream API의 반사적(reflective) 타입 추출 기능을 통해 타입 정보가 올바르게 감지되었는지 확인하려면 DataStream.getType()을 사용하세요. 가장 바깥쪽 레코드의 TypeInformation이 CompositeType이면, 테이블 스키마를 파생할 때 첫 번째 레벨에서 평면화(flatten)됩니다.
DataStream API가 리플렉션을 기반으로 더 구체적인
TypeInformation을 항상 추출할 수 있는 것은 아닙니다. 이는 종종 조용히 발생하며, 일반 Kryo 직렬화기에 의해 뒷받침되는GenericTypeInfo로 이어집니다.예를 들어,
Row클래스는 반사적으로 분석될 수 없으며 항상 명시적 타입 정보 선언이 필요합니다. DataStream API에 적절한 타입 정보가 선언되지 않으면 행이RAW데이터 타입으로 표시되고 Table API는 해당 필드에 접근할 수 없습니다. Java에서는.map(...).returns(TypeInformation)을, Scala에서는.map(...)(TypeInformation)을 사용해 타입 정보를 명시적으로 선언하세요.
TypeInformation → DataType
TypeInformation을 DataType으로 변환할 때 다음 규칙이 적용됩니다:
TypeInformation의 모든 하위 클래스는 Flink 내장 직렬화기와 정렬된 nullability를 포함해 논리적 타입에 매핑됩니다.TupleTypeInfoBase의 하위 클래스는 (Row의 경우) 행 또는 (튜플, POJO, case class의 경우) 구조화된 타입으로 변환됩니다.BigDecimal은 기본적으로DECIMAL(38, 18)로 변환됩니다.PojoTypeInfo필드의 순서는 모든 필드를 매개변수로 갖는 생성자에 의해 결정됩니다. 변환 중에 해당 생성자가 발견되지 않으면 필드 순서는 알파벳순이 됩니다.GenericTypeInfo및 나열된org.apache.flink.table.api.DataTypes중 하나로 표현될 수 없는 기타TypeInformation은 블랙박스RAW타입으로 처리됩니다. 현재 세션 구성이 raw 타입의 직렬화기를 구체화하는 데 사용됩니다. 그러면 복합 중첩 필드에 접근할 수 없습니다.- 전체 변환 로직은 TypeInfoDataTypeConverter를 참조하세요.
커스텀 스키마 선언 또는 UDF에서 위 로직을 호출하려면 DataTypes.of(TypeInformation)을 사용하세요.
DataType → TypeInformation
테이블 런타임은 출력 레코드를 DataStream API의 첫 번째 연산자로 올바르게 직렬화하도록 보장합니다.
이후에는 DataStream API의 타입 정보 의미론을 고려해야 합니다.
레거시 변환(Legacy Conversion)
다음 섹션은 향후 버전에서 제거될 API의 오래된 부분을 설명합니다.
특히 이러한 부분은 최근의 많은 새 기능과 리팩토링에 잘 통합되지 않을 수 있습니다(예:
RowKind가 올바르게 설정되지 않음, 타입 시스템이 매끄럽게 통합되지 않음).
DataStream을 Table로 변환
DataStream은 StreamTableEnvironment에서 Table로 직접 변환될 수 있습니다. 결과 뷰의 스키마는 등록된 컬렉션의 데이터 타입에 따라 달라집니다.
Java
StreamTableEnvironment tableEnv = ...;
DataStream<Tuple2<Long, String>> stream = ...;
Table table2 = tableEnv.fromDataStream(stream, $("myLong"), $("myString"));
Scala
val tableEnv: StreamTableEnvironment = ???
val stream: DataStream[(Long, String)] = ???
val table2: Table = tableEnv.fromDataStream(stream, $"myLong", $"myString")
Python
t_env = ... # type: StreamTableEnvironment
stream = ... # type: DataStream of Types.TUPLE([Types.LONG(), Types.STRING()])
table2 = t_env.from_data_stream(stream, col('my_long'), col('my_stream'))
Table을 DataStream으로 변환
Table의 결과는 DataStream으로 변환될 수 있습니다. 이렇게 하면 Table API 또는 SQL 쿼리 결과에서 커스텀 DataStream 프로그램을 실행할 수 있습니다.
Table을 DataStream으로 변환할 때 결과 레코드의 데이터 타입, 즉 Table의 행이 변환될 데이터 타입을 지정해야 합니다. 대개 가장 편리한 변환 타입은 Row입니다. 다음 목록은 각 옵션의 기능 개요를 제공합니다:
- Row: 필드가 위치로 매핑되고, 필드 수는 무제한이며,
null값을 지원하지만 타입 안전 접근은 없습니다. - POJO: 필드가 이름으로 매핑되고(POJO 필드는
Table필드와 동일한 이름이어야 함), 필드 수는 무제한이며,null값을 지원하고 타입 안전 접근이 가능합니다. - Case Class: 필드가 위치로 매핑되고,
null값을 지원하지 않으며, 타입 안전 접근이 가능합니다. - Tuple: 필드가 위치로 매핑되고, 22(Scala) 또는 25(Java)개 필드로 제한되며,
null값을 지원하지 않고, 타입 안전 접근이 가능합니다. - Atomic Type:
Table에 단일 필드가 있어야 하고,null값을 지원하지 않으며, 타입 안전 접근이 가능합니다.
Table을 DataStream으로 변환
스트리밍 쿼리의 결과인 Table은 동적으로 업데이트됩니다. 즉, 쿼리의 입력 스트림에 새 레코드가 도착함에 따라 변경됩니다. 따라서 이러한 동적 쿼리가 변환되는 DataStream은 테이블의 업데이트를 인코딩해야 합니다.
Table을 DataStream으로 변환하는 두 가지 모드가 있습니다:
- Append 모드: 이 모드는 동적
Table이INSERT변경으로만 수정되는 경우, 즉 append-only이고 이전에 방출된 결과가 업데이트되지 않는 경우에만 사용할 수 있습니다. - Retract 모드: 이 모드는 항상 사용할 수 있습니다.
boolean플래그로INSERT와DELETE변경을 인코딩합니다.
Java
StreamTableEnvironment tableEnv = ...;
Table table = tableEnv.fromValues(
DataTypes.Row(
DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("age", DataTypes.INT()),
row("john", 35),
row("sarah", 32));
// Convert the Table into an append DataStream of Row by specifying the class
DataStream<Row> dsRow = tableEnv.toAppendStream(table, Row.class);
// Convert the Table into an append DataStream of Tuple2<String, Integer> with TypeInformation
TupleTypeInfo<Tuple2<String, Integer>> tupleType = new TupleTypeInfo<>(Types.STRING(), Types.INT());
DataStream<Tuple2<String, Integer>> dsTuple = tableEnv.toAppendStream(table, tupleType);
// Convert the Table into a retract DataStream of Row.
// A retract stream of type X is a DataStream<Tuple2<Boolean, X>>.
// The boolean field indicates the type of the change.
// True is INSERT, false is DELETE.
DataStream<Tuple2<Boolean, Row>> retractStream = tableEnv.toRetractStream(table, Row.class);
Scala
val tableEnv: StreamTableEnvironment = ???
// Table with two fields (String name, Integer age)
val table: Table = tableEnv.fromValues(
DataTypes.Row(
DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("age", DataTypes.INT()),
row("john", 35),
row("sarah", 32))
// Convert the Table into an append DataStream of Row by specifying the class
val dsRow: DataStream[Row] = tableEnv.toAppendStream[Row](table)
// Convert the Table into an append DataStream of (String, Integer) with TypeInformation
val dsTuple: DataStream[(String, Int)] dsTuple =
tableEnv.toAppendStream[(String, Int)](table)
// Convert the Table into a retract DataStream of Row.
// A retract stream of type X is a DataStream<Tuple2<Boolean, X>>.
// The boolean field indicates the type of the change.
// True is INSERT, false is DELETE.
val retractStream: DataStream[(Boolean, Row)] = tableEnv.toRetractStream[Row](table)
Python
from pyflink.table import DataTypes
from pyflink.common.typeinfo import Types
t_env = ...
table = t_env.from_elements([("john", 35), ("sarah", 32)],
DataTypes.ROW([DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("age", DataTypes.INT())]))
# Convert the Table into an append DataStream of Row by specifying the type information
ds_row = t_env.to_append_stream(table, Types.ROW([Types.STRING(), Types.INT()]))
# Convert the Table into an append DataStream of Tuple[str, int] with TypeInformation
ds_tuple = t_env.to_append_stream(table, Types.TUPLE([Types.STRING(), Types.INT()]))
# Convert the Table into a retract DataStream of Row by specifying the type information
# A retract stream of type X is a DataStream of Tuple[bool, X].
# The boolean field indicates the type of the change.
# True is INSERT, false is DELETE.
retract_stream = t_env.to_retract_stream(table, Types.ROW([Types.STRING(), Types.INT()]))
참고: 동적 테이블과 그 속성에 대한 자세한 논의는 Dynamic Tables 문서에서 확인할 수 있습니다.
Table이 DataStream으로 변환되면, DataStream 프로그램을 실행하기 위해
StreamExecutionEnvironment.execute()메서드를 사용하세요.
데이터 타입을 테이블 스키마로 매핑
Flink의 DataStream API는 매우 다양한 타입을 지원합니다. 튜플(내장 Scala, Flink Java 튜플, Python 튜플), POJO, Scala case class, Flink의 Row 타입 같은 복합 타입은 테이블 표현식에서 접근할 수 있는 여러 필드를 가진 중첩 데이터 구조를 허용합니다. 다른 타입은 원자적(atomic) 타입으로 처리됩니다. 다음에서는 Table API가 이러한 타입을 내부 행 표현으로 변환하는 방법을 설명하고 DataStream을 Table로 변환하는 예시를 보여줍니다.
데이터 타입을 테이블 스키마로 매핑하는 방법은 두 가지입니다: 필드 위치 기반(position-based) 또는 필드 이름 기반(name-based).
위치 기반 매핑(Position-based Mapping)
위치 기반 매핑은 필드 순서를 유지하면서 필드에 더 의미 있는 이름을 부여하는 데 사용할 수 있습니다. 이 매핑은 정의된 필드 순서가 있는 복합 데이터 타입과 원자적 타입에 사용할 수 있습니다. 튜플, 행, case class 같은 복합 데이터 타입에는 이러한 필드 순서가 있습니다. 그러나 POJO의 필드는 필드 이름을 기반으로 매핑해야 합니다(다음 섹션 참조). 필드는 프로젝션할 수 있지만 as(Java와 Scala) 또는 alias(Python) 별칭으로 이름을 바꿀 수는 없습니다.
위치 기반 매핑을 정의할 때 지정된 이름은 입력 데이터 타입에 존재하지 않아야 합니다. 그렇지 않으면 API는 매핑이 필드 이름을 기반으로 이루어져야 한다고 가정합니다. 필드 이름이 지정되지 않으면 복합 타입의 기본 필드 이름과 필드 순서가 사용되거나, 원자적 타입의 경우 f0이 사용됩니다.
Java
StreamTableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section;
DataStream<Tuple2<Long, Integer>> stream = ...;
// convert DataStream into Table with field "myLong" only
Table table = tableEnv.fromDataStream(stream, $("myLong"));
// convert DataStream into Table with field names "myLong" and "myInt"
Table table = tableEnv.fromDataStream(stream, $("myLong"), $("myInt"));
Scala
// get a TableEnvironment
val tableEnv: StreamTableEnvironment = ... // see "Create a TableEnvironment" section
val stream: DataStream[(Long, Int)] = ...
// convert DataStream into Table with field "myLong" only
val table: Table = tableEnv.fromDataStream(stream, $"myLong")
// convert DataStream into Table with field names "myLong" and "myInt"
val table: Table = tableEnv.fromDataStream(stream, $"myLong", $"myInt")
Python
from pyflink.table.expressions import col
# get a TableEnvironment
t_env = ... # see "Create a TableEnvironment" section
stream = ... # type: DataStream of Types.Tuple([Types.LONG(), Types.INT()])
# convert DataStream into Table with field "my_long" only
table = t_env.from_data_stream(stream, col('my_long'))
# convert DataStream into Table with field names "my_long" and "my_int"
table = t_env.from_data_stream(stream, col('my_long'), col('my_int'))
이름 기반 매핑(Name-based Mapping)
이름 기반 매핑은 POJO를 포함한 모든 데이터 타입에 사용할 수 있습니다. 테이블 스키마 매핑을 정의하는 가장 유연한 방법입니다. 매핑의 모든 필드는 이름으로 참조되며 as 별칭을 사용해 이름을 바꿀 수 있습니다. 필드는 재정렬하고 프로젝션할 수 있습니다.
필드 이름이 지정되지 않으면 복합 타입의 기본 필드 이름과 필드 순서가 사용되거나, 원자적 타입의 경우 f0이 사용됩니다.
Java
StreamTableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section
DataStream<Tuple2<Long, Integer>> stream = ...;
// convert DataStream into Table with field "f1" only
Table table = tableEnv.fromDataStream(stream, $("f1"));
// convert DataStream into Table with swapped fields
Table table = tableEnv.fromDataStream(stream, $("f1"), $("f0"));
// convert DataStream into Table with swapped fields and field names "myInt" and "myLong"
Table table = tableEnv.fromDataStream(stream, $("f1").as("myInt"), $("f0").as("myLong"));
Scala
// get a TableEnvironment
val tableEnv: StreamTableEnvironment = ... // see "Create a TableEnvironment" section
val stream: DataStream[(Long, Int)] = ...
// convert DataStream into Table with field "_2" only
val table: Table = tableEnv.fromDataStream(stream, $"_2")
// convert DataStream into Table with swapped fields
val table: Table = tableEnv.fromDataStream(stream, $"_2", $"_1")
// convert DataStream into Table with swapped fields and field names "myInt" and "myLong"
val table: Table = tableEnv.fromDataStream(stream, $"_2" as "myInt", $"_1" as "myLong")
Python
from pyflink.table.expressions import col
# get a TableEnvironment
t_env = ... # see "Create a TableEnvironment" section
stream = ... # type: DataStream of Types.Tuple([Types.LONG(), Types.INT()])
# convert DataStream into Table with field "f1" only
table = t_env.from_data_stream(stream, col('f1'))
# convert DataStream into Table with swapped fields
table = t_env.from_data_stream(stream, col('f1'), col('f0'))
# convert DataStream into Table with swapped fields and field names "my_int" and "my_long"
table = t_env.from_data_stream(stream, col('f1').alias('my_int'), col('f0').alias('my_long'))
원자적(Atomic) 타입
Flink는 원시 타입(Integer, Double, String) 또는 제네릭 타입(분석하고 분해할 수 없는 타입)을 원자적 타입으로 취급합니다. 원자적 타입의 DataStream은 단일 컬럼을 가진 Table로 변환됩니다. 컬럼의 타입은 원자적 타입에서 유추됩니다. 컬럼의 이름은 지정할 수 있습니다.
Java
StreamTableEnvironment tableEnv = ...;
DataStream<Long> stream = ...;
// Convert DataStream into Table with field name "myLong"
Table table = tableEnv.fromDataStream(stream, $("myLong"));
Scala
val tableEnv: StreamTableEnvironment = ???
val stream: DataStream[Long] = ...
// Convert DataStream into Table with default field name "f0"
val table: Table = tableEnv.fromDataStream(stream)
// Convert DataStream into Table with field name "myLong"
val table: Table = tableEnv.fromDataStream(stream, $"myLong")
Python
from pyflink.table.expressions import col
t_env = ...
stream = ... # types: DataStream of Types.Long()
# Convert DataStream into Table with default field name "f0"
table = t_env.from_data_stream(stream)
# Convert DataStream into Table with field name "my_long"
table = t_env.from_data_stream(stream, col('my_long'))
튜플(Scala, Java, Python)과 Case Class(Scala 전용)
Java
Flink는 Java용 자체 튜플 클래스를 제공합니다. Java 튜플 클래스의 DataStream은 테이블로 변환될 수 있습니다. 모든 필드에 이름을 제공하면 필드 이름을 바꿀 수 있습니다(위치 기반 매핑). 필드 이름이 지정되지 않으면 기본 필드 이름이 사용됩니다. 원래 필드 이름(Flink 튜플의 f0, f1, …)이 참조되면 API는 위치 기반이 아닌 이름 기반 매핑이라고 가정합니다. 이름 기반 매핑은 별칭(as)으로 필드 재정렬과 프로젝션을 허용합니다.
StreamTableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section
DataStream<Tuple2<Long, String>> stream = ...;
// convert DataStream into Table with renamed field names "myLong", "myString" (position-based)
Table table = tableEnv.fromDataStream(stream, $("myLong"), $("myString"));
// convert DataStream into Table with reordered fields "f1", "f0" (name-based)
Table table = tableEnv.fromDataStream(stream, $("f1"), $("f0"));
// convert DataStream into Table with projected field "f1" (name-based)
Table table = tableEnv.fromDataStream(stream, $("f1"));
// convert DataStream into Table with reordered and aliased fields "myString", "myLong" (name-based)
Table table = tableEnv.fromDataStream(stream, $("f1").as("myString"), $("f0").as("myLong"));
Scala
Flink는 Scala 내장 튜플을 지원합니다. Scala 내장 튜플의 DataStream은 테이블로 변환될 수 있습니다. 모든 필드에 이름을 제공하면 필드 이름을 바꿀 수 있습니다(위치 기반 매핑). 필드 이름이 지정되지 않으면 기본 필드 이름이 사용됩니다. 원래 필드 이름(Scala 튜플의 _1, _2, …)이 참조되면 API는 위치 기반이 아닌 이름 기반 매핑이라고 가정합니다. 이름 기반 매핑은 별칭(as)으로 필드 재정렬과 프로젝션을 허용합니다.
// get a TableEnvironment
val tableEnv: StreamTableEnvironment = ... // see "Create a TableEnvironment" section
val stream: DataStream[(Long, String)] = ...
// convert DataStream into Table with field names "myLong", "myString" (position-based)
val table: Table = tableEnv.fromDataStream(stream, $"myLong", $"myString")
// convert DataStream into Table with reordered fields "_2", "_1" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"_2", $"_1")
// convert DataStream into Table with projected field "_2" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"_2")
// convert DataStream into Table with reordered and aliased fields "myString", "myLong" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"_2" as "myString", $"_1" as "myLong")
// define case class
case class Person(name: String, age: Int)
val streamCC: DataStream[Person] = ...
// convert DataStream into Table with field names 'myName, 'myAge (position-based)
val table = tableEnv.fromDataStream(streamCC, $"myName", $"myAge")
// convert DataStream into Table with reordered and aliased fields "myAge", "myName" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"age" as "myAge", $"name" as "myName")
Python
Flink는 Python 내장 튜플을 지원합니다. 튜플의 DataStream은 테이블로 변환될 수 있습니다. 모든 필드에 이름을 제공하면 필드 이름을 바꿀 수 있습니다(위치 기반 매핑). 필드 이름이 지정되지 않으면 기본 필드 이름이 사용됩니다. 원래 필드 이름(f0, f1, …)이 참조되면 API는 위치 기반이 아닌 이름 기반 매핑이라고 가정합니다. 이름 기반 매핑은 별칭(alias)으로 필드 재정렬과 프로젝션을 허용합니다.
from pyflink.table.expressions import col
stream = ... # type: DataStream of Types.TUPLE([Types.LONG(), Types.STRING()])
# convert DataStream into Table with renamed field names "my_long", "my_string" (position-based)
table = t_env.from_data_stream(stream, col('my_long'), col('my_string'))
# convert DataStream into Table with reordered fields "f1", "f0" (name-based)
table = t_env.from_data_stream(stream, col('f1'), col('f0'))
# convert DataStream into Table with projected field "f1" (name-based)
table = t_env.from_data_stream(stream, col('f1'))
# convert DataStream into Table with reordered and aliased fields "my_string", "my_long" (name-based)
table = t_env.from_data_stream(stream, col('f1').alias('my_string'), col('f0').alias('my_long'))
POJO(Java와 Scala)
Flink는 POJO를 복합 타입으로 지원합니다. 무엇이 POJO를 결정하는지에 대한 규칙은 여기에 문서화되어 있습니다.
필드 이름을 지정하지 않고 POJO DataStream을 Table로 변환하면 원래 POJO 필드의 이름이 사용됩니다. 이름 매핑은 원래 이름을 요구하며 위치로는 수행할 수 없습니다. 필드는 별칭(as 키워드)으로 이름을 바꾸고, 재정렬하고, 프로젝션할 수 있습니다.
Java
StreamTableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section
// Person is a POJO with fields "name" and "age"
DataStream<Person> stream = ...;
// convert DataStream into Table with renamed fields "myAge", "myName" (name-based)
Table table = tableEnv.fromDataStream(stream, $("age").as("myAge"), $("name").as("myName"));
// convert DataStream into Table with projected field "name" (name-based)
Table table = tableEnv.fromDataStream(stream, $("name"));
// convert DataStream into Table with projected and renamed field "myName" (name-based)
Table table = tableEnv.fromDataStream(stream, $("name").as("myName"));
Scala
// get a TableEnvironment
val tableEnv: StreamTableEnvironment = ... // see "Create a TableEnvironment" section
// Person is a POJO with field names "name" and "age"
val stream: DataStream[Person] = ...
// convert DataStream into Table with renamed fields "myAge", "myName" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"age" as "myAge", $"name" as "myName")
// convert DataStream into Table with projected field "name" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"name")
// convert DataStream into Table with projected and renamed field "myName" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"name" as "myName")
Python
현재 PyFlink에서는 커스텀 POJO 클래스가 지원되지 않습니다.
Row
Row 데이터 타입은 무제한의 필드와 null 값을 가진 필드를 지원합니다. 필드 이름은 RowTypeInfo를 통해 또는 Row DataStream을 Table로 변환할 때 지정할 수 있습니다. 행 타입은 위치와 이름으로 필드를 매핑하는 것을 모두 지원합니다. 모든 필드에 이름을 제공하면 필드 이름을 바꿀 수 있고(위치 기반 매핑), 프로젝션/정렬/이름 변경을 위해 개별적으로 선택할 수도 있습니다(이름 기반 매핑).
Java
StreamTableEnvironment tableEnv = ...;
// DataStream of Row with two fields "name" and "age" specified in `RowTypeInfo`
DataStream<Row> stream = ...;
// Convert DataStream into Table with renamed field names "myName", "myAge" (position-based)
Table table = tableEnv.fromDataStream(stream, $("myName"), $("myAge"));
// Convert DataStream into Table with renamed fields "myName", "myAge" (name-based)
Table table = tableEnv.fromDataStream(stream, $("name").as("myName"), $("age").as("myAge"));
// Convert DataStream into Table with projected field "name" (name-based)
Table table = tableEnv.fromDataStream(stream, $("name"));
// Convert DataStream into Table with projected and renamed field "myName" (name-based)
Table table = tableEnv.fromDataStream(stream, $("name").as("myName"));
Scala
val tableEnv: StreamTableEnvironment = ???
// DataStream of Row with two fields "name" and "age" specified in `RowTypeInfo`
val stream: DataStream[Row] = ...
// Convert DataStream into Table with renamed field names "myName", "myAge" (position-based)
val table: Table = tableEnv.fromDataStream(stream, $"myName", $"myAge")
// Convert DataStream into Table with renamed fields "myName", "myAge" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"name" as "myName", $"age" as "myAge")
// Convert DataStream into Table with projected field "name" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"name")
// Convert DataStream into Table with projected and renamed field "myName" (name-based)
val table: Table = tableEnv.fromDataStream(stream, $"name" as "myName")
Python
from pyflink.table.expressions import col
t_env = ...;
# DataStream of Row with two fields "name" and "age" specified in `RowTypeInfo`
stream = ...
# Convert DataStream into Table with renamed field names "my_name", "my_age" (position-based)
table = t_env.from_data_stream(stream, col('my_name'), col('my_age'))
# Convert DataStream into Table with renamed fields "my_name", "my_age" (name-based)
table = t_env.from_data_stream(stream, col('name').alias('my_name'), col('age').alias('my_age'))
# Convert DataStream into Table with projected field "name" (name-based)
table = t_env.from_data_stream(stream, col('name'))
# Convert DataStream into Table with projected and renamed field "my_name" (name-based)
table = t_env.from_data_stream(stream, col('name').alias("my_name"))