Temporal Table Function
Temporal Table Function (시간 테이블 함수)
Temporal table function은 특정 시점에서의 시간 테이블(temporal table) 버전에 접근할 수 있게 해줍니다. 시간 테이블의 데이터에 접근하려면 반환될 테이블의 버전을 결정하는 시간 속성(time attribute)을 전달해야 합니다. Flink는 이를 표현하기 위해 테이블 함수(table functions)의 SQL 문법을 사용합니다.
출처: 문서
본문
시간 테이블 함수는 특정 시점에서 시간 테이블의 버전에 접근할 수 있게 해줍니다. 시간 테이블의 데이터에 접근하려면 반환될 테이블의 버전을 결정하는 시간 속성을 전달해야 합니다. Flink는 테이블 함수의 SQL 문법을 사용해 이를 표현합니다.
버전화 테이블(versioned table)과 달리, temporal table function은 append-only 스트림 위에서만 정의할 수 있습니다. 즉 changelog 입력을 지원하지 않습니다. 또한 temporal table function은 순수 SQL DDL로는 정의할 수 없습니다.
Temporal Table Function 정의하기
Temporal table function은 Table API를 사용해 append-only 스트림 위에서 정의할 수 있습니다. 테이블은 하나 이상의 키 컬럼과 버전화에 사용되는 시간 속성으로 등록됩니다.
버전화를 위해 temporal table function으로 등록하려는 통화 환율(currency rates)의 append-only 테이블이 있다고 가정해 봅시다.
SELECT * FROM currency_rates;
update_time currency rate
============= ========= ====
09:00:00 Yen 102
09:00:00 Euro 114
09:00:00 USD 1
11:15:00 Euro 119
11:49:00 Pounds 108
Table API를 사용해 키로 currency를, 버전화 시간 속성으로 update_time을 사용해 이 스트림을 등록할 수 있습니다.
Java
TemporalTableFunction rates = tEnv
.from("currency_rates")
.createTemporalTableFunction("update_time", "currency");
tEnv.createTemporarySystemFunction("rates", rates);
Scala
rates = tEnv
.from("currency_rates")
.createTemporalTableFunction("update_time", "currency")
tEnv.createTemporarySystemFunction("rates", rates)
Python
Still not supported in Python API.
Temporal Table Function 조인 (Temporal Table Function Join)
정의된 후에는 temporal table function이 표준 테이블 함수로 사용됩니다. Append-only 테이블(왼쪽 입력/probe 쪽)은 시간에 따라 변하고 그 변경을 추적하는 시간 테이블(오른쪽 입력/build 쪽)과 조인할 수 있습니다. 이를 통해 특정 시점에 존재했던 키의 값을 검색할 수 있습니다.
고객의 서로 다른 통화 주문을 추적하는 append-only 테이블 orders를 고려해 봅시다.
SELECT * FROM orders;
order_time amount currency
========== ====== =========
10:15 2 Euro
10:30 1 USD
10:32 50 Yen
10:52 3 Euro
11:04 5 USD
이 테이블들이 주어졌을 때, 주문을 공통 통화인 USD로 변환하고 싶습니다.
SQL
SELECT
SUM(amount * rate) AS amount
FROM
orders,
LATERAL TABLE (rates(order_time))
WHERE
rates.currency = orders.currency
Java
Table result = orders
.joinLateral(call("rates", $("o_proctime")), $("o_currency").isEqual($("r_currency")))
.select($("(o_amount").times($("r_rate")).sum().as("amount"));
Scala
val result = orders
.joinLateral($"rates(order_time)", $"orders.currency = rates.currency")
.select($"((o_amount * r_rate).sum as amount)")
Python
Still not supported in Python API.