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_BEFOREUPDATE_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")
);

더 알아보기 (Learn more)