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.

더 알아보기 (Learn more)