버전 테이블

버전 테이블 (Versioned Tables)

Flink SQL은 진화하는 동적 테이블 위에서 동작하며, 이 테이블은 append-only(추가 전용)이거나 updating(갱신형)일 수 있습니다. 버전 테이블(versioned table)은 각 키에 대해 과거의 값을 기억하는 특별한 종류의 updating 테이블입니다.

출처: 문서

본문

개념 (Concept)

동적 테이블은 시간에 따른 관계를 정의합니다. 특히 메타데이터를 다룰 때 키의 이전 값이 바뀌었다고 해서 무의미해지지는 않는 경우가 많습니다.

Flink SQL은 PRIMARY KEY 제약 조건과 시간 속성(time attribute)을 가진 모든 동적 테이블 위에 버전 테이블을 정의할 수 있습니다.

Flink에서 기본 키 제약 조건은 테이블이나 뷰의 한 컬럼 또는 컬럼 집합이 고유하고 null이 아니라는 뜻입니다. upsert 테이블에서 기본 키 의미론은 특정 키에 대한 물리화된 변경(INSERT/UPDATE/DELETE)이 시간에 따른 단일 행의 변경을 나타낸다는 의미입니다. upsert 테이블의 시간 속성은 각 변경이 언제 발생했는지를 정의합니다.

이를 종합하면 Flink는 시간에 따른 행의 변경을 추적하고 각 키에 대해 각 값이 유효했던 기간을 유지할 수 있습니다.

한 상점에서 서로 다른 제품의 가격을 추적하는 테이블을 가정해 보겠습니다.

(changelog kind)  update_time  product_id product_name price
================= ===========  ========== ============ =====
+(INSERT)         00:01:00     p_001      scooter      11.11
+(INSERT)         00:02:00     p_002      basketball   23.11
-(UPDATE_BEFORE)  12:00:00     p_001      scooter      11.11
+(UPDATE_AFTER)   12:00:00     p_001      scooter      12.99
-(UPDATE_BEFORE)  12:00:00     p_002      basketball   23.11
+(UPDATE_AFTER)   12:00:00     p_002      basketball   19.99
-(DELETE)         18:00:00     p_001      scooter      12.99

이 변경 집합이 주어지면 스쿠터(scooter)의 가격이 시간에 따라 어떻게 변하는지 추적합니다. 처음 00:01:00에 카탈로그에 추가되었을 때 가격은 $11.11입니다. 이후 12:00:00에 $12.99로 오르고 18:00:00에 카탈로그에서 삭제됩니다.

서로 다른 시간에 테이블에서 다양한 제품의 가격을 조회하면 서로 다른 결과를 얻게 됩니다. 10:00:00에는 한 집합의 가격이 보입니다:

update_time  product_id product_name price
===========  ========== ============ =====
00:01:00     p_001      scooter      11.11
00:02:00     p_002      basketball   23.11

반면 13:00:00에는 또 다른 집합의 가격이 보입니다:

update_time  product_id product_name price
===========  ========== ============ =====
12:00:00     p_001      scooter      12.99
12:00:00     p_002      basketball   19.99

버전 테이블 소스 (Versioned Table Sources)

버전 테이블은 기본 소스나 포맷이 changelog를 직접 정의하는 모든 테이블에 대해 암시적으로 정의됩니다. 예로는 upsert Kafka 소스와 debezium, canal 같은 데이터베이스 changelog 포맷이 있습니다. 앞서 논의한 대로 유일하게 추가되는 요구 사항은 CREATE 테이블 문이 PRIMARY KEY와 이벤트-시간 속성을 포함해야 한다는 것입니다.

CREATE TABLE products (
	product_id    STRING,
	product_name  STRING,
	price         DECIMAL(32, 2),
	update_time   TIMESTAMP(3) METADATA FROM 'value.source.timestamp' VIRTUAL,
	PRIMARY KEY (product_id) NOT ENFORCED,
	WATERMARK FOR update_time AS update_time
) WITH (...);

