Table API 배우기
Table API 배우기
이 교육의 초점은 스트리밍 분석 애플리케이션 작성을 시작할 수 있을 만큼 Table API를 폭넓게 다루는 것입니다.
Table API는 Flink의 선언적(declarative), 관계형(relational) API입니다. 무엇을 계산할지 기술하면 Flink가 어떻게 효율적으로 계산할지 알아냅니다. 동일한 쿼리는 수정 없이 배치와 스트리밍 데이터 모두에서 동작하며, Flink의 최적화기가 효율적인 실행 계획을 자동으로 선택합니다.
출처: 문서
본문
테이블로 표현할 수 있는 것
Table API는 정의된 스키마가 있는 구조화된 데이터로 동작합니다. 모든 테이블은 특정 데이터 타입을 가진 명명된 열을 가집니다.
Flink는 테이블 열에 대해 풍부한 데이터 타입 집합을 지원합니다.
- 기본 타입:
STRING,INT,BIGINT,DOUBLE,BOOLEAN,TIMESTAMP - 복합 타입:
ARRAY,MAP,ROW(중첩 구조용) - 구조화 타입: 명명된 필드를 가진 Java POJO
- 특수 타입:
RAW(UDF의 불투명 바이트 데이터용),INTERVAL,NULL
테이블은 외부 시스템(Kafka, 파일, 데이터베이스), DataStream, 또는 값을 사용해 인라인으로 만들 수 있습니다.
완전한 예시
이 예시는 사람에 대한 레코드 테이블을 입력으로 받아 성인만 포함하도록 필터링합니다.
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.TableEnvironment;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.types.Row;
import static org.apache.flink.table.api.Expressions.*;
public class Example {
public static void main(String[] args) {
TableEnvironment tableEnv = TableEnvironment.create(
EnvironmentSettings.inStreamingMode());
Table people = tableEnv.fromValues(
DataTypes.ROW(
DataTypes.FIELD("name", DataTypes.STRING()),
DataTypes.FIELD("age", DataTypes.INT())
),
Row.of("Alice", 35),
Row.of("Bob", 35),
Row.of("Charlie", 2)
);
Table adults = people
.filter($("age").isGreaterOrEqual(18))
.select($("name"), $("age"));
adults.execute().print();
}
}
출력은 대략 다음과 같습니다.
+----+--------------------------------+-------------+
| op | name | age |
+----+--------------------------------+-------------+
| +I | Alice | 35 |
| +I | Bob | 35 |
+----+--------------------------------+-------------+
op 열은 연산 유형(+I는 INSERT를 의미)을 보여줍니다. +I 플래그는 이 행들이 결과 테이블에 삽입되고 있음을 나타냅니다. 아래 changelog 섹션에서 이 플래그를 더 자세히 설명합니다.
Table Environment
모든 Table API 프로그램에는 TableEnvironment가 필요합니다. 이것은 다음의 중심 진입점입니다.
- 테이블 생성 및 등록
- 쿼리 실행
- 사용자 정의 함수 등록
- Table API와 DataStream API 간 변환
// For streaming applications
TableEnvironment tableEnv = TableEnvironment.create(
EnvironmentSettings.inStreamingMode());
// For batch applications
TableEnvironment tableEnv = TableEnvironment.create(
EnvironmentSettings.inBatchMode());
테이블에서 execute()를 호출하면 Flink가 테이블 프로그램을 최적화된 작업 그래프로 컴파일하고 실행을 위해 제출합니다.
기본 연산
Select 및 Projection
select()를 사용해 포함할 열을 선택하고 새 계산 열을 만들 수 있습니다. $() 함수는 이름으로 열을 참조합니다.
Table result = orders
.select($("product"), $("amount"), $("price").times($("amount")).as("total"));
필터링
filter() 또는 where()(둘은 동일)를 사용해 조건과 일치하는 행만 유지합니다.
Table filtered = orders
.filter($("amount").isGreater(0))
.filter($("status").isNotEqual("cancelled"));
열 추가 및 수정
addColumns()로 새 계산 열을 추가하고, renameColumns()로 기존 열 이름을 바꾸고, dropColumns()로 열을 제거합니다.
Table result = orders
.addColumns($("price").times($("amount")).as("total"))
.renameColumns($("product").as("item"))
.dropColumns($("internal_id"));
내장 함수
Table API는 일반적인 작업을 위한 많은 내장 함수를 포함합니다.
import static org.apache.flink.table.api.Expressions.*;
Table result = products
.select(
$("name").upperCase(), // String functions
$("price").round(2), // Math functions
$("created_at").extract(TimeIntervalUnit.HOUR) // Temporal functions
);
내장 함수의 전체 목록은 Built-in Functions 참조를 보세요.
집계 및 그룹화
groupBy()를 집계 함수와 함께 사용해 요약 통계를 계산합니다.
Table counts = orders
.groupBy($("product"))
.select($("product"), $("amount").sum().as("total_sold"), $("product").count().as("order_count"));
일반적인 집계 함수에는 sum(), count(), avg(), min(), max()가 있습니다.
스트리밍 모드에서 집계는 갱신되는 결과를 생성합니다. 그룹에 새 행이 도착할 때마다 Flink가 집계를 갱신하고 해당 그룹에 대한 새 결과를 방출합니다.
Changelog 이해하기
Table API를 스트리밍에 사용할 때 이해해야 할 가장 중요한 개념 중 하나는 *stream-table duality(스트림-테이블 이중성)*입니다. 스트림은 테이블로 볼 수 있고(시간이 지남에 따라 행이 삽입됨), 테이블은 changelog 스트림으로 볼 수 있습니다(테이블의 각 변경은 이벤트).
Append-Only 테이블
filter()와 select() 같은 단순 연산은 append-only 테이블을 생성합니다. 새 행은 삽입되지만 기존 행은 수정되거나 삭제되지 않습니다. 이들은 +I(INSERT) 플래그만 생성합니다.
갱신되는 테이블
집계와 다른 상태 기반 연산은 갱신되는(updating) 테이블을 생성합니다. 그룹의 집계 결과가 변경되면 Flink는 이전 값을 철회(retract)하고 새 값을 삽입해야 합니다.
볼 수 있는 changelog 플래그는 다음과 같습니다.
| Flag | Meaning |
|---|---|
+I |
INSERT: 새 행이 추가됨 |
-D |
DELETE: 기존 행이 제거됨 |
-U |
UPDATE_BEFORE: 갱신 전의 이전 값 (철회) |
+U |
UPDATE_AFTER: 갱신 후의 새 값 |
예를 들어 제품별 주문을 세고 "widget"에 대한 주문 3개가 도착한다면:
+I [widget, 1] -- First order: count is 1
-U [widget, 1] -- Retract previous count
+U [widget, 2] -- Second order: count is now 2
-U [widget, 2] -- Retract previous count
+U [widget, 3] -- Third order: count is now 3
이것은 싱크에 연결할 때 중요합니다. Append-only 싱크(예: 파일)는 +I 연산만 수락할 수 있습니다. 갱신 싱크(예: 데이터베이스 또는 Kafka upsert 토픽)는 모든 연산을 처리할 수 있습니다.
윈도우 집계
무한 집계는 상태를 영원히 유지하므로 많은 스트리밍 사용 사례에 실용적이지 않습니다. 윈도우는 스트림의 유계 부분에 대해 집계할 수 있게 해줍니다.
시간당 주문을 세기 위해 텀블링 윈도우를 사용하는 예시입니다.
import static org.apache.flink.table.api.Expressions.*;
Table hourlyStats = orders
.window(Tumble.over(lit(1).hours()).on($("order_time")).as("w"))
.groupBy($("product"), $("w"))
.select(
$("product"),
$("w").start().as("window_start"),
$("w").end().as("window_end"),
$("amount").sum().as("total_sold")
);
윈도우 집계는 각 윈도우가 한 번 최종 결과를 생성하고 결코 갱신하지 않으므로 append-only 결과를 생성합니다. 따라서 갱신을 지원하지 않는 싱크에 이상적입니다.
윈도우에 대한 자세한 내용은 Group Aggregation 문서를 참조하세요.
사용자 정의 함수
내장 함수가 요구를 충족하지 못할 때 사용자 정의 함수(UDF)를 만들 수 있습니다.
스칼라 함수
ScalarFunction은 한 행을 받아 하나의 값을 생성합니다.
import org.apache.flink.table.functions.ScalarFunction;
public class HashFunction extends ScalarFunction {
public String eval(String input) {
return Integer.toHexString(input.hashCode());
}
}
// Register and use
tableEnv.createTemporaryFunction("hash", HashFunction.class);
Table result = orders.select($("id"), call("hash", $("customer_name")).as("hashed_name"));
테이블 함수
TableFunction은 한 행을 받아 0개 이상의 행을 생성합니다. 적용하려면 joinLateral()을 사용하세요.
import org.apache.flink.table.functions.TableFunction;
import org.apache.flink.types.Row;
public class SplitFunction extends TableFunction<Row> {
public void eval(String str) {
for (String s : str.split(",")) {
collect(Row.of(s.trim()));
}
}
}
// Use with joinLateral
tableEnv.createTemporaryFunction("split", SplitFunction.class);
Table result = orders.joinLateral(call("split", $("tags")).as("tag"));
UDF에 대한 자세한 내용은 User-defined Functions 문서를 참조하세요.
Process Table Functions
선언적 연산으로 충분하지 않을 때 Process Table Functions(PTF)가 처리에 대한 완전한 제어를 제공합니다. PTF는 DataStream API의 ProcessFunction이 작동하는 방식과 유사하게 복잡한 로직에 대한 Table API의 "탈출구(escape hatch)"입니다.
PTF는 다음을 할 수 있습니다.
- 여러 행에 걸쳐 **상태(state)**에 접근하고 관리
- 특정 시간에 처리를 트리거할 타이머(timer) 등록
- 사용자 지정 파티션별 로직으로 파티션된(partitioned) 데이터 처리
키별 발생 횟수를 세는 간단한 상태 기반 PTF입니다.
import org.apache.flink.table.annotation.*;
import org.apache.flink.table.functions.ProcessTableFunction;
import org.apache.flink.types.Row;
public class CountingFunction extends ProcessTableFunction<String> {
public static class CountState {
public long count = 0L;
}
public void eval(
@StateHint CountState state,
@ArgumentHint(ArgumentTrait.SET_SEMANTIC_TABLE) Row input) {
state.count++;
collect("Seen " + state.count + " rows for this partition");
}
}
// Use with partitionBy for per-key state
Table result = events
.partitionBy($("user_id"))
.process(CountingFunction.class);
PTF는 선언적 Table API와 저수준 스트림 처리 사이의 간극을 메웁니다. 표준 연산으로 표현할 수 없는 상태, 타이머 또는 이벤트별 처리 로직에 대한 세밀한 제어가 필요할 때 사용하세요.
완전한 PTF 가이드는 Process Table Functions를 참조하세요.
실습
이 시점에서 간단한 Table API 애플리케이션을 코딩하고 실행하기 시작할 수 있을 만큼 알게 되었습니다. flink-training-repo를 클론하고 README의 지침을 따른 후 Table API 연습을 시도해 보세요.
추가 읽기
- Table API 연산 - 모든 연산에 대한 완전한 참조
- 내장 함수 - 모든 사용 가능한 함수
- 사용자 정의 함수 - 사용자 정의 함수 만들기
- Process Table Functions - 고급 상태 기반 처리
- SQL & Table API 커넥터 - 외부 시스템 읽기/쓰기
- 공통 개념 - Table API와 SQL 사이의 공유 개념