버전 테이블

버전 테이블 (Versioned Tables)

Flink SQL은 진화하는 동적 테이블(dynamic tables) 위에서 동작하며, 이는 append-only이거나 updating일 수 있어요. 버전 테이블(versioned tables)은 각 key의 과거 값을 기억하는 특별한 유형의 updating 테이블을 나타내요.

출처: 문서

본문

개념 (Concept)

동적 테이블은 시간에 따른 관계(relations)를 정의해요. 특히 메타데이터를 다룰 때는 key의 이전 값이 변경될 때 관련이 없어지지 않는 경우가 많아요.

Flink SQL은 PRIMARY KEY 제약과 시간 속성(time attribute)이 있는 모든 동적 테이블 위에 버전 테이블을 정의할 수 있어요. Flink에서 기본 키 제약은 테이블 또는 뷰의 컬럼 또는 컬럼 집합이 고유하고 null이 아님을 의미해요.

upserting 테이블의 기본 키 의미론은 특정 key에 대한 구체화된 변경(INSERT/UPDATE/DELETE)이 시간에 따른 단일 행에 대한 변경을 나타낸다는 뜻이에요. upserting 테이블의 시간 속성은 각 변경이 언제 발생했는지 정의해요.

함께 취하면, Flink는 시간에 따른 행의 변경을 추적하고 각 값이 그 key에 대해 유효했던 기간을 유지할 수 있어요.

가게의 여러 제품 가격을 추적하는 테이블을 가정해 봅시다.

(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 

이 변경 집합이 주어지면 스쿠터 가격이 시간에 따라 어떻게 변하는지 추적할 수 있어요. 처음에는 카탈로그에 추가되었을 때 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 table 문이 PRIMARY KEY와 event-time 속성을 포함해야 한다는 것이에요.

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)

기본 쿼리가 고유 key 제약과 event-time 속성을 포함한다면 Flink는 버전 뷰의 정의도 지원해요. 통화 환율의 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)