Changelog Conversion
Changelog Conversion (체인지로그 변환)
Flink SQL은 changelog 스트림을 다루기 위한 내장 프로세스 테이블 함수(PTF)를 제공합니다. FROM_CHANGELOG는 연산 코드가 있는 append-only 테이블을 (잠재적으로 갱신되는) 동적 테이블로, TO_CHANGELOG는 동적 테이블을 명시적 연산 코드가 있는 append-only 테이블로 변환합니다. 이를 통해 CDC(Change Data Capture) 스트림을 유연하게 소비하고 다시 물리화할 수 있습니다.
출처: 문서
본문
Flink SQL은 changelog 스트림을 다루기 위한 내장 프로세스 테이블 함수(PTF)를 제공합니다.
| 함수 | 설명 |
|---|---|
| FROM_CHANGELOG | 연산 코드가 있는 append-only 테이블을 (잠재적으로 갱신되는) 동적 테이블로 변환 |
| TO_CHANGELOG | 동적 테이블을 명시적 연산 코드가 있는 append-only 테이블로 변환 |
FROM_CHANGELOG
FROM_CHANGELOG PTF는 명시적 연산 코드 컬럼이 있는 append-only 테이블을 (잠재적으로 갱신되는) 동적 테이블로 변환합니다. 각 입력 행에는 변경 연산을 나타내는 문자열 컬럼이 있어야 합니다. 연산 컬럼은 엔진이 해석하고 출력에서 제거됩니다.
이는 Debezium 같은 시스템의 CDC 스트림을 소비할 때 유용합니다. 이러한 시스템에서는 이벤트가 명시적 연산 필드가 있는 평평한(flat) append-only 레코드로 도착합니다. 또한 이벤트에 특정 변환을 적용한 후 append-only 테이블을 다시 갱신 테이블로 변환할 때 TO_CHANGELOG 함수와 함께 사용하는 데도 유용합니다.
참고: 이 버전은 CDC 데이터가 업데이트를 전체 이미지(full image)로 인코딩(즉, 업데이트 전과 후에 대해 별도의 이벤트 제공)해야 합니다. 소스가
UPDATE_BEFORE와UPDATE_AFTER이벤트를 모두 제공하는지 확인하세요.FROM_CHANGELOG는 매우 강력한 함수이지만 올바르게 구성되지 않으면 이후 연산과 테이블에서 잘못된 결과를 생성할 수 있습니다.
구문 (Syntax)
SELECT * FROM FROM_CHANGELOG(
input => TABLE source_table,
[op => DESCRIPTOR(op_column_name),]
[op_mapping => MAP[
'c, r', 'INSERT',
'ub', 'UPDATE_BEFORE',
'ua', 'UPDATE_AFTER',
'd', 'DELETE'
]]
)
파라미터 (Parameters)
| 파라미터 | 필수 | 설명 |
|---|---|---|
input |
예 | 입력 테이블. append-only여야 합니다. |
op |
아니오 | 연산 코드 컬럼의 단일 컬럼 이름을 가진 DESCRIPTOR. 기본값은 op입니다. 이 컬럼은 입력 테이블에 존재해야 하며 타입이 STRING이어야 합니다. 컬럼은 nullable로 선언될 수 있지만, 런타임에 NULL 값이면 TableRuntimeException으로 작업이 실패합니다 — 모든 changelog 행은 연산 코드를 가져야 합니다. |
op_mapping |
아니오 | 사용자 정의 코드를 Flink 변경 연산 이름에 매핑하는 MAP<STRING, STRING>. 키는 사용자 정의 코드(예: 'c', 'u', 'd')이고, 값은 Flink 변경 연산 이름(INSERT, UPDATE_BEFORE, UPDATE_AFTER, DELETE)입니다. 키는 여러 코드를 같은 연산으로 매핑하기 위해 쉼표로 구분된 코드를 포함할 수 있습니다(예: 'c, r'). 매핑에 없는 op 코드를 받으면 런타임에 TableRuntimeException으로 작업이 실패합니다. 각 변경 연산은 모든 항목에 걸쳐 최대 한 번만 나타날 수 있습니다. |
기본 op_mapping
op_mapping을 생략하면 다음 표준 이름이 사용됩니다. 이들은 기본적으로 TO_CHANGELOG로부터의 역변환을 허용합니다.
| 입력 코드 | 변경 연산 |
|---|---|
'INSERT' |
INSERT |
'UPDATE_BEFORE' |
UPDATE_BEFORE |
'UPDATE_AFTER' |
UPDATE_AFTER |
'DELETE' |
DELETE |
활성 매핑(기본 또는 사용자 정의)에 없는 op 코드를 가진 입력 행은 런타임에 TableRuntimeException으로 작업이 실패합니다.
출력 스키마 (Output Schema)
출력은 연산 코드(예: op) 컬럼을 제외한 모든 입력 컬럼을 포함하며, 연산 코드 컬럼은 Flink SQL 엔진이 해석하고 제거합니다. 각 출력 행은 적절한 변경 연산(INSERT, UPDATE_BEFORE, UPDATE_AFTER, DELETE)을 가집니다.
[all_input_columns_without_op]
예시 (Examples)
표준 op 이름과의 기본 사용
-- Input (append-only):
-- +I[id:1, op:'INSERT', name:'Alice']
-- +I[id:2, op:'INSERT', name:'Bob']
-- +I[id:1, op:'UPDATE_BEFORE', name:'Alice']
-- +I[id:1, op:'UPDATE_AFTER', name:'Alice2']
-- +I[id:2, op:'DELETE', name:'Bob']
SELECT * FROM FROM_CHANGELOG(
input => TABLE cdc_stream
)
-- Output (updating table):
-- +I[id:1, name:'Alice']
-- +I[id:2, name:'Bob']
-- -U[id:1, name:'Alice']
-- +U[id:1, name:'Alice2']
-- -D[id:2, name:'Bob']
-- Table state after all events:
-- | id | name |
-- |----|--------|
-- | 1 | Alice2 |
사용자 정의 연산 컬럼 이름
-- Source schema: id INT, operation STRING, name STRING
SELECT * FROM FROM_CHANGELOG(
input => TABLE cdc_stream,
op => DESCRIPTOR(operation)
)
-- The operation column named 'operation' is used instead of 'op'
Table API
Table cdcStream = ...;
// Default: reads 'op' column with standard change operation names
Table result = cdcStream.fromChangelog();
// With custom op column name
Table result = cdcStream.fromChangelog(
descriptor("operation").asArgument("op")
);
// With custom op_mapping
Table result = cdcStream.fromChangelog(
descriptor("op").asArgument("op"),
map("c, r", "INSERT",
"ub", "UPDATE_BEFORE",
"ua", "UPDATE_AFTER",
"d", "DELETE").asArgument("op_mapping")
);
TO_CHANGELOG
TO_CHANGELOG PTF는 동적 테이블(즉, 갱신 테이블)을 명시적 연산 코드 컬럼이 있는 append-only 테이블로 변환합니다. 각 입력 행은 원래 변경 연산(INSERT, UPDATE_BEFORE, UPDATE_AFTER, DELETE)에 관계없이 원래 연산을 나타내는 문자열 컬럼이 있는 INSERT-only 행으로 내보내집니다.
이는 changelog 이벤트를 append만 지원하는 다운스트림 시스템(예: 메시지 큐, 로그 저장소, append-only 파일 싱크)에 물리화해야 할 때 유용합니다. 또한 예를 들어 DELETE와 같은 특정 유형의 업데이트를 걸러내는 데도 유용합니다.
구문 (Syntax)
SELECT * FROM TO_CHANGELOG(
input => TABLE source_table,
[op => DESCRIPTOR(op_column_name),]
[op_mapping => MAP['INSERT', 'I', 'DELETE', 'D', ...]]
)
파라미터 (Parameters)
| 파라미터 | 필수 | 설명 |
|---|---|---|
input |
예 | 입력 테이블. insert-only, retract, upsert 테이블을 허용합니다. |
op |
아니오 | 연산 코드 컬럼의 단일 컬럼 이름을 가진 DESCRIPTOR. 기본값은 op입니다. |
op_mapping |
아니오 | 변경 연산 이름을 사용자 정의 출력 코드에 매핑하는 MAP<STRING, STRING>. 키는 여러 연산을 같은 코드로 매핑하기 위해 쉼표로 구분된 이름을 포함할 수 있습니다(예: 'INSERT, UPDATE_AFTER'). 제공되면 매핑된 연산만 전달됩니다 — 매핑되지 않은 이벤트는 버려집니다. 각 변경 연산은 모든 항목에 걸쳐 최대 한 번만 나타날 수 있습니다. |
기본 op_mapping
op_mapping을 생략하면 네 가지 변경 연산 모두 표준 이름으로 매핑됩니다.
| 변경 연산 | 출력 값 |
|---|---|
| INSERT | 'INSERT' |
| UPDATE_BEFORE | 'UPDATE_BEFORE' |
| UPDATE_AFTER | 'UPDATE_AFTER' |
| DELETE | 'DELETE' |
출력 스키마 (Output Schema)
출력 컬럼은 다음과 같이 정렬됩니다.
[op_column, all_input_columns]
모든 출력 행은 INSERT입니다 — 테이블은 항상 append-only입니다.
예시 (Examples)
기본 사용
-- Input: retract table from an aggregation
-- +I[name:'Alice', cnt:1]
-- +U[name:'Alice', cnt:2]
-- -D[name:'Bob', cnt:1]
SELECT * FROM TO_CHANGELOG(
input => TABLE my_aggregation
)
-- Output (append-only):
-- +I[op:'INSERT', name:'Alice', cnt:1]
-- +I[op:'UPDATE_AFTER', name:'Alice', cnt:2]
-- +I[op:'DELETE', name:'Bob', cnt:1]
사용자 정의 연산 컬럼 이름
SELECT * FROM TO_CHANGELOG(
input => TABLE my_aggregation,
op => DESCRIPTOR(operation)
)
-- The op column is now named 'operation' instead of 'op'
필터링이 있는 사용자 정의 연산 코드
SELECT * FROM TO_CHANGELOG(
input => TABLE my_aggregation,
op => DESCRIPTOR(op_code),
op_mapping => MAP['INSERT', 'I', 'UPDATE_AFTER', 'U']
)
-- Only INSERT and UPDATE_AFTER events are forwarded
-- DELETE events are dropped (not in the mapping)
-- op_code values are 'I' and 'U' instead of full names
삭제 플래그 패턴
SELECT * FROM TO_CHANGELOG(
input => TABLE my_aggregation,
op => DESCRIPTOR(deleted),
op_mapping => MAP['INSERT, UPDATE_AFTER', 'false', 'DELETE', 'true']
)
-- INSERT and UPDATE_AFTER produce deleted='false'
-- DELETE produces deleted='true'
-- UPDATE_BEFORE is dropped (not in the mapping)
Table API
// Default: adds 'op' column and supports all changelog modes
Table result = myTable.toChangelog();
// With custom parameters
Table result = myTable.toChangelog(
descriptor("op_code").asArgument("op"),
map("INSERT", "I", "UPDATE_AFTER", "U").asArgument("op_mapping")
);
// Deletion flag pattern: comma-separated keys map multiple change operations to the same code
Table result = myTable.toChangelog(
descriptor("deleted").asArgument("op"),
map("INSERT, UPDATE_AFTER", "false", "DELETE", "true").asArgument("op_mapping")
);