Table API
Table API
Table API는 스트림과 배치 처리를 위한 통합된 관계형 API예요. Table API 쿼리는 수정 없이 배치 또는 스트리밍 입력에서 실행될 수 있어요. Table API는 SQL 언어의 상위 집합이며 Apache Flink와 함께 작업하기 위해 특별히 설계됐어요. Table API는 Scala, Java, Python을 위한 언어 통합(language-integrated) API예요. SQL에서 흔한 것처럼 쿼리를 String 값으로 지정하는 대신, Table API 쿼리는 자동 완성과 문법 검증 같은 IDE 지원과 함께 Java, Scala 또는 Python에서 언어 임베디드 스타일로 정의돼요.
출처: 문서
본문
Table API는 Flink의 SQL 통합과 많은 개념과 API 부분을 공유해요. 테이블을 등록하거나 Table 객체를 만드는 방법을 배우려면 Common Concepts & API를 참조하세요. Streaming Concepts 페이지는 동적 테이블과 시간 속성 같은 스트리밍 특정 개념을 다뤄요.
다음 예시들은 (a, b, c, rowtime) 속성을 가진 Orders라는 등록된 테이블을 가정해요. rowtime 필드는 스트리밍에서 논리적 시간 속성(time attribute)이거나, 배치에서 일반적인 타임스탬프 필드예요.
개요 & 예시 (Overview & Examples)
Table API는 Scala, Java, Python에서 사용 가능해요. Scala Table API는 Scala 표현식을 활용하고, Java Table API는 Expression DSL과 동등한 표현식으로 파싱되고 변환되는 문자열을 모두 지원하며, Python Table API는 현재 동등한 표현식으로 파싱되고 변환되는 문자열만 지원해요.
다음 예시는 Scala, Java, Python Table API의 차이를 보여줘요. 테이블 프로그램은 배치 환경에서 실행돼요. Orders 테이블을 스캔하고 필드 a로 그룹화하며, 그룹별 결과 행 수를 세어요.
Java Table API는 org.apache.flink.table.api.java.*를 임포트해 활성화돼요. 다음 예시는 Java Table API 프로그램이 어떻게 구성되고 표현식이 문자열로 어떻게 지정되는지 보여줘요. Expression DSL을 위해서는 org.apache.flink.table.api.Expressions.*를 정적으로 임포트하는 것도 필요해요.
import org.apache.flink.table.api.*;
import static org.apache.flink.table.api.Expressions.*;
EnvironmentSettings settings = EnvironmentSettings
.newInstance()
.inStreamingMode()
.build();
TableEnvironment tEnv = TableEnvironment.create(settings);
// register Orders table in table environment
// ...
// specify table program
Table orders = tEnv.from("Orders"); // schema (a, b, c, rowtime)
Table counts = orders
.groupBy($("a"))
.select($("a"), $("b").count().as("cnt"));
// print
counts.execute().print();
Scala Table API는 org.apache.flink.table.api._, org.apache.flink.api.scala._, org.apache.flink.table.api.bridge.scala._(DataStream과의 브리징용)를 임포트해 활성화돼요.
import org.apache.flink.api.scala._
import org.apache.flink.table.api._
import org.apache.flink.table.api.bridge.scala._
// environment configuration
val settings = EnvironmentSettings
.newInstance()
.inStreamingMode()
.build()
val tEnv = TableEnvironment.create(settings)
// register Orders table in table environment
// ...
// specify table program
val orders = tEnv.from("Orders") // schema (a, b, c, rowtime)
val result = orders
.groupBy($"a")
.select($"a", $"b".count as "cnt")
.execute()
.print()
from pyflink.table import *
from pyflink.table.expressions import col
# environment configuration
t_env = TableEnvironment.create(
environment_settings=EnvironmentSettings.in_batch_mode())
# register Orders table and Result table sink in table environment
source_data_path = "/path/to/source/directory/"
result_data_path = "/path/to/result/directory/"
source_ddl = f"""
create table Orders(
a VARCHAR,
b BIGINT,
c BIGINT,
rowtime TIMESTAMP(3),
WATERMARK FOR rowtime AS rowtime - INTERVAL '1' SECOND
) with (
'connector' = 'filesystem',
'format' = 'csv',
'path' = '{source_data_path}'
)
"""
t_env.execute_sql(source_ddl)
sink_ddl = f"""
create table `Result`(
a VARCHAR,
cnt BIGINT
) with (
'connector' = 'filesystem',
'format' = 'csv',
'path' = '{result_data_path}'
)
"""
t_env.execute_sql(sink_ddl)
# specify table program
orders = t_env.from_path("Orders") # schema (a, b, c, rowtime)
orders.group_by(col("a")).select(col("a"), col("b").count.alias('cnt')).execute_insert("result").wait()
다음 예시는 더 복잡한 Table API 프로그램을 보여줘요. Orders 테이블을 다시 스캔해요. null 값을 필터링하고, String 타입의 a 필드를 정규화하며, 각 시간과 제품 a에 대해 평균 청구 금액 b를 계산해요.
// environment configuration
// ...
// specify table program
Table orders = tEnv.from("Orders"); // schema (a, b, c, rowtime)
Table result = orders
.filter(
and(
$("a").isNotNull(),
$("b").isNotNull(),
$("c").isNotNull()
))
.select($("a").lowerCase().as("a"), $("b"), $("rowtime"))
.window(Tumble.over(lit(1).hours()).on($("rowtime")).as("hourlyWindow"))
.groupBy($("hourlyWindow"), $("a"))
.select($("a"), $("hourlyWindow").end().as("hour"), $("b").avg().as("avgBillingAmount"));
연산 (Operations)
Table API는 다양한 연산을 지원해요. 주요 연산 카테고리는 다음과 같아요:
- 스캔, 프로젝션, 필터 (Scan, Projection, Filter):
from,select,filter등. - 집계 (Aggregations):
groupBy,groupBy창, distinct 집계, over 창 집계 등. 예를 들어 GROUP BY, 시간 윈도우, over 윈도우에 대한 distinct 집계를 지원해요. - 조인 (Joins): inner/outer 조인, 시간 간격 조인, temporal 조인, lookup 조인, array expansion 등.
- set 연산 (Set Operations): union, intersect, minus 등.
- 정렬, 오프셋, 페치 (OrderBy, Offset & Fetch):
orderBy,offset,fetch등. - 윈도우 (Windows): 텀블링, 슬라이딩, 세션 윈도우와 over 윈도우. 예시:
- 텀블링 이벤트 타임/프로세싱 타임/행-카운트 윈도우
- 슬라이딩 이벤트 타임/프로세싱 타임/행-카운트 윈도우
- 세션 윈도우
- 무한(Unbounded)/유한(Bounded) over 윈도우
- 사용자 정의 함수 (User-defined Functions): UDF, UDTF(테이블 함수), UDAF(집계 함수). Python general/vectorized 스칼라 함수와 집계 함수를 매핑/플랫맵/조인/집계에 사용할 수 있어요.
- 데이터 타입 (Data Types): Table API는 풍부한 데이터 타입 집합을 지원하며, 자세한 내용은 Data Types 문서를 참고하세요.
각 연산의 자세한 문법과 예시는 Operations 섹션을 참고하세요. 이 페이지에는 Java, Scala, Python 예시가 포함되어 있으며, 윈도우 집계, 조인, distinct 집계, 사용자 정의 함수 사용법 등을 설명해요.
예를 들어, 텀블링 이벤트 타임 윈도우는 다음과 같이 정의할 수 있어요:
tableEnv.from("Orders")
.window(Tumble.over(lit(1).hour()).on($("rowtime")).as("w"))
.groupBy($("w"))
.select($("b").sum().as("sum"));
Python Table API는 아직 모든 연산을 지원하지 않아요. 지원되지 않는 연산은 "Not yet supported in Python Table API"로 표시돼요.
데이터 타입 (Data Types)
자세한 내용은 Flink SQL 데이터 타입을 참고하세요.