Process Table Functions
Process Table Functions (PTFs)
Process Table Functions(PTFs) 는 Flink SQL 및 Table API 에서 가장 강력한 함수 종류입니다. 내장 연산만큼 풍부한 기능을 가진 사용자 정의 operator 를 구현할 수 있게 합니다. PTF 는 (파티셔닝된) 테이블을 받아 새 테이블을 생성할 수 있습니다. Flink 의 관리 상태, 이벤트 시간 및 타이머 서비스, 기반 테이블 changelog 에 접근할 수 있습니다.
출처: 문서
본문
Process Table Functions(PTFs) 는 Flink SQL 및 Table API 에서 가장 강력한 함수 종류입니다. 내장 연산만큼 풍부한 기능을 가진 사용자 정의 operator 를 구현할 수 있게 합니다. PTF 는 (파티셔닝된) 테이블을 받아 새 테이블을 생성할 수 있습니다. Flink 의 관리 상태, 이벤트 시간 및 타이머 서비스, 기반 테이블 changelog 에 접근할 수 있습니다.
개념적으로 PTF 는 그 자체로 사용자 정의 함수 이며 다른 모든 사용자 정의 함수의 초집합(superset) 입니다. 0개, 1개 또는 여러 개의 테이블을 0개, 1개 또는 여러 개의 행(또는 구조적 타입)에 매핑합니다. 스칼라 인수가 지원됩니다. 상태 저장적인 특성 덕분에 집계 동작을 구현하는 것도 가능합니다.
PTF 는 다음 작업을 가능하게 합니다:
- 테이블의 각 행에 변환을 적용.
- 테이블을 논리적으로 서로 다른 집합으로 파티셔닝하고 집합별로 변환을 적용.
- 반복 접근을 위해 본 이벤트를 저장.
- 대기, 동기화 또는 타임아웃을 가능하게 하는 더 나중 시점에 처리를 계속.
- 복잡한 상태 머신 또는 규칙 기반 조건 로직을 사용해 이벤트를 버퍼링하고 집계.
다형성 테이블 함수 (Polymorphic Table Functions)
PTF 쿼리 문법과 의미론은 SQL:2016 의 Polymorphic Table Functions 에서 파생되었습니다. 다형성 테이블 함수의 기대 동작 및 SQL 언어 내 통합에 대한 자세한 정보는 ISO/IEC 19075-7:2021 (Part 7) 에서 찾을 수 있습니다. 공개적으로 이용 가능한 요약은 관련 SIGMOD 논문의 Section 3 에 제공됩니다.
둘 다 같은 약어(PTF)를 공유하지만, Flink 의 process table functions 은 상태 관리, 시간, 타이머 서비스 같은 Flink 고유 기능을 통합함으로써 다형성 테이블 함수를 향상시킵니다. 행 또는 집합 의미론이 있는 테이블 인수, 디스크립터 인수, 가상 프로세서(virtual processors) 와 관련된 처리 개념을 포함한 호출 특성은 SQL 표준과 정렬됩니다.
동기 부여 예제 (Motivating Examples)
다음 예제는 PTF 가 테이블을 받아 변환할 수 있는 방법을 보여줍니다. @ArgumentHint 는 함수가 단순한 스칼라 행 값이 아닌 테이블을 인수로 받는다고 지정합니다. 두 예제 모두 입력 테이블의 각 행에 대해 eval() 메서드가 호출됩니다. 또한 @ArgumentHint 는 함수가 테이블을 처리할 수 있다는 것뿐만 아니라 행 또는 집합 의미론이든 간에 함수가 테이블을 어떻게 해석하는지도 정의합니다.
Greeting
들어오는 각 고객에게 인사를 추가하는 예제입니다.
Java:
import org.apache.flink.table.annotation.*;
import org.apache.flink.table.api.*;
import org.apache.flink.table.functions.ProcessTableFunction;
import org.apache.flink.types.Row;
import static org.apache.flink.table.api.Expressions.*;
// A PTF that takes a table argument, conceptually viewing the table as a row.
// The result is never stateful and derived purely based on the current row.
public static class Greeting extends ProcessTableFunction<String> {
public void eval(@ArgumentHint(ArgumentTrait.ROW_SEMANTIC_TABLE) Row input) {
collect("Hello " + input.getFieldAs("name") + "!");
}
}
TableEnvironment env = TableEnvironment.create(...);
// Call the PTF with row semantics "inline" (without registration) in Table API
env.fromValues("Bob", "Alice", "Bob")
.as("name")
.process(Greeting.class)
.execute()
.print();
// For SQL, register the PTF upfront
env.executeSql("CREATE VIEW Names(name) AS VALUES ('Bob'), ('Alice'), ('Bob')");
env.createFunction("Greeting", Greeting.class);
// Call the PTF with row semantics in SQL
env.executeSql("SELECT * FROM Greeting(TABLE Names)").print();
Table API 와 SQL 두 결과 모두 다음과 유사합니다:
+----+--------------------------------+
| op | EXPR$0 |
+----+--------------------------------+
| +I | Hello Bob! |
| +I | Hello Alice! |
| +I | Hello Bob! |
+----+--------------------------------+
Greeting with Memory
들어오는 각 고객에 대해, 그 고객이 이전에 인사받았는지 여부를 고려해 인사를 추가하는 예제입니다.
Java:
import org.apache.flink.table.annotation.*;
import org.apache.flink.table.api.*;
import org.apache.flink.table.functions.ProcessTableFunction;
import org.apache.flink.types.Row;
import static org.apache.flink.table.api.Expressions.*;
// A PTF that takes a table argument, conceptually viewing the table as a set.
// The result can be stateful and derived based on the current row and/or
// previous rows in the set.
// The call's partitioning defines the size of the set.
public static class GreetingWithMemory extends ProcessTableFunction<String> {
public static class CountState {
public long counter = 0L;
}
public void eval(@StateHint CountState state, @ArgumentHint(ArgumentTrait.SET_SEMANTIC_TABLE) Row input) {
state.counter++;
collect("Hello " + input.getFieldAs("name") + ", your " + state.counter + " time?");
}
}
TableEnvironment env = TableEnvironment.create(...);
// Call the PTF with set semantics "inline" (without registration) in Table API
env.fromValues("Bob", "Alice", "Bob")
.as("name")
.partitionBy($("name"))
.process(GreetingWithMemory.class)
.execute()
.print();
// For SQL, register the PTF upfront
env.executeSql("CREATE VIEW Names(name) AS VALUES ('Bob'), ('Alice'), ('Bob')");
env.createFunction("GreetingWithMemory", GreetingWithMemory.class);
// Call the PTF with set semantics in SQL
env.executeSql("SELECT * FROM GreetingWithMemory(TABLE Names PARTITION BY name)").print();
Table API 와 SQL 두 결과 모두 다음과 유사합니다:
+----+--------------------------------+--------------------------------+
| op | name | EXPR$0 |
+----+--------------------------------+--------------------------------+
| +I | Bob | Hello Bob, your 1 time? |
| +I | Alice | Hello Alice, your 1 time? |
| +I | Bob | Hello Bob, your 2 time? |
+----+--------------------------------+--------------------------------+
Greeting with Follow Up
들어오는 각 고객에게 인사를 추가하고 일정 시간 후 후속 알림을 보내는 예제입니다. 이 예제는 시간과 이름 없는(unnamed) 타이머의 사용을 보여줍니다.
Java:
import org.apache.flink.table.annotation.*;
import org.apache.flink.table.api.*;
import org.apache.flink.table.functions.ProcessTableFunction;
import org.apache.flink.types.Row;
import static org.apache.flink.table.api.Expressions.*;
// A stateful PTF with event-time timers.
// Every incoming row not only greets the customer but also registers a timer to follow
// up on the customer after 1 minute.
public static class GreetingWithFollowUp extends ProcessTableFunction<String> {
public static class StayState {
public String name;
public long counter = 0L;
}
public void eval(
Context ctx,
@StateHint StayState state,
@ArgumentHint(ArgumentTrait.SET_SEMANTIC_TABLE) Row input
) {
state.name = input.getFieldAs("name");
state.counter++;
collect("Hello " + state.name + ", your " + state.counter + " time?");
TimeContext<Instant> timeCtx = ctx.timeContext(Instant.class);
timeCtx.registerOnTime(timeCtx.time().plus(Duration.ofMinutes(1)));
}
public void onTimer(StayState state) {
collect("Hello " + state.name + ", I hope you enjoyed your stay! "
+ "Please let us know if there's anything we could have done better.");
}
}
TableEnvironment env = TableEnvironment.create(...);
// Create a watermarked table with events from Bob and Alice
env.createTable(
"Names",
TableDescriptor.forConnector("datagen")
.schema(
Schema.newBuilder()
.column("ts", "TIMESTAMP_LTZ(3)")
.column("random", "INT")
.columnByExpression("name", "IF(random % 2 = 0, 'Bob', 'Alice')")
.watermark("ts", "ts - INTERVAL '1' SECOND")
.build())
.option("rows-per-second", "1")
.build());
// Call the PTF in Table API and pass the watermarked time column
env.from("Names")
.partitionBy($("name"))
.process(GreetingWithCatchingUp.class, descriptor("ts").asArgument("on_time"))
.execute()
.print();
// For SQL, register the PTF upfront
env.createFunction("GreetingWithFollowUp", GreetingWithFollowUp.class);
// Call the PTF in SQL and pass the watermarked time column
env.executeSql("SELECT * FROM GreetingWithFollowUp(input => TABLE Names PARTITION BY name, on_time => DESCRIPTOR(ts))").print();
Table API 와 SQL 두 결과 모두 다음과 유사합니다:
+----+--------------------------------+--------------------------------+-------------------------+
| op | name | EXPR$0 | rowtime |
+----+--------------------------------+--------------------------------+-------------------------+
| +I | Bob | Hello Bob, your 1 time? | 2025-03-27 07:38:31.134 |
| +I | Bob | Hello Bob, your 2 time? | 2025-03-27 07:38:31.134 |
| +I | Alice | Hello Alice, your 1 time? | 2025-03-27 07:38:31.134 |
| +I | Alice | Hello Alice, your 2 time? | 2025-03-27 07:38:31.134 |
...
| +I | Bob | Hello Bob, I hope you enjoy... | 2025-03-27 07:39:31.134 |
| +I | Alice | Hello Alice, I hope you enj... | 2025-03-27 07:39:31.134 |
| +I | Bob | Hello Bob, I hope you enjoy... | 2025-03-27 07:39:31.135 |
테이블 의미론과 가상 프로세서 (Table Semantics and Virtual Processors)
PTF 는 테이블을 인수로 소비함으로써 새 테이블을 생성할 수 있습니다. 확장성을 위해 입력 테이블은 소위 "가상 프로세서(virtual processors)" 간에 분산됩니다. SQL 표준이 정의하는 가상 프로세서는 PTF 인스턴스를 실행하며 전체 테이블의 일부에만 접근할 수 있습니다. 인수 선언이 그 일부의 크기와 데이터의 동일 위치(co-location) 를 결정합니다. 개념적으로 테이블은 "행별로"(즉 행 의미론) 또는 "집합별로"(즉 집합 의미론) 처리될 수 있습니다.
참고: 이 맥락에서 접근한다는 것은 개별 행이 가상 프로세서를 통해 스트리밍됨을 의미합니다. 반복 접근을 위해 과거 이벤트를 상태로 저장하는 것은 PTF 의 책임입니다.
행 의미론을 가진 테이블 인수
행 의미론을 가진 테이블을 받는 PTF 는 행 사이에 상관관계가 없고 각 행이 독립적으로 처리될 수 있다고 가정합니다. 프레임워크는 행을 가상 프로세서 간에 어떻게 분산할지 자유로우며, 각 가상 프로세서는 현재 처리 중인 행에만 접근할 수 있습니다.
집합 의미론을 가진 테이블 인수
집합 의미론을 가진 테이블을 받는 PTF 는 행 사이에 상관관계가 있다고 가정합니다. 함수를 호출할 때 PARTITION BY 절이 상관관계를 위한 열을 정의합니다. 프레임워크는 같은 집합에 속하는 모든 행이 동일 위치에 있도록 보장합니다. PTF 인스턴스는 같은 집합에 속하는 모든 행에 접근할 수 있습니다. 다시 말해: 가상 프로세서는 키 컨텍스트로 범위가 지정됩니다.
키를 제공하지 않는 것도 가능합니다(인수가 ArgumentTrait.OPTIONAL_PARTITION_BY 로 선언된 경우). 이 경우 하나의 가상 프로세서만 전체 테이블을 처리하므로 확장성 이점을 잃게 됩니다.
호출 문법 (Call Syntax)
PTF 를 호출할 때 시스템은 사용자 정의 입력 인수와 함께 상태 및 시간 관리를 위한 암시적 인수를 자동으로 추가합니다.
TableFilter PTF 가 다음과 같이 구현되었다고 가정합니다:
Java:
@DataTypeHint("ROW<threshold INT, score BIGINT>")
public static class TableFilter extends ProcessTableFunction<Row> {
public void eval(@ArgumentHint(ArgumentTrait.ROW_SEMANTIC_TABLE) Row input, int threshold) {
long score = input.getFieldAs("score");
if (score > threshold) {
collect(Row.of(threshold, score));
}
}
}
유효한 호출 서명은 다음과 같습니다:
TableFilter(input => {TABLE, ROW SEMANTIC TABLE}, threshold => INT NOT NULL, on_time => DESCRIPTOR, uid => STRING)
on_time 과 uid 는 모두 기본적으로 선택 사항입니다. 시간 의미론이 필요하면 on_time 이 필요합니다. 상태 저장변환의 경우 uid 를 제공해야 합니다.
SQL 과 Table API 사용자는 모두 위치 기반 또는 이름 기반 문법으로 함수를 호출할 수 있습니다. 더 나은 가독성과 향후 함수 발전을 위해 이름 기반 문법을 권장합니다. 이는 선택 인수에 대한 더 나은 지원을 제공하고 특정 인수 순서를 유지할 필요를 없앱니다.
SQL:
-- Position-based
SELECT * FROM TableFilter(TABLE t, 100)
SELECT * FROM TableFilter(TABLE t, 100, DEFAULT, 'my-ptf')
SELECT * FROM TableFilter(TABLE t, 100, DEFAULT, DEFAULT)
-- Name-based
SELECT * FROM TableFilter(input => TABLE t, threshold => 100)
SELECT * FROM TableFilter(input => TABLE t, uid => 'my-ptf')
Table API:
// Position-based
env.from("t").process(TableFilter.class, 100)
env.from("t").process(TableFilter.class, 100, null, "my-ptf")
env.from("t").process(TableFilter.class, 100, null, null)
// Name-based
env.from("t").process(TableFilter.class, lit(100).asArgument("threshold"))
env.from("t").process(TableFilter.class, lit(100).asArgument("threshold"), lit("my-ptf").asArgument("uid"))
// Fully name-based (including the table itself)
env.fromCall(
TableFilter.class,
env.from("t").asArgument("input"),
lit(100).asArgument("threshold"),
lit("my-ptf").asArgument("uid"))
함수 연결 (Function Chaining)
여러 PTF 를 연속으로 호출할 수 있습니다.
SQL 에서는 Divide and Conquer 전략을 적용해 문제를 더 작고 관리하기 쉬운 부분으로 나눔으로써 코드 가독성을 높이기 위해 공통 테이블 표현식(즉 WITH) 을 사용하는 것을 권장합니다:
WITH
ptf1 AS (
SELECT * FROM f1(input => TABLE t PARTITION BY name, on_time => DESCRIPTOR(ts))
),
ptf2 AS (
SELECT * FROM f2(input => TABLE ptf1 PARTITION BY name, on_time => DESCRIPTOR(rowtime))
)
SELECT * FROM ptf2;
Table API 에서 프레임워크는 함수의 연속 적용을 가능하게 합니다:
env.from("t")
.partitionBy($("name"))
.process("f1", descriptor("ts").asArgument("on_time"))
.partitionBy($("name"))
.process("f2", descriptor("rowtime").asArgument("on_time"))
구현 가이드 (Implementation Guide)
PTF 는 Flink 의 다른 사용자 정의 함수와 유사한 구현 원칙을 따릅니다. 자세한 내용은 사용자 정의 함수 구현 가이드 를 참고하세요. 이 페이지는 PTF 특유의 사항에 초점을 맞춥니다.
process table function 을 정의하려면 org.apache.flink.table.functions 의 기본 클래스 ProcessTableFunction 을 확장하고 eval(...) 이라는 이름의 평가 메서드를 구현해야 합니다. eval() 메서드는 지원되는 입력 인수와, 상태 저장 PTF 의 경우 상태 항목을 선언합니다.
eval() 의 시그니처는 다음 패턴을 따라야 합니다:
eval( <context>? , <state entry>* , <call argument>* )
평가 메서드는 public 으로 선언되어야 하며 static 이 아니어야 합니다. 오버로딩은 지원되지 않습니다.
사용자 정의 함수를 카탈로그에 저장하려면 클래스에 기본 생성자가 있어야 하며 런타임에 인스턴스화할 수 있어야 합니다. Table API 의 익명 인라인 함수는 함수 객체가 상태 저장적이지 않은 경우(즉 transient 및 static 필드만 포함)에만 영속화될 수 있습니다.
데이터 타입 (Data Types)
기본적으로 입력 및 출력 데이터 타입은 리플렉션을 사용해 자동으로 추출됩니다. 여기에는 출력 데이터 타입을 결정하기 위한 클래스의 제네릭 인수 T 가 포함됩니다. 입력 인수는 eval() 메서드에서 파생됩니다. 리플렉션 정보가 충분하지 않으면 @FunctionHint, @ArgumentHint, @DataTypeHint 어노테이션으로 지원되고 보강될 수 있습니다.
스칼라 함수와 달리 평가 메서드 자체에는 반환 타입이 없어야 합니다. 대신 테이블 함수는 0개, 1개 또는 그 이상의 레코드를 내보내기 위해 평가 메서드 내에서 호출할 수 있는 collect(T) 메서드를 제공합니다. 반환된 레코드는 하나 이상의 필드로 구성될 수 있습니다. 출력 레코드가 단일 필드로만 구성되면 구조적 레코드를 생략하고 런타임에 암시적으로 row 로 감싸질 스칼라 값을 내보낼 수 있습니다.
다음 예제는 데이터 타입을 지정하는 방법을 보여줍니다:
Java:
// Function that accepts two scalar INT arguments and emits them as an implicit ROW<INT>
class AdditionFunction extends ProcessTableFunction<Integer> {
public void eval(Integer a, Integer b) {
collect(a + b);
}
}
// Function that produces an explicit ROW<i INT, s STRING> from scalar arguments,
// the function hint helps in declaring the row's fields
@DataTypeHint("ROW<i INT, s STRING>")
class DuplicatorFunction extends ProcessTableFunction<Row> {
public void eval(Integer i, String s) {
collect(Row.of(i, s));
collect(Row.of(i, s));
}
}
// Function that accepts a scalar DECIMAL(10, 4) and emits it as
// an explicit ROW<DECIMAL(10, 4)>
@FunctionHint(output = @DataTypeHint("ROW<d DECIMAL(10, 4)>"))
class DuplicatorFunction extends ProcessTableFunction<Row> {
public void eval(@DataTypeHint("DECIMAL(10, 4)") BigDecimal d) {
collect(Row.of(d));
collect(Row.of(d));
}
}
인수 (Arguments)
@ArgumentHint 어노테이션은 각 인수의 이름, 데이터 타입, 특성(traits) 을 선언할 수 있게 합니다.
대부분의 경우 시스템이 이름과 데이터 타입을 리플렉션으로 자동 추론할 수 있으므로 지정할 필요가 없습니다. 그러나 특히 인수의 종류를 정의할 때 특성은 명시적으로 제공되어야 합니다. PTF 의 인수는 ArgumentTrait.SCALAR, ArgumentTrait.ROW_SEMANTIC_TABLE, 또는 ArgumentTrait.SET_SEMANTIC_TABLE 로 설정될 수 있습니다. 기본적으로 인수는 스칼라 값으로 취급됩니다.
다음 예제는 @ArgumentHint 어노테이션의 사용법을 보여줍니다:
Java:
// Function that has two arguments:
// "input_table" (a table with set semantics) and "threshold" (a scalar value)
class ThresholdFunction extends ProcessTableFunction<Integer> {
public void eval(
// For table arguments, a data type for Row is optional (leading to polymorphic behavior)
@ArgumentHint(value = ArgumentTrait.SET_SEMANTIC_TABLE, name = "input_table") Row t,
// Scalar arguments require a data type either explicit or via reflection
@ArgumentHint(value = ArgumentTrait.SCALAR, name = "threshold") Integer threshold
) {
int amount = t.getFieldAs("amount");
if (amount >= threshold) {
collect(amount);
}
}
}
테이블 인수
특성 ArgumentTrait.SET_SEMANTIC_TABLE 및 ArgumentTrait.ROW_SEMANTIC_TABLE 은 테이블 인수를 정의합니다.
테이블 인수는 구체적인 데이터 타입(row 또는 structured type) 을 선언하거나 다형성 방식으로 어떤 유형의 row 든 받아들일 수 있습니다.
Row 클래스는 테이블 인수 또는 스칼라 인수로 선언될 수 있습니다.
스칼라 인수의 경우 데이터 타입이 완전히 지정되어야 하고 값이 제공되어야 합니다. 예: f(my_scalar_arg => ROW(12).
테이블 인수의 경우 전체 데이터 타입은 선택 사항이며 한 행 대신 테이블을 기대합니다. 예: f(my_table_arg => TABLE t).
Java:
// Function with explicit table argument type of row
class MyPTF extends ProcessTableFunction<String> {
public void eval(
Context ctx,
@ArgumentHint(value = ArgumentTrait.SET_SEMANTIC_TABLE, type = @DataTypeHint("ROW<s STRING>")) Row t
) {
TableSemantics semantics = ctx.tableSemanticsFor("t");
// Always returns "ROW < s STRING >"
semantics.dataType();
...
}
}
// Function with explicit table argument type of structured type "Customer"
class MyPTF extends ProcessTableFunction<String> {
public void eval(
Context ctx,
@ArgumentHint(value = ArgumentTrait.SET_SEMANTIC_TABLE) Customer c
) {
TableSemantics semantics = ctx.tableSemanticsFor("c");
// Always returns structured type of "Customer"
semantics.dataType();
...
}
}
// Function with polymorphic table argument
class MyPTF extends ProcessTableFunction<String> {
public void eval(
Context ctx,
@ArgumentHint(value = ArgumentTrait.SET_SEMANTIC_TABLE) Row t
) {
TableSemantics semantics = ctx.tableSemanticsFor("t");
// Always returns "ROW" but content depends on the table that is passed into the call
semantics.dataType();
...
}
}
컨텍스트 (Context)
입력 테이블에 대한 추가 정보와 프레임워크가 제공하는 기타 서비스를 위해 Context 를 eval() 메서드의 첫 번째 인수로 추가할 수 있습니다.
Java:
// Function that accesses the Context for reading the PARTITION BY columns and
// excluding them when building a result string
class ConcatNonKeysFunction extends ProcessTableFunction<String> {
public void eval(Context ctx, @ArgumentHint(ArgumentTrait.SET_SEMANTIC_TABLE) Row inputTable) {
TableSemantics semantics = ctx.tableSemanticsFor("inputTable");
List<Integer> keys = Arrays.asList(semantics.partitionByColumns());
return IntStream.range(0, inputTable.getArity())
.filter(pos -> !keys.contains(pos))
.mapToObj(inputTable::getField)
.map(Object::toString)
.collect(Collectors.joining(", "));
}
}
상태 (State)
집합 의미론 테이블을 받는 PTF 는 상태 저장적일 수 있습니다. 중간 결과는 버퍼링, 캐싱, 집계되거나 반복 접근을 위해 단순히 저장될 수 있습니다. 함수는 프레임워크가 관리하는 하나 이상의 상태 항목을 가질 수 있습니다. Flink 는 실패나 재시작 중에 이를 저장하고 복원하는 것을 처리합니다(즉 Flink 관리 상태).
상태 항목은 키로 파티셔닝되며 전역으로 접근할 수 없습니다. 파티셔닝(또는 파티셔닝이 없는 경우 단일 파티션) 은 해당 함수 호출에 의해 정의됩니다. 다시 말해: 가상 프로세서가 전체 테이블의 일부에만 접근할 수 있는 것처럼, PTF 는 PARTITION BY 절이 정의하는 전체 상태의 일부에만 접근할 수 있습니다. Flink 에서 이 개념은 keyed state 로도 알려져 있습니다.
상태 항목은 eval() 메서드의 가변 매개변수로 추가될 수 있습니다. 호출 인수와 구분하기 위해 Context 매개변수 뒤, 다른 모든 인수 앞에 선언되어야 합니다. 또한 @StateHint 로 어노테이션하거나 @FunctionHint(state = ...) 의 일부로 선언해야 합니다.
읽기 및 쓰기 접근을 위해 row 또는 structured type(기본 생성자가 있는 POJO) 만 데이터 타입으로 적합합니다. 상태가 없으면 모든 필드는 row type 의 경우 null, structured type 의 경우 기본값으로 설정됩니다. 상태 효율성을 위해 모든 필드를 nullable 로 유지하는 것이 권장됩니다.
Java:
// Function that counts and stores its intermediate result in the CountState object
// which will be persisted by Flink
class CountingFunction extends ProcessTableFunction<String> {
public static class CountState {
public long count = 0L;
}
public void eval(
@StateHint CountState memory,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
memory.count++;
collect("Seen rows: " + memory.count);
}
}
// Function that waits for a second event coming in
class CountingFunction extends ProcessTableFunction<String> {
public static class SeenState {
public String first;
}
public void eval(
@StateHint SeenState memory,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
if (memory.first == null) {
memory.first = input.toString();
} else {
collect("Event 1: " + memory.first + " and Event 2: " + input.toString());
}
}
}
// Function that uses Row for state
class CountingFunction extends ProcessTableFunction<String> {
public void eval(
@StateHint(type = @DataTypeHint("ROW<count BIGINT>")) Row memory,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
Long newCount = 1L;
if (memory.getField("count") != null) {
newCount += memory.getFieldAs("count");
}
memory.setField("count", newCount);
collect("Seen rows: " + newCount);
}
}
상태 TTL
각 상태 항목에 대해 TTL(Time-to-live) 기간을 지정할 수 있으며, Flink 의 상태 백엔드는 TTL 이 만료되면 항목을 자동으로 정리합니다.
@StateHint(ttl = "...") 어노테이션은 유휴 상태(즉 create 또는 write 연산으로 업데이트되지 않은 상태) 가 얼마나 오래 유지될지의 최소 시간 간격을 지정합니다. 상태는 최소 시간보다 덜 유휴했을 동안에는 절대 지워지지 않으며, 유휴 이후 어느 시점에는 지워집니다.
계속 성장하는 상태 크기를 효율적으로 관리하거나 데이터 보호 요구사항을 준수하기 위해 TTL 을 사용하세요.
정리는 processing time 을 기반으로 하며, 이는 System.currentTimeMillis() 가 정의하는 벽시계 시간에 해당합니다.
제공된 문자열은 Flink 의 duration 문법(예: "3 days", "45 min", "3 hours", "60 s") 을 사용해야 합니다. 단위가 지정되지 않으면 값은 밀리초로 해석됩니다. 상태 항목의 TTL 설정은 전체 파이프라인의 전역 상태 TTL 구성 table.exec.state.ttl 보다 높은 우선순위를 가집니다.
기본적으로 TTL 은 Long.MAX_VALUE 로 설정되어 상태 레이아웃에서 합리적인 값의 향후 조정을 허용합니다. 상태 크기가 우려되고 TTL 이 불필요하면 상태 레이아웃에서 TTL 을 효과적으로 제외하는 0 으로 설정할 수 있습니다.
Java:
// Function with 3 state entries each using a different TTL.
class CountingFunction extends ProcessTableFunction<String> {
public void eval(
Context ctx,
@StateHint(ttl = "1 hour") SomeState shortTermState,
@StateHint(ttl = "1 day") SomeState longTermState,
@StateHint SomeState infiniteState, // potentially influenced by table.exec.state.ttl
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
...
}
}
큰 상태 (Large State)
Flink 의 상태 백엔드는 큰 상태를 효율적으로 처리하기 위한 다양한 유형의 상태를 제공합니다.
현재 PTF 는 세 가지 유형의 상태를 지원합니다:
- Value state: 단일 값을 나타냅니다.
- List state: 값을 리스트로 나타내며, 추가, 제거, 반복 같은 연산을 지원합니다.
- Map state: 개별 항목의 효율적인 조회, 수정, 제거를 위한 map(키-값 쌍) 을 나타냅니다.
기본적으로 PTF 의 상태 항목은 value state 로 표현됩니다. 즉 평가 메서드가 호출될 때마다 모든 상태 항목이 상태 백엔드에서 완전히 읽히고, 평가 메서드가 끝나면 값이 상태 백엔드에 다시 기록됩니다.
상태 접근을 최적화하고 불필요한 (역)직렬화를 피하기 위해 상태 항목은 다음으로 선언될 수 있습니다:
org.apache.flink.table.api.dataview.ListView(리스트 상태용)org.apache.flink.table.api.dataview.MapView(맵 상태용)
이들은 기반 Flink 상태 백엔드에 대한 직접 뷰를 제공합니다.
예를 들어 MapView 를 사용할 때 MapView#get 으로 값에 접근하면 지정된 키와 연관된 값만 역직렬화합니다. 전체 맵을 로드하지 않고 개별 항목에 효율적으로 접근할 수 있습니다. 이 접근 방식은 맵이 전체적으로 메모리에 들어가지 않을 때 특히 유용합니다.
State TTL 은 리스트 또는 맵의 각 항목에 개별적으로 적용되어 상태 요소에 대한 세밀한 만료 제어를 가능하게 합니다.
다음 예제는 MapView 를 선언하고 사용하는 방법을 보여줍니다. PTF 가 스키마 (userId, eventId, ...) 를 가진 테이블을 처리하고 userId 로 파티셔닝되며 distinct eventId 값의 카디널리티가 높다고 가정합니다. 이 사용 사례에서는 일반적으로 테이블을 userId 와 eventId 둘 다로 파티셔닝하는 것이 좋습니다. 예시 목적으로 큰 상태는 map state 로 저장됩니다.
Java:
// Function that uses a map view for storing a large map for an event history per user
class LargeHistoryFunction extends ProcessTableFunction<String> {
public void eval(
@StateHint MapView<String, Integer> largeMemory,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
String eventId = input.getFieldAs("eventId");
Integer count = largeMemory.get(eventId);
if (count == null) {
largeMemory.put(eventId, 1);
} else {
if (count > 1000) {
collect("Anomaly detected: " + eventId);
}
largeMemory.put(eventId, count + 1);
}
}
}
다른 데이터 타입과 유사하게 리플렉션을 사용해 필요한 타입 정보를 추출합니다. 리플렉션이 불가능한 경우—Row 객체가 관련될 때처럼—타입 힌트를 제공할 수 있습니다. 리스트 뷰에는 ARRAY 데이터 타입을, 맵 뷰에는 MAP 데이터 타입을 사용하세요.
Java:
// Function that uses a list view of rows
class LargeHistoryFunction extends ProcessTableFunction<String> {
public void eval(
@StateHint(type = @DataTypeHint("ARRAY<ROW<s STRING, i INT>>")) ListView<Row> largeMemory,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
...
}
}
효율성과 설계 원칙
상태 저장 함수는 데이터 레이아웃과 데이터 보존을 잘 생각해야 함을 의미합니다. 무한한 파티션 수(개방된 키스페이스) 또는 파티션 내에서도 상태가 계속 커질 수 있습니다. @StateHint(ttl = ... ) 을 설정하거나 결국 Context.clearAllState() 를 호출하는 것을 고려하세요.
Java:
// Function that waits for a second event coming in BUT with better state efficiency
class CountingFunction extends ProcessTableFunction<String> {
public static class SeenState {
public String first;
}
public void eval(
Context ctx,
@StateHint(ttl = "1 day") SeenState memory,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
if (memory.first == null) {
memory.first = input.toString();
} else {
collect("Event 1: " + memory.first + " and Event 2: " + input.toString());
ctx.clearAllState();
}
}
}
시간과 타이머 (Time and Timers)
PTF 는 이벤트 시간을 기본적으로 지원합니다. 시간 기반 연산은 Context#timeContext(Class) 로 접근할 수 있습니다.
시간 컨텍스트는 항상 현재 처리 중인 이벤트로 범위가 지정됩니다. 이벤트는 현재 입력 행 또는 발화하는 타이머일 수 있습니다.
시간과 타이머의 타임스탬프는 java.time.Instant, java.time.LocalDateTime, 또는 Long 으로 표현될 수 있습니다. 이러한 타임스탬프는 epoch 이후 밀리초를 기반으로 하며 로컬 세션 시간대를 고려하지 않습니다. 시간 클래스는 timeContext() 의 인수로 전달될 수 있습니다.
시간 (Time)
모든 PTF 는 선택적 on_time 인수를 받습니다. 함수 호출의 on_time 인수는 워터마크가 선언된 시간 속성 열을 선언합니다. 테이블의 행을 처리할 때 이 타임스탬프는 TimeContext#time() 으로, 워터마크는 각각 TimeContext#currentWatermark()/TimeContext#tableWatermark() 로 접근할 수 있습니다.
함수 호출에서 on_time 인수를 지정하면 프레임워크가 이후 시간 기반 연산을 위해 함수의 출력에 rowtime 열을 반환하도록 지시합니다.
on_time 속성 선언을 위한 SQL 문법:
SELECT * FROM f(..., on_time => DESCRIPTOR(`my_timestamp`));
on_time 속성 선언을 위한 Table API:
.process(MyFunction.class, ..., descriptor("my_timestamp").asArgument("on_time"));
ArgumentTrait.REQUIRE_ON_TIME 은 필요하면 on_time 인수를 필수로 만듭니다.
on_time 인수가 제공되면 타이머를 사용할 수 있습니다. 다음 동기 부여 예제는 eval() 과 onTimer() 가 함께 작동하는 방식을 보여줍니다:
Java:
// Function that sends out a ping for the given key.
// The ping is sent one minute after the last event for this key was observed.
public static class PingLaterFunction extends ProcessTableFunction<String> {
public void eval(
Context ctx,
@ArgumentHint({ArgumentTrait.SET_SEMANTIC_TABLE, ArgumentTrait.REQUIRE_ON_TIME}) Row input
) {
TimeContext<Instant> timeCtx = ctx.timeContext(Instant.class);
// Replaces an existing timer and thus potentially resets the minute if necessary
timeCtx.registerOnTime("ping", timeCtx.time().plus(Duration.ofMinutes(1)));
}
public void onTimer(OnTimerContext onTimerCtx) {
collect("ping");
}
}
TableEnvironment env = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
// Create a table of watermarked events
env.executeSql(
"CREATE TABLE Events (ts TIMESTAMP_LTZ(3), id STRING, val INT, WATERMARK FOR ts AS ts - INTERVAL '2' SECONDS) " +
"WITH ('connector' = 'datagen')");
// Use Expressions.descriptor("ts") to mark the "on_time" argument in Table API
env.from("Events")
.partitionBy($("id"))
.process(PingLaterFunction.class, descriptor("ts").asArgument("on_time"))
.execute()
.print();
// For SQL register the function and pass the DESCRIPTOR argument
env.createFunction("PingLaterFunction", PingLaterFunction.class);
env.executeSql("SELECT * FROM PingLaterFunction(input => TABLE Events PARTITION BY id, on_time => DESCRIPTOR(`ts`))").print();
결과는 다음과 유사합니다:
| op | id | EXPR$0 | rowtime |
+----+--------------------------------+--------------------------------+-------------------------+
| +I | a1b9204c341c2d136c7d494fe29... | ping | 2025-03-26 16:56:47.842 |
| +I | 45a864eed208f4c0eebb40d55fd... | ping | 2025-03-26 16:56:47.842 |
| +I | 2039c6d08896a3255b9bfc38ad0... | ping | 2025-03-26 16:56:47.842 |
| +I | 41a9269ed057793a0ea97c31f61... | ping | 2025-03-26 16:56:47.842 |
on_time 속성이 선언되면 출력 끝에 rowtime 열이 포함됩니다. 이 열은 이후 PTF 호출이나 시간 기반 연산에 사용될 수 있는 또 다른 워터마크 시간 속성을 나타냅니다. rowtime 의 데이터 타입은 입력의 시간 속성에서 파생됩니다.
현재 타임스탬프
TimeContext#time() 는 현재 처리 중인 이벤트의 타임스탬프를 반환합니다.
이벤트는 테이블의 행 또는 발화하는 타이머일 수 있습니다:
1. 행 이벤트 타임스탬프 — eval() 메서드 내에서 현재 처리 중인 행의 타임스탬프. 함수 호출의 on_time 인수에 의해 구동되며, 이 메서드는 참조된 시간 속성 열의 내용을 반환합니다. on_time 인수가 현재 처리 중인 테이블의 시간 속성 열을 참조하지 않으면 null 을 반환합니다.
2. 타이머 이벤트 타임스탬프 — onTimer() 메서드 내에서 현재 처리 중인 발화 타이머의 타임스탬프.
테이블 워터마크
TimeContext#tableWatermark() 는 현재 처리 중인 입력 테이블의 이벤트 시간 워터마크를 반환합니다.
워터마크는 소스에서 생성되어 논리적 클록을 진행하기 위해 토폴로지를 통해 전송됩니다. 입력 테이블의 현재 워터마크는 테이블을 생성하는 모든 업스트림 Flink subtask 의 최소 워터마크입니다.
다중 입력 시나리오에서 각 입력 테이블은 자체 독립적인 워터마크를 가질 수 있습니다. 이 메서드는 eval() 메서드에서 현재 처리 중인 입력 테이블에 특정된 워터마크를 반환합니다. (currentWatermark() 가 반환하는) 모든 입력 테이블 간의 전역 최소 워터마크가 아닙니다.
이는 입력별로 늦은(late) 이벤트 감지에 특히 유용합니다.
처리 중인 입력 테이블의 현재 워터마크를 반환합니다. onTimer() 메서드 내에서 호출되거나 테이블을 생성하는 모든 업스트림 Flink subtask 로부터 워터마크를 아직 받지 못했으면 null 입니다.
현재 워터마크
TimeContext#currentWatermark() 는 이 PTF 인스턴스의 현재 이벤트 시간 워터마크를 반환합니다.
워터마크는 소스에서 생성되어 논리적 클록을 진행하기 위해 토폴로지를 통해 전송됩니다. PTF 인스턴스의 현재 워터마크는 모든 입력 테이블의 전역 최소 워터마크(즉 모든 업스트림 Flink subtask 와 테이블 파티션에 걸친) 입니다.
이 메서드는 PTF 를 평가하는 Flink subtask 의 현재 워터마크를 반환합니다. 따라서 반환된 타임스탬프는 현재 처리 중인 입력 테이블과 파티션과 무관하게 전체 Flink subtask 를 나타냅니다. 이 동작은 SQL 의 SELECT CURRENT_WATERMARK(...) 호출과 유사합니다.
모든 업스트림 Flink subtask 와 테이블 파티션에 걸친 PTF 인스턴스의 현재 워터마크를 반환합니다. 모든 입력에 걸쳐 최소 논리적 시간을 계산할 수 없으면 null 값이 반환됩니다. 이는 시작 또는 복구 중 하나 이상의 활성(즉 유휴하지 않은) 입력이 아직 워터마크를 보내지 않았을 때 발생합니다.
타이머 (Timers)
집합 의미론 테이블을 받는 PTF 는 타이머를 지원할 수 있습니다. 타이머는 더 나중 시점에 처리를 계속할 수 있게 합니다. 이는 대기, 동기화 또는 타임아웃을 가능하게 합니다. 타이머는 워터마크가 논리적 클록을 진행할 때 등록된 시간에 발화합니다.
타이머는 이름을 가질 수 있고(TimeContext#registerOnTime(String, TimeType)) 이름이 없을 수 있습니다(TimeContext#registerOnTime(TimeType)). 타이머의 이름은 기존 타이머를 교체하거나 삭제하는 데, 또는 발화할 때 OnTimerContext#currentTimer() 로 여러 타이머를 식별하는 데 유용할 수 있습니다.
타이머 이벤트에 반응하기 위해 eval() 메서드 옆에 onTimer() 메서드를 선언해야 합니다. onTimer() 메서드의 시그니처는 선택적 OnTimerContext 다음에 (eval() 메서드에 선언된) 모든 상태 항목을 포함해야 합니다.
onTimer() 의 시그니처는 다음 패턴을 따라야 합니다:
onTimer( <on timer context>? , <state entry>* )
Flink 는 실패나 재시작 중에 타이머를 저장하고 복원하는 것을 처리합니다. 따라서 타이머는 특별한 종류의 상태입니다. 마찬가지로 타이머는 PARTITION BY 절이 정의하는 가상 프로세서로 범위가 지정됩니다. 타이머는 현재 가상 프로세서에서만 등록되고 삭제될 수 있습니다.
다음 예제는 타이머를 등록하고 지우는 방법을 보여줍니다:
Java:
// Function that waits for a second event or timeouts after 60 seconds
class TimerFunction extends ProcessTableFunction<String> {
public static class SeenState {
public String seen = null;
}
public void eval(
Context ctx,
@StateHint SeenState memory,
@ArgumentHint({SET_SEMANTIC_TABLE, REQUIRE_ON_TIME}) Row input
) {
TimeContext<Instant> timeCtx = ctx.timeContext(Instant.class);
if (memory.seen == null) {
memory.seen = input.getField(0).toString();
timeCtx.registerOnTimer("timeout", timeCtx.time().plusSeconds(60));
} else {
collect("Second event arrived for: " + memory.seen);
ctx.clearAll();
}
}
public void onTimer(OnTimerContext onTimerCtx, SeenState memory) {
collect("Timeout for: " + memory.seen);
}
}
늦은 레코드 처리
늦은 레코드는 현재 워터마크보다 작거나 같은 시간 속성 값을 가진 레코드입니다. PTF 는 eval() 메서드를 호출함으로써 늦은 레코드를 늦지 않은 레코드처럼 처리합니다. on_time 인수가 지정되면 늦은 타임스탬프가 출력에 보존됩니다. 이 동작은 행 및 집합 의미론 PTF 에 대해 동일합니다.
현재 워터마크보다 작거나 같은 시간에 타이머를 등록하는 것은 허용됩니다. eval() 내에서 등록되면 타이머는 다음 워터마크 진행 시 발화합니다. onTimer() 내에서 등록되면 타이머는 현재 타이머가 끝난 직후 발화합니다. onTimer() 내에서 과거 시간 타이머를 무조건 재등록하면 무한 루프가 발생함에 유의하세요.
효율성과 설계 원칙
너무 많은 타이머를 등록하면 성능에 영향을 줄 수 있습니다. 무한한 파티션 수(개방된 키스페이스) 또는 파티션 내에서도 타이머 상태가 계속 커질 수 있습니다. 따라서 등록된 타이머 수를 최소한으로 줄이고 더 이상 필요하지 않으면 Context#clearAllTimers() 또는 TimeContext#clearTimer(String) 로 타이머를 정리하는 것을 고려하세요.
정렬 (Ordering)
집합 의미론을 가진 테이블을 받는 PTF 는 함수 호출에서 선택적 ORDER BY 절을 지정해 각 파티션 내에서 행이 처리되는 순서를 정의할 수 있습니다. ORDER BY 절은 행이 지정된 순서로 eval() 메서드에 전달되도록 보장합니다.
ORDER BY 절은 첫 번째 열이 시간 속성 열(즉 워터마크 선언이 있는 TIMESTAMP 또는 TIMESTAMP_LTZ 열) 이어야 함을 요구합니다. 첫 번째 ORDER BY 열은 오름차순으로 지정되어야 합니다. 이는 행이 이벤트 시간 순서로 처리되도록 보장합니다. 추가 열은 같은 타임스탬프를 가진 행의 순서를 정의하는 보조 정렬 키로 지정될 수 있습니다.
SQL:
SELECT * FROM my_ptf(
input_table => TABLE source_table
PARTITION BY user_id
ORDER BY (event_time ASC, priority DESC NULLS FIRST)
)
Java:
env.from("source_table")
.partitionBy($("user_id"))
.orderBy($("event_time").asc(), $("priority").desc())
.process(MyPTF.class, descriptor("event_time").asArgument("on_time"))
ORDER BY 와 on_time 인수의 차이
ORDER BY 와 on_time 인수 모두 시간 속성과 관련이 있지만 서로 다른 목적을 제공합니다:
- on_time: 시간 컨텍스트(
TimeContext#time()) 와 출력 타임스탬프를 뒷받침하는 시간 속성 열을 선언합니다. 행의 처리 순서에는 영향을 주지 않습니다. - ORDER BY: 각 파티션 내의 행을 물리적으로 버퍼링하고 정렬해 eval() 메서드로의 순서 있는 전달을 보장합니다. 같은 테이블 인수에 ORDER BY 와
on_time이 모두 지정되면 같은 시간 속성 열을 참조해야 합니다.
정렬 보장과 늦은 이벤트
ORDER BY 가 시간 속성 열에 지정되면 프레임워크는 무순서 이벤트를 재정렬하기 위해 파티션 및 입력 테이블당 정렬 버퍼를 유지합니다. 정렬 버퍼는 주어진 입력 테이블의 워터마크가 진행될 때 플러시되며, 그 시점에 워터마크보다 작거나 같은 타임스탬프를 가진 모든 버퍼된 행이 정렬된 순서로 eval() 메서드에 전달됩니다. 늦은 이벤트(워터마크 후 도착) 는 정렬 보장을 유지하기 위해 버려집니다.
다음 예제는 보조 정렬이 있는 정렬된 처리를 보여줍니다. 먼저 함수 구현입니다:
// Function that processes events in order and captures the ordering
public static class OrderedProcessor extends ProcessTableFunction<List<Event>> {
public record Event(Integer score, Instant ts) {}
public static class BufferState {
// Stores all input events that enter the eval() after sorting
public List<Event> events = new ArrayList<>();
}
public void eval(
Context ctx,
@StateHint BufferState state,
@ArgumentHint(SET_SEMANTIC_TABLE) Row input
) {
// Optional: Access ordering information at runtime
TableSemantics semantics = ctx.tableSemanticsFor("input");
int[] orderColumns = semantics.orderByColumns();
SortDirection[] directions = semantics.orderByDirections();
// Buffer the incoming row
state.events.add(new Event(input.getFieldAs("score"), input.getFieldAs("ts")));
// Emit current buffer state
collect(state.events);
}
}
이 함수는 SQL 또는 Table API 로 호출될 수 있습니다:
SQL:
-- Create a watermarked table
CREATE TABLE Events (
name STRING,
score INT,
ts TIMESTAMP_LTZ(3),
WATERMARK FOR ts AS ts - INTERVAL '1' SECOND
) WITH (
'connector' = 'datagen'
);
-- Register the function
CREATE FUNCTION OrderedProcessor AS 'org.example.OrderedProcessor';
-- Use ORDER BY with primary and secondary sort columns
SELECT * FROM OrderedProcessor(
input => TABLE Events
PARTITION BY name
ORDER BY (ts ASC, score DESC)
);
Java:
TableEnvironment env = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
// Create a watermarked table
env.executeSql(
"CREATE TABLE Events (" +
"name STRING, " +
"score INT, " +
"ts TIMESTAMP_LTZ(3), " +
"WATERMARK FOR ts AS ts - INTERVAL '1' SECOND" +
") WITH ('connector' = 'datagen')"
);
// Use orderBy() with primary and secondary sort columns
env.from("Events")
.partitionBy($("name"))
.orderBy($("ts").asc(), $("score").desc())
.process(OrderedProcessor.class)
.execute()
.print();
이 예제에서:
- 이벤트는 먼저
ts(오름차순) 로 정렬되어 이벤트 시간 순서를 보장합니다. - 같은 타임스탬프를 가진 이벤트는 그 다음
score로 정렬됩니다. - 늦은 이벤트(워터마크보다 작은 타임스탬프) 는 자동으로 버려집니다.
TableSemanticsAPI 는 런타임에 정렬 구성에 대한 접근을 제공합니다.- 출력은 정렬된 입력 이벤트의 계속 커지는 목록입니다.
여러 테이블 (Multiple Tables)
PTF 는 여러 테이블을 동시에 처리할 수 있습니다. 이는 다양한 사용 사례를 가능하게 합니다:
- 상태를 효율적으로 관리하는 custom joins 구현.
- 사이드 입력(side inputs) 으로서 차원 테이블의 정보로 메인 테이블을 풍부화.
- 런타임에 키가 지정된 가상 프로세서에 제어 이벤트 보내기.
eval() 메서드는 여러 입력을 지원하기 위해 여러 테이블 인수를 지정할 수 있습니다. 모든 테이블 인수는 집합 의미론으로 선언되고 일관된 파티셔닝을 사용해야 합니다. 다시 말해, PARTITION BY 절의 열 수와 데이터 타입이 관련된 모든 테이블 인수에 걸쳐 일치해야 합니다.
어느 입력의 행이든 한 번에 하나씩 함수에 전달됩니다. 따라서 한 번에 하나의 테이블 인수만 null 이 아닙니다. 현재 처리 중인 입력을 결정하려면 null 검사를 사용하세요.
시스템은 다음에 어떤 입력 행이 가상 프로세서를 통해 스트리밍될지 결정합니다. PTF 에서 제대로 처리되지 않으면 입력 간의 경쟁 조건(race condition) 을 초래하고 결과적으로 비결정적 결과를 만들 수 있습니다. 주어진 워터마크까지 모든 행이 도착하기를 기다리는 시간 기반 조인 또는 PTF 가 특정 조건이 충족될 때까지 하나 이상의 입력 행을 버퍼링하는 조건 기반으로 함수를 설계하는 것이 권장됩니다.
예제: Custom Join
다음 예제는 두 테이블 사이의 custom join 을 구현하는 방법을 보여줍니다:
Java:
TableEnvironment env = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
env.executeSql("CREATE VIEW Visits(name) AS VALUES ('Bob'), ('Alice'), ('Bob')");
env.executeSql("CREATE VIEW Purchases(customer, item) AS VALUES ('Alice', 'milk')");
env.createFunction("Greeting", GreetingWithLastPurchase.class);
env
.executeSql("SELECT * FROM Greeting(TABLE Visits PARTITION BY name, TABLE Purchases PARTITION BY customer)")
.print();
// --------------------
// Function declaration
// --------------------
// Function that greets a customer and suggests the last purchase made, if available.
public static class GreetingWithLastPurchase extends ProcessTableFunction<String> {
// Keep the last purchased item in state
public static class LastItemState {
public String lastItem;
}
// The eval() method takes two @ArgumentHint(SET_SEMANTIC_TABLE) arguments
public void eval(
@StateHint LastItemState state,
@ArgumentHint(SET_SEMANTIC_TABLE) Row visit,
@ArgumentHint(SET_SEMANTIC_TABLE) Row purchase) {
// Process row from table Purchases
if (purchase != null) {
state.lastItem = purchase.getFieldAs("item");
}
// Process row from table Visits
else if (visit != null) {
if (state.lastItem == null) {
collect("Hello " + visit.getFieldAs("name") + ", let me know if I can help!");
} else {
collect("Hello " + visit.getFieldAs("name") + ", here to buy " + state.lastItem + " again?");
}
}
}
}
결과는 다음과 유사합니다:
+----+--------------------------------+--------------------------------+--------------------------------+
| op | name | customer | EXPR$0 |
+----+--------------------------------+--------------------------------+--------------------------------+
| +I | Bob | Bob | Hello Bob, let me know if I... |
| +I | Alice | Alice | Hello Alice, here to buy Pr... |
| +I | Bob | Bob | Hello Bob, let me know if I... |
+----+--------------------------------+--------------------------------+--------------------------------+
효율성과 설계 원칙
입력 테이블 수가 많으면 단일 TaskManager 또는 subtask 에 부정적인 영향을 줄 수 있습니다. 각 입력에 대해 네트워크 버퍼를 할당해야 하므로 메모리 소비가 증가하여 테이블 인수 수는 최대 20 개로 제한됩니다.
고르지 않게 분포된 키는 단일 가상 프로세서에 과부하를 일으켜 백프레셔로 이어질 수 있습니다. 적절한 파티션 키를 선택하는 것이 중요합니다.
UID 를 통한 쿼리 진화 (Query Evolution with UIDs)
다른 SQL operator 와 달리 PTF 는 상태 저장 쿼리 진화를 지원합니다.
플래너의 관점에서 PTF 는 최적화되지 않은 상태 저장 빌딩 블록이며, 플래너는 쿼리의 주변 operator 를 최적화합니다. PTF 의 상태 항목은 주변 쿼리나 PTF 자체가 변경되더라도 Flink 에 의해 영속화되고 복원될 수 있습니다. 상태 항목의 스키마가 변경되지 않는 한 말입니다.
향후 쿼리 진화를 위해 프레임워크는 집합 의미론 테이블에서 동작하는 모든 PTF 에 고유 식별자(UID) 를 강제합니다. UID 는 암시적 uid 문자열 인수로 제공될 수 있습니다. 이는 PTF 의 상태 항목을 체크포인트나 savepoint 에 영속화할 때 사용됩니다. uid 인수가 지정되지 않으면 프레임워크가 함수 이름을 사용해 문당 하나의 고유한 PTF 호출을 보장합니다. PTF 가 여러 번 호출되면 검증에서 전체 Flink 작업에 걸쳐 고유하도록 수동으로 지정된 UID 가 필요합니다.
또한 UID 는 최적화 도구가 파이프라인의 공통 부분을 병합할지 결정하는 데 도움이 됩니다. 공유 UID 는 단일 상태 저장 PTF 를 유지하면서 fan-out 동작을 가능하게 합니다.
Fan-out 예제
다음 예제에서 최적화 도구는 두 INSERT INTO 문에서 공유 파이프라인 부분 SELECT * FROM f(..., uid => 'same') 을 감지해 단일 상태 저장 PTF operator 를 유지할 수 있습니다. 그런 다음 PTF 의 결과는 filter 조건에 따라 두 대상으로 분할되어 전송됩니다.
EXECUTE STATEMENT SET
BEGIN
INSERT INTO bob_sink SELECT * FROM f(r => TABLE t PARTITION BY name, uid => 'same') WHERE name = 'Bob';
INSERT INTO alice_sink SELECT * FROM f(r => TABLE t PARTITION BY name, uid => 'same') WHERE name = 'Alice';
END;
다른 UID 는 이 최적화를 비활성화하고 공유 테이블 t 를 소비하는 두 개의 상태 저장 블록이 유지됩니다.
EXECUTE STATEMENT SET
BEGIN
INSERT INTO bob_sink SELECT * FROM f(r => TABLE t PARTITION BY name, uid => 'ptf1') WHERE name = 'Bob';
INSERT INTO alice_sink SELECT * FROM f(r => TABLE t PARTITION BY name, uid => 'ptf2') WHERE name = 'Alice';
END;
통과 열 (Pass-Through Columns)
테이블 의미론과 on_time 인수가 정의되었는지 여부에 따라 시스템은 모든 함수 출력에 추가 열을 추가합니다.
집합 의미론을 가진 테이블 인수의 경우 출력은 PARTITION BY 열로 접두됩니다.
on_time 인수가 있는 호출의 경우 출력은 rowtime 으로 접미됩니다.
요약하면 기본 패턴은 다음과 같습니다:
<PARTITION BY keys> | <function output> | <rowtime>
ArgumentTrait.PASS_COLUMNS_THROUGH 는 테이블 인수의 모든 열을 PTF 의 출력에 포함하도록 시스템에 지시합니다.
열 k 와 v 를 포함하는 테이블 t 와 열 c1 과 c2 를 생성하는 PTF f() 가 주어지면 SELECT * FROM f(table_arg => TABLE t PARTITION BY k) 의 출력은 다음 순서를 사용합니다:
Default: | k | c1 | c2 |
With pass-through columns: | k | v | c1 | c2 |
이를 통해 PTF 는 입력 열을 수동으로 전달할 필요 없이 메인 집계에 집중할 수 있습니다.
참고: 통과 열은 단일 테이블 인수를 받고 타이머를 사용하지 않는 append-only PTF 에서만 사용할 수 있습니다.
업데이트와 Changelog
기본적으로 PTF 는 테이블 인수가 append-only 테이블에 의해 뒷받침된다고 가정합니다. 여기서 새 레코드는 기존 레코드를 업데이트하지 않고 테이블에 삽입됩니다. 그러면 PTF 는 새 append-only 테이블을 출력으로 생성합니다.
append-only 테이블은 이상적이며 이벤트 시간과 워터마크에서 원활하게 작동하지만, 업데이트 테이블로 작업해야 하는 시나리오가 있습니다. 이 경우 레코드는 초기 삽입 후 업데이트되거나 삭제될 수 있습니다. 이는 여러 측면에 영향을 미칩니다:
- 상태 관리: 연산은 어떤 레코드든 다시 업데이트될 수 있으므로 더 큰 상태 공간(footprint) 이 필요할 수 있는 가능성을 수용해야 합니다.
- 파이프라인 복잡성: 레코드가 최종적이지 않고 이후에 변경될 수 있으므로 전체 파이프라인 결과가 in-flight 상태로 남습니다.
- 다운스트림 시스템: in-flight 데이터는 Flink 뿐만 아니라 일관성과 최종성이 중요한 다운스트림 시스템에서도 문제를 일으킬 수 있습니다.
효율적이고 고성능의 데이터 처리를 위해 가능하면 append-only 테이블로 파이프라인을 설계해 상태 관리를 단순화하고 업데이트 테이블과 관련된 복잡성을 피할 것을 권장합니다.
PTF 는 그렇게 구성되면 업데이트 테이블을 소비 및/또는 생성할 수 있습니다. 이 섹션은 PTF 를 사용한 CDC(Change Data Capture) 의 간략한 개요를 제공합니다.
CDC 기초
내부적으로 Flink SQL 엔진의 테이블은 changelog 에 의해 뒷받침됩니다. 이러한 changelog 는 INSERT(+I), UPDATE_BEFORE(-U), UPDATE_AFTER(+U), 또는 DELETE(-D) 메시지를 포함하는 CDC(Change Data Capture) 정보를 인코딩합니다.
changelog 에 이러한 플래그가 존재하는 것이 소비자 또는 생산자의 Changelog Mode 를 구성합니다:
Append Mode {+I}
- 모든 메시지는 insert-only 입니다.
- 모든 삽입 메시지는 불변의 사실입니다.
- 메시지는 관련이 없으므로 파티션과 프로세서에 걸쳐 임의의 방식으로 분산될 수 있습니다.
Upsert Mode {+I, +U, -D}
- 메시지는 업데이트 테이블로 이어지는 업데이트를 포함할 수 있습니다.
- 업데이트는 키(즉 upsert key) 로 관련됩니다.
- 모든 메시지는 업서트 키 아래의 결과에 대한 upsert 또는 delete 메시지입니다.
- 같은 upsert 키의 메시지는 같은 파티션과 프로세서에 도착해야 합니다.
- 삭제는 upsert 키 열의 값만(즉 partial deletes) 또는 모든 열의 값(즉 full deletes) 을 포함할 수 있습니다.
- 이 모드는
-U메시지가 없으므로 문헌에서 partial image 로도 알려져 있습니다.
Retract Mode {+I, -U, +U, -D}
- 메시지는 업데이트 테이블로 이어지는 업데이트를 포함할 수 있습니다.
- 모든 삽입 또는 업데이트 이벤트는 "되돌릴 수 있는"(즉 재시도된) 사실입니다.
- 업데이트는 모든 열로 관련됩니다. 단순화하면 전체 행이 일종의 키이지만 중복은 지원됩니다. 예:
+I['Bob', 42]는-D['Bob', 42]와 관련되고,+U['Alice', 13]은-U['Alice', 13]과 관련됩니다. - 따라서 모든 메시지는 삽입(
+) 또는 그 재시도(-) 입니다. - 이 모드는 문헌에서 full image 로 알려져 있습니다.
업데이트 입력 테이블
ArgumentTrait.SUPPORTS_UPDATES 는 주어진 테이블 인수에 업데이트가 입력으로 허용됨을 시스템에 지시합니다. 기본적으로 테이블 인수는 insert-only 이며 업데이트는 거부됩니다.
입력 테이블은 집계나 outer join 같은 하위 쿼리가 증분 계산을 강제할 때 업데이트 테이블이 됩니다. 예를 들어 다음 쿼리는 함수가 재시도 메시지를 소화할 수 있을 때만 작동합니다:
// The change +I[1] followed by -U[1], +U[2], -U[2], +U[3] will enter the function
// if `table_arg` is declared with SUPPORTS_UPDATES
WITH UpdatingTable AS (
SELECT COUNT(*) FROM (VALUES 1, 2, 3)
)
SELECT * FROM f(table_arg => TABLE UpdatingTable)
업데이트를 지원해야 한다면 테이블 인수의 데이터 타입이 변경을 인코딩할 수 있는 방식으로 선택되는지 확인하세요. 다시 말해: RowKind 변경 플래그를 노출하는 Row 타입을 선택하세요.
기반 입력 테이블의 changelog 가 어떤 종류의 변경이 함수로 들어가는지 결정합니다. 입력 테이블이 append-only 이면 함수는 {+I} 를 받습니다. 입력 테이블이 파티션 키와 같은 upsert 키를 사용해 업서트하면 함수는 {+I,+U,-D} 를 받습니다. 그렇지 않으면 재시도 {+I,-U,+U,-D}(즉 RowKind.UPDATE_BEFORE 포함) 가 함수로 들어갑니다. 모든 업데이트 경우에 재시도를 강제하려면 ArgumentTrait.REQUIRE_UPDATE_BEFORE 를 사용하세요.
업서트 테이블의 경우 changelog 가 키 전용 삭제(partial deletions 라고도 함) 를 포함하면 행이 함수로 들어갈 때 upsert 키 필드만 설정됩니다. 비키 필드는 NOT NULL 제약과 관계없이 null 로 설정됩니다. 전체 삭제만 함수로 들어가도록 강제하려면 ArgumentTrait.REQUIRE_FULL_DELETE 를 사용하세요.
SUPPORTS_UPDATES 특성은 고급 사용 사례를 위한 것입니다. 배치 모드에서는 입력이 항상 insert-only 임에 유의하세요. 따라서 PTF 가 배치와 스트리밍 모드에서 같은 결과를 생성해야 한다면 결과는 워터마크와 이벤트 시간을 기반으로 내보내져야 합니다.
Retract 모드 강제
ArgumentTrait.REQUIRE_UPDATE_BEFORE 는 SUPPORT_UPDATES 인 테이블 인수가 업데이트를 인코딩할 때 RowKind.UPDATE_BEFORE 메시지를 포함해야 함을 시스템에 지시합니다. 다시 말해: 업데이트 테이블을 retract changelog mode 로 표현하는 것을 강제합니다.
기본적으로 업데이트는 입력 연산이 내보낸 대로 인코딩됩니다. 따라서 업데이트 테이블은 upsert changelog mode 로 인코딩될 수 있고 삭제는 키만 포함할 수 있습니다.
다음 예제는 입력 changelog 가 업데이트를 다르게 인코딩하는 방법을 보여줍니다:
// Given a table UpdatingTable(name STRING PRIMARY KEY, score INT)
// backed by upsert changelog with changes
// +I[Alice, 42], +I[Bob, 0], +U[Bob, 2], +U[Bob, 100], -D[Bob, NULL].
// Given a function `f` that declares `table_arg` with REQUIRE_UPDATE_BEFORE.
SELECT * FROM f(table_arg => TABLE UpdatingTable PARTITION BY name)
// The following changes will enter the function:
// +I[Alice, 42], +I[Bob, 0], -U[Bob, 0], +U[Bob, 2], -U[Bob, 2], +U[Bob, 100], -U[Bob, 100]
// In both encodings, a materialized table would only contain a row for Alice.
전체 삭제가 있는 Upserts 강제
ArgumentTrait.REQUIRE_FULL_DELETE 는 SUPPORT_UPDATES 인 테이블 인수가 업데이트 테이블이 upsert changelog 에 의해 뒷받침될 때 RowKind.DELETE 메시지에 모든 필드를 포함해야 함을 시스템에 지시합니다.
업서트 테이블의 경우 changelog 가 키 전용 삭제(partial deletes) 를 포함하면 행이 함수로 들어갈 때 upsert 키 필드만 설정됩니다. 비키 필드는 NOT NULL 제약과 관계없이 null 로 설정됩니다.
다음 예제는 입력 changelog 가 업데이트를 다르게 인코딩하는 방법을 보여줍니다:
// Given a table UpdatingTable(name STRING PRIMARY KEY, score INT)
// backed by upsert changelog with changes
// +I[Alice, 42], +I[Bob, 0], +U[Bob, 2], +U[Bob, 100], -D[Bob, NULL].
// Given a function `f` that declares `table_arg` with REQUIRE_FULL_DELETE.
SELECT * FROM f(table_arg => TABLE UpdatingTable PARTITION BY name)
// The following changes will enter the function:
// +I[Alice, 42], +I[Bob, 0], +U[Bob, 2], +U[Bob, 100], -D[Bob, 100].
// In both encodings, a materialized table would only contain a row for Alice.
예제: Changelog 필터링
다음 함수는 PTF 가 업데이트 테이블을 append-only 테이블로 변환할 수 있는 방법을 보여줍니다. 각 Row 에 인코딩된 업데이트를 적용하는 대신 changelog 플래그를 페이로드에 통합합니다. PTF 가 내보낸 행은 RowKind.INSERT 임이 보장됩니다. 원래 changelog 플래그를 페이로드에 보존함으로써 특정 업데이트 유형의 필터링을 허용합니다. 이 예제에서는 모든 삭제를 걸러냅니다.
Java:
TableEnvironment env = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
Table data = env
.fromValues(
Row.of("Bob", 23),
Row.of("Alice", 42),
Row.of("Alice", 2))
.as("name", "score");
// Since the aggregation is not windowed and potentially unbounded,
// the result is an updating table. Usually, this means that all following
// operations and sinks need to support updates.
Table aggregated = data
.groupBy($("name"))
.select($("name"), $("score").sum().as("sum"));
// However, the PTF will convert the updating table into an insert-only result.
// Subsequent operations and sinks can easily digest the resulting table.
Table changelog = aggregated
.partitionBy($("name"))
.process(ToChangelogFunction.class);
// For event-driven applications, filtering on certain CDC events is possible.
Table upsertsOnly = changelog.filter($("flag").in("INSERT", "UPDATE_AFTER"));
upsertsOnly.execute().print();
// --------------------
// Function declaration
// --------------------
@DataTypeHint("ROW<flag STRING, sum INT>")
public static class ToChangelogFunction extends ProcessTableFunction<Row> {
public void eval(@ArgumentHint({SET_SEMANTIC_TABLE, SUPPORT_UPDATES}) Row input) {
// Forwards the sum column and includes the row's kind as a string column.
Row changelogRow =
Row.of(
input.getKind().toString(),
input.getField("sum"));
collect(changelogRow);
}
}
PTF 는 콘솔에서 디버깅할 때 다음 출력을 생성합니다. op 섹션은 결과가 append-only 임을 나타냅니다. 원래 플래그는 flag 열에 인코딩됩니다.
+----+--------------------------------+--------------------------------+-------------+
| op | name | flag | sum |
+----+--------------------------------+--------------------------------+-------------+
| +I | Bob | INSERT | 23 |
| +I | Alice | INSERT | 42 |
| +I | Alice | UPDATE_AFTER | 44 |
+----+--------------------------------+--------------------------------+-------------+
제한 사항
ArgumentTrait.PASS_COLUMNS_THROUGH는ArgumentTrait.SUPPORTS_UPDATES가 선언되면 지원되지 않습니다.- PTF 가 업데이트를 받으면
on_time인수는 지원되지 않습니다.
업데이트 함수 출력
ChangelogFunction 인터페이스는 함수가 내보낼 수 있는 변경 유형(예: insert, update, delete) 을 선언할 수 있게 해 플래너가 쿼리 플래닝 중 정보에 기반한 결정을 내릴 수 있습니다.
이 인터페이스는 고급 사용 사례를 위한 것이며 주의해서 구현해야 합니다. PTF 에서 잘못된 changelog 를 내보내면 전체 쿼리에서 정의되지 않은 동작이 발생할 수 있습니다.
결과 changelog mode 는 다음에 의해 영향을 받을 수 있습니다:
ChangelogContext.getTableChangelogMode(int)로 접근할 수 있는 입력 테이블 인수의 changelog mode.ChangelogContext.getRequiredChangelogMode()로 접근할 수 있는 다운스트림 operator 가 요구하는 changelog mode.
플래너의 changelog mode 추론은 여러 단계를 포함합니다. getChangelogMode(ChangelogContext) 메서드는 각 단계마다 호출됩니다:
- 플래너는 PTF 가 업데이트 또는 insert-only 를 내보내는지 확인합니다.
- 업데이트가 내보내지면 플래너는 업데이트에 {@link RowKind#UPDATE_BEFORE} 메시지(retract mode) 가 포함되는지, 아니면 {@link RowKind#UPDATE_AFTER} 메시지(upsert mode) 로 충분한지 결정합니다. 이를 위해 {@link #getChangelogMode} 는 {@link ChangelogContext#getRequiredChangelogMode()} 가 나타내는 retract mode 및 upsert mode 능력을 모두 질의하기 위해 두 번 호출될 수 있습니다.
- upsert mode 에서 플래너는 {@link RowKind#DELETE} 메시지가 모든 필드(full deletes) 또는 키 필드만(partial deletes) 포함하는지 확인합니다. partial deletes 의 경우 행이 제거될 때 upsert 키 필드만 설정되며, 모든 비키 필드는 nullability 제약과 관계없이 null 입니다. {@link ChangelogContext#getRequiredChangelogMode()} 는 다운스트림 operator 가 full deletes 를 요구하는지 나타냅니다.
changelog 를 내보내는 것은 집합 의미론(ArgumentTrait.SET_SEMANTIC_TABLE 참고) 을 가진 테이블 인수를 받는 PTF 에만 유효합니다. upsert 의 경우 upsert 키는 PARTITION BY 키와 같아야 합니다.
ChangelogFunction 구현이 ChangelogContext 와 관계없이 고정된 ChangelogMode 를 반환하는 것은 완전히 유효합니다. 이 접근 방식은 PTF 가 특정 시나리오나 파이프라인 설정을 위해 설계되었고 입력 모드에 동적으로 적응할 필요가 없을 때 적절할 수 있습니다. 이 경우 PTF 의 적용성이 제한되어 설계된 사전 정의된 컨텍스트 내에서만 올바르게 작동할 수 있음에 유의하세요.
일부 경우 이 인터페이스는 특정 호출 위치의 최종 changelog mode 가 결정된 후 PTF 를 재구성하기 위해 SpecializedFunction 과 함께 사용되어야 합니다. 최종 changelog mode 는 런타임에 ProcessTableFunction.Context.getChangelogMode() 로도 사용할 수 있습니다.
예제: Custom Aggregation
다음 함수는 PTF 가 사용자 정의 조건 로직에 기반한 업데이트를 내보낼 수 있는 집계 함수를 구현할 수 있는 방법을 보여줍니다. 이 함수는 name 으로 파티셔닝된 score 결과 테이블을 받고 파티션당 합계를 유지합니다. 0 보다 낮은 score 는 잘못된 것으로 취급되어 이 키에 대한 전체 집계를 무효화합니다.
Java:
TableEnvironment env = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
Table data = env
.fromValues(
Row.of("Bob", 23),
Row.of("Alice", 42),
Row.of("Alice", 2),
Row.of("Bob", -1),
Row.of("Bob", 45))
.as("name", "score");
// Call the PTF on an append-only table
Table aggregated = data
.partitionBy($("name"))
.process(CustomAggregation.class);
aggregated.execute().print();
// --------------------
// Function declaration
// --------------------
@DataTypeHint("ROW<sum INT>")
public static class CustomAggregation
extends ProcessTableFunction<Row>
implements ChangelogFunction {
@Override
public ChangelogMode getChangelogMode(ChangelogContext changelogContext) {
// Tells the system that the PTF produces updates encoded as retractions
return ChangelogMode.all();
}
public static class Accumulator {
public Integer sum = 0;
}
public void eval(@StateHint Accumulator state, @ArgumentHint(SET_SEMANTIC_TABLE) Row input) {
int score = input.getFieldAs("score");
// A negative state indicates that the partition
// key has been marked as invalid before
if (state.sum == -1) {
return;
}
// A negative score marks the entire aggregation result as invalid.
if (score < 0) {
// Send out a -D for the affected partition key and
// mark the invalidation in state. All subsequent operations
// and sinks will remove the aggregation result.
collect(Row.ofKind(RowKind.DELETE, state.sum));
state.sum = -1;
} else {
if (state.sum == 0) {
// Emit +I for the first valid aggregation result.
state.sum += score;
collect(Row.ofKind(RowKind.INSERT, state.sum));
} else {
// Emit -U (with old aggregation result) and +U (with new aggregation result)
// for encoding the update.
collect(Row.ofKind(RowKind.UPDATE_BEFORE, state.sum));
state.sum += score;
collect(Row.ofKind(RowKind.UPDATE_AFTER, state.sum));
}
}
}
}
PTF 는 콘솔에서 디버깅할 때 다음 출력을 생성합니다. op 섹션은 결과가 업데이트 중임을 나타냅니다. 그러나 Bob 에 대한 잘못된 -1 이 수신된 후에는 Bob 에 대한 업데이트가 전달되지 않아 PTF 가 값 45 의 업데이트를 무시합니다. Alice 에 대한 집계 결과는 유효한 score 만 포함하며 materialized table 에 보존됩니다.
+----+--------------------------------+-------------+
| op | name | sum |
+----+--------------------------------+-------------+
| +I | Bob | 23 |
| +I | Alice | 42 |
| -U | Alice | 42 |
| +U | Alice | 44 |
| -D | Bob | 23 |
+----+--------------------------------+-------------+
제한 사항
- PTF 가 업데이트를 내보내면
on_time인수는 지원되지 않습니다. - 현재 upsert PTF 를 테스트하는 것은 어렵습니다.
collect()같은 디버깅 sink 는 retract mode 로 동작하고 어떤 PTF 에서든 retract 지원을 요청하기 때문입니다. Upserting PTF 는 upserting sink(예:kafka-upsert커넥터) 로 테스트해야 합니다.
고급 예제 (Advanced Examples)
쇼핑 카트 (Shopping Cart)
다음 예제는 쇼핑 카트를 모델링하는 전형적인 PTF 사용 사례를 보여줍니다. 서로 다른 이벤트가 카트의 내용에 영향을 줍니다. 이 예제에서 각 사용자는 아이템을 ADD 하거나 REMOVE 할 수 있습니다. 성공적인 경우 사용자는 CHECKOUT 으로 거래를 완료합니다.
다음 입력 테이블을 생각해 보세요:
+----+--------------------------------+--------------------------------+----------------------+-------------------------+
| op | user | eventType | productId | ts |
+----+--------------------------------+--------------------------------+----------------------+-------------------------+
| +I | Bob | ADD | 1 | 2025-03-27 12:00:11.000 |
| +I | Alice | ADD | 1 | 2025-03-27 12:00:21.000 |
| +I | Bob | REMOVE | 1 | 2025-03-27 12:00:51.000 |
| +I | Bob | ADD | 2 | 2025-03-27 12:00:55.000 |
| +I | Bob | ADD | 5 | 2025-03-27 12:00:56.000 |
| +I | Bob | CHECKOUT | <NULL> | 2025-03-27 12:01:50.000 |
CheckoutProcessor PTF 는 이러한 이벤트를 처리하고 체크아웃이 완료될 때까지 쇼핑 카트 내용을 상태에 저장하도록 설계되었습니다. 또한 알림(reminder) 과 타임아웃 로직을 통합합니다. 사용자가 지정된 기간 동안 비활성이면 현재 카트 내용이 있는 REMINDER 이벤트가 내보내집니다. CHECKOUT 이벤트를 수신하면 PTF 가 정리되고 체크아웃 이벤트가 전송됩니다.
Java:
// Function that implements the core business logic of a shopping cart.
// The PTF takes ADD, REMOVE, CHECKOUT events and can send out either a REMINDER or CHECKOUT event.
@DataTypeHint("ROW<checkout_type STRING, items MAP<BIGINT, INT>>")
public static class CheckoutProcessor extends ProcessTableFunction<Row> {
// Object that is stored in state.
public static class ShoppingCart {
// The system needs to be able to access all fields for persistence and restore.
// A map for product IDs to number of items.
public Map<Long, Integer> content = new HashMap<>();
// Arbitrary helper methods can be added for structuring the code.
public void addItem(long productId) {
content.compute(productId, (k, v) -> (v == null) ? 1 : v + 1);
}
public void removeItem(long productId) {
content.compute(productId, (k, v) -> (v == null || v == 1) ? null : v - 1);
}
public boolean hasContent() {
return !content.isEmpty();
}
}
// Main processing logic
public void eval(
Context ctx,
@StateHint ShoppingCart cart,
@ArgumentHint({SET_SEMANTIC_TABLE, REQUIRE_ON_TIME}) Row events,
Duration reminderInterval,
Duration timeoutInterval
) {
String eventType = events.getFieldAs("eventType");
Long productId = events.getFieldAs("productId");
switch (eventType) {
// ADD item
case "ADD":
cart.addItem(productId);
updateTimers(ctx, reminderInterval, timeoutInterval);
break;
// REMOVE item
case "REMOVE":
cart.removeItem(productId);
if (cart.hasContent()) {
updateTimers(ctx, reminderInterval, timeoutInterval);
} else {
ctx.clearAll();
}
break;
// CHECKOUT process
case "CHECKOUT":
if (cart.hasContent()) {
collect(Row.of("CHECKOUT", cart.content));
}
ctx.clearAll();
break;
}
}
// Executes REMINDER and TIMEOUT events
public void onTimer(OnTimerContext ctx, ShoppingCart cart) {
switch (ctx.currentTimer()) {
// Send reminder event
case "REMINDER":
collect(Row.of("REMINDER", cart.content));
break;
// Cancel transaction
case "TIMEOUT":
ctx.clearAll();
break;
}
}
// Helper method that sets or replaces timers for REMINDER and TIMEOUT
private void updateTimers(Context ctx, Duration reminderInterval, Duration timeoutInterval) {
TimeContext<Instant> timeCtx = ctx.timeContext(Instant.class);
timeCtx.registerOnTime("REMINDER", timeCtx.time().plus(reminderInterval));
timeCtx.registerOnTime("TIMEOUT", timeCtx.time().plus(timeoutInterval));
}
}
출력은 다음과 유사할 수 있습니다. 여기서는 1초라는 매우 짧은 알림 간격을 가정합니다.
+----+--------------------------------+--------------------------------+--------------------------------+-------------------------+
| op | user | checkout_type | items | rowtime |
+----+--------------------------------+--------------------------------+--------------------------------+-------------------------+
| +I | Bob | REMINDER | {1=1} | 2025-03-27 12:00:12.000 |
| +I | Alice | REMINDER | {1=1} | 2025-03-27 12:00:22.000 |
| +I | Bob | CHECKOUT | {2=1, 5=1} | 2025-03-27 12:01:50.000 |
실제 시나리오에서 출력은 두 개의 별도 시스템에 분산될 가능성이 높습니다. 알림은 이메일 알림 큐에 넣을 수 있고, 체크아웃 프로세스는 별도의 다운스트림 시스템에 의해 최종화될 것입니다.
CREATE VIEW Checkouts AS SELECT * FROM CheckoutProcessor(
events => TABLE Events PARTITION BY `user`,
on_time => DESCRIPTOR(ts),
reminderInterval => INTERVAL '1' DAY,
timeoutInterval => INTERVAL '2' DAY, uid => 'cart-processor'
)
EXECUTE STATEMENT SET
BEGIN
INSERT INTO EmailNotifications SELECT * FROM Checkouts WHERE `checkout_type` = 'REMINDER';
INSERT INTO CheckoutEvents SELECT * FROM Checkouts WHERE `checkout_type` = 'CHECKOUT';
END;
뷰를 정의하고 두 INSERT INTO 경로에서 같은 uid 를 사용함으로써, 결과 Flink 작업은 PTF 가 파이프라인에 한 번 존재하면서 split 동작을 사용합니다.
결제 조인 (Payment Joining)
다음 예제는 PTF 가 조인에 사용될 수 있는 방법을 보여줍니다. 또한 PTF 가 더미 데이터로 bounded 테이블을 만드는 데이터 생성기로 사용될 수 있는 방법도 보여줍니다.
Java:
// ---------------------------
// Table program
// ---------------------------
TableEnvironment env = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
// Generate data with a unified schema
Table orders = env.fromCall(OrderGenerator.class);
Table payments = env.fromCall(PaymentGenerator.class);
// Partition orders and payments and pass them into the Joiner function
Table joined = env.fromCall(
Joiner.class,
orders.partitionBy($("id")).asArgument("order"),
payments.partitionBy($("orderId")).asArgument("payment"));
joined.execute().print();
// ---------------------------
// Data Generation
// ---------------------------
// A PTF that generates Orders
public static class OrderGenerator extends ProcessTableFunction<Order> {
public void eval() {
Stream.of(
Order.of("Bob", 1000001, 23.46, "USD"),
Order.of("Bob", 1000021, 6.99, "USD"),
Order.of("Alice", 1000601, 0.79, "EUR"),
Order.of("Charly", 1000703, 100.60, "EUR")
)
.forEach(this::collect);
}
}
// A PTF that generates Payments
public static class PaymentGenerator extends ProcessTableFunction<Payment> {
public void eval() {
Stream.of(
Payment.of(999997870, 1000001),
Payment.of(999997870, 1000001),
Payment.of(999993331, 1000021),
Payment.of(999994111, 1000601)
)
.forEach(this::collect);
}
}
// Order POJO
public static class Order {
public String userId;
public int id;
public double amount;
public String currency;
public static Order of(String userId, int id, double amount, String currency) {
Order order = new Order();
order.userId = userId;
order.id = id;
order.amount = amount;
order.currency = currency;
return order;
}
}
// Payment POJO
public static class Payment {
public int id;
public int orderId;
public static Payment of(int id, int orderId) {
Payment payment = new Payment();
payment.id = id;
payment.orderId = orderId;
return payment;
}
}
데이터를 생성한 후 상태 저장 Joiner 는 일치하는 쌍을 찾을 때까지 이벤트를 버퍼링합니다. 입력 테이블의 중복은 무시됩니다.
Java:
// Function that buffers one object of each side to find exactly one join result.
// The function expects that a payment event enters within 1 hour.
// Otherwise, state is discarded using TTL.
public static class Joiner extends ProcessTableFunction<JoinResult> {
public void eval(
Context ctx,
@StateHint(ttl = "1 hour") JoinResult seen,
@ArgumentHint(SET_SEMANTIC_TABLE) Order order,
@ArgumentHint(SET_SEMANTIC_TABLE) Payment payment
) {
if (order != null) {
if (seen.order != null) {
// skip duplicates
return;
} else {
// wait for matching payment
seen.order = order;
}
} else if (payment != null) {
if (seen.payment != null) {
// skip duplicates
return;
} else {
// wait for matching order
seen.payment = payment;
}
}
if (seen.order != null && seen.payment != null) {
// Send out the final join result
collect(seen);
}
}
}
// POJO for the output of Joiner
public static class JoinResult {
public Order order;
public Payment payment;
}
출력은 다음과 유사할 수 있습니다. 결제 999997870 에 대한 중복 이벤트가 필터링되었습니다. Charly 에 대한 일치는 찾을 수 없었습니다.
+----+-------------+-------------+--------------------------------+--------------------------------+
| op | id | orderId | order | payment |
+----+-------------+-------------+--------------------------------+--------------------------------+
| +I | 1000021 | 1000021 | (amount=6.99, currency=USD,... | (id=999993331, orderId=1000... |
| +I | 1000601 | 1000601 | (amount=0.79, currency=EUR,... | (id=999994111, orderId=1000... |
| +I | 1000001 | 1000001 | (amount=23.46, currency=USD... | (id=999997870, orderId=1000... |
+----+-------------+-------------+--------------------------------+--------------------------------+
제한 사항 (Limitations)
PTF 는 초기 단계에 있습니다. 다음 제한 사항이 적용됩니다:
- PTF 는 배치 모드에서 실행될 수 없습니다.
- Broadcast state (브로드캐스트 상태)