버전 테이블 뷰 (Versioned Table Views)

기본 쿼리에 고유 키 제약 조건과 이벤트-시간 속성이 포함되어 있으면 Flink는 버전 뷰(versioned view)를 정의하는 것도 지원합니다. 통화 환율의 append-only 테이블을 상상해 보세요.

CREATE TABLE currency_rates (
	currency      STRING,
	rate          DECIMAL(32, 10),
	update_time   TIMESTAMP(3),
	WATERMARK FOR update_time AS update_time
) WITH (
	'connector' = 'kafka',
	'topic'	    = 'rates',
	'properties.bootstrap.servers' = 'localhost:9092',
	'format'    = 'json'
);

currency_rates 테이블은 USD를 기준으로 각 통화에 대한 행을 포함하며, 환율이 변할 때마다 새 행을 수신합니다. JSON 포맷은 네이티브 changelog 의미론을 지원하지 않으므로 Flink는 이 테이블을 append-only로만 읽을 수 있습니다.

(changelog kind) update_time   currency   rate
================ ============= =========  ====
+(INSERT)        09:00:00      Yen        102
+(INSERT)        09:00:00      Euro       114
+(INSERT)        09:00:00      USD        1
+(INSERT)        11:15:00      Euro       119
+(INSERT)        11:45:00      Pounds     107
+(INSERT)        11:49:00      Pounds     108

Flink는 각 행을 테이블에 대한 INSERT로 해석하므로 currency에 PRIMARY KEY를 정의할 수 없습니다. 그러나 쿼리 개발자인 우리에게 이 테이블에는 버전 테이블을 정의하는 데 필요한 모든 정보가 들어 있다는 것은 분명합니다. Flink는 추론된 기본 키(currency)와 이벤트 시간(update_time)을 가진 정렬된 changelog 스트림을 생성하는 중복 제거(deduplication) 쿼리를 정의함으로써 이 테이블을 버전 테이블로 재해석할 수 있습니다.

-- Define a versioned view
CREATE VIEW versioned_rates AS
SELECT currency, rate, update_time              -- (1) `update_time` keeps the event time
  FROM (
      SELECT *,
      ROW_NUMBER() OVER (PARTITION BY currency  -- (2) the inferred unique key `currency` can be a primary key
         ORDER BY update_time DESC) AS rownum
      FROM currency_rates)
WHERE rownum = 1;

-- the view `versioned_rates` will produce a changelog as the following.
(changelog kind) update_time currency   rate
================ ============= =========  ====
+(INSERT)        09:00:00      Yen        102
+(INSERT)        09:00:00      Euro       114
+(INSERT)        09:00:00      USD        1
+(UPDATE_AFTER)  11:15:00      Euro       119
+(INSERT)        11:45:00      Pounds     107
+(UPDATE_AFTER)  11:49:00      Pounds     108

Flink는 이 쿼리를 효율적으로 후속 쿼리에서 사용 가능한 버전 테이블로 변환하는 특별한 최적화 단계를 가지고 있습니다. 일반적으로 다음 형식의 쿼리 결과는 버전 테이블을 생성합니다:

SELECT [column_list]
FROM table_name
QUALIFY ROW_NUMBER() OVER ([PARTITION BY col1[, col2...]] ORDER BY time_attr [asc|desc]) = 1

파라미터 설명:

  • ROW_NUMBER(): 각 행에 1부터 시작하는 고유한 순차 번호를 할당합니다.
  • PARTITION BY col1[, col2...]: 파티션 컬럼, 즉 중복 제거 키를 지정합니다. 이 컬럼들이 후속 버전 테이블의 기본 키를 형성합니다.
  • ORDER BY time_attr DESC: 순서를 지정하는 컬럼으로, 반드시 시간 속성이어야 합니다.
  • WHERE rownum = 1: Flink가 이 쿼리가 버전 테이블을 생성하는 것으로 인식하려면 rownum = 1이 필요합니다.

더 알아보기 (Learn more)