조인
조인 (Joins)
Flink SQL은 동적 테이블(dynamic tables)에 대한 복잡하고 유연한 조인 연산을 지원합니다. 쿼리가 요구할 수 있는 다양한 의미론을 처리하기 위해 여러 유형의 조인이 있습니다.
기본적으로 조인의 순서는 최적화되지 않습니다. 테이블은 FROM 절에 지정된 순서대로 조인됩니다. 갱신 빈도가 가장 낮은 테이블을 먼저, 갱신 빈도가 가장 높은 테이블을 마지막에 나열하여 조인 쿼리 성능을 조정할 수 있습니다. 지원되지 않으며 쿼리 실패를 유발하는 교차 조인(Cartesian product)이 나오지 않는 순서로 테이블을 지정해야 합니다.
출처: 문서
본문
일반 조인 (Regular Joins)
Regular joins은 가장 일반적인 유형의 조인으로, 조인의 어느 한쪽에 대한 새로운 레코드나 변경 사항이 모두 보이고 전체 조인 결과에 영향을 미칩니다. 예를 들어 왼쪽에 새 레코드가 있으면 제품 id가 같을 때 오른쪽의 모든 이전 및 이후 레코드와 조인됩니다.
SELECT * FROM Orders
INNER JOIN Product
ON Orders.productId = Product.id
스트리밍 쿼리의 경우 일반 조인의 문법은 가장 유연하며 모든 종류의 갱신(insert, update, delete) 입력 테이블을 허용합니다. 그러나 이 연산은 중요한 운영상의 의미를 가집니다: 조인 입력의 양쪽을 Flink 상태에 영원히 유지해야 합니다. 따라서 쿼리 결과 계산에 필요한 상태는 모든 입력 테이블과 중간 조인 결과의 고유 입력 행 수에 따라 무한히 커질 수 있습니다. 과도한 상태 크기를 방지하려면 적절한 상태 time-to-live(TTL)가 있는 쿼리 구성을 제공할 수 있습니다. 이는 쿼리 결과의 정확성에 영향을 줄 수 있습니다. 자세한 내용은 query configuration을 참조하세요.
스트리밍 쿼리의 경우 쿼리 결과를 계산하는 데 필요한 상태는 집계 유형과 고유 그룹 키 수에 따라 무한히 커질 수 있습니다. 과도한 상태 크기를 방지하려면 idle 상태 보존 시간을 제공하세요. 자세한 내용은 Idle State Retention Time을 참조하세요.
INNER Equi-JOIN
조인 조건으로 제한된 단순 Cartesian product를 반환합니다. 현재는 equi-join만 지원됩니다. 즉 동등 술어가 있는 접속 조건이 하나 이상 있는 조인입니다. 임의의 cross 또는 theta 조인은 지원되지 않습니다.
SELECT *
FROM Orders
INNER JOIN Product
ON Orders.product_id = Product.id
OUTER Equi-JOIN
자격이 있는 Cartesian product의 모든 행(즉 조인 조건을 통과하는 모든 결합 행)과, 조인 조건이 다른 테이블의 어떤 행과도 일치하지 않은 외부 테이블의 각 행의 사본 하나를 반환합니다. Flink는 LEFT, RIGHT, FULL outer join을 지원합니다. 현재는 equi-join만 지원됩니다. 즉 동등 술어가 있는 접속 조건이 하나 이상 있는 조인입니다. 임의의 cross 또는 theta 조인은 지원되지 않습니다.
SELECT *
FROM Orders
LEFT JOIN Product
ON Orders.product_id = Product.id
SELECT *
FROM Orders
RIGHT JOIN Product
ON Orders.product_id = Product.id
SELECT *
FROM Orders
FULL OUTER JOIN Product
ON Orders.product_id = Product.id
다중 일반 조인
여러 연결된 조인의 성능 문제가 있고 많은 상태를 생성한다면 MultiJoin 연산자 사용을 고려하세요. 자세한 내용은 tuning multiple regular joins을 참조하세요.
Interval Joins
조인 조건과 시간 제약으로 제한된 단순 Cartesian product를 반환합니다. interval join은 하나 이상의 equi-join 술어와 양쪽의 시간을 제한하는 조인 조건이 필요합니다. 두 개의 적절한 범위 술어(<, <=, >=, >), BETWEEN 술어, 또는 두 입력 테이블의 시간 속성(즉 processing time 또는 event time)을 비교하는 단일 동등 술어가 그러한 조건을 정의할 수 있습니다.
예를 들어 이 쿼리는 주문이 접수된 후 4시간 후에 배송된 경우 모든 주문을 해당 배송(shipment)과 조인합니다.
SELECT *
FROM Orders o, Shipments s
WHERE o.id = s.order_id
AND o.order_time BETWEEN s.ship_time - INTERVAL '4' HOUR AND s.ship_time
다음 술어는 유효한 interval join 조건의 예입니다.
ltime = rtimeltime >= rtime AND ltime < rtime + INTERVAL '10' MINUTEltime BETWEEN rtime - INTERVAL '10' SECOND AND rtime + INTERVAL '5' SECOND
스트리밍 쿼리의 경우 일반 조인과 비교해 interval join은 시간 속성이 있는 append-only 테이블만 지원합니다. 시간 속성은 유사 단조 증가이므로 Flink는 결과의 정확성에 영향을 주지 않고 상태에서 오래된 값을 제거할 수 있습니다.
Temporal Joins
Temporal table은 시간에 따라 변하는 테이블입니다. Flink에서는 동적 테이블(dynamic table)이라고도 합니다. Temporal table의 행은 하나 이상의 시간 기간과 연관되며 모든 Flink 테이블은 temporal(동적)입니다. Temporal table은 하나 이상의 버전화된 테이블 스냅샷을 포함하며, 변경을 추적하는 변경 이력 테이블(예: 모든 스냅샷을 포함하는 데이터베이스 changelog)이거나 변경을 구체화하는 변경 차원 테이블(예: 최신 스냅샷을 포함하는 데이터베이스 테이블)일 수 있습니다.
이벤트 시간 Temporal Join
이벤트 시간 temporal join은 버전화된 테이블(versioned table)에 대한 조인을 허용합니다. 이는 테이블이 변경되는 메타데이터로 강화(rich)될 수 있고 특정 시점에서 그 값을 검색할 수 있음을 의미합니다.
Temporal joins은 임의의 테이블(왼쪽 입력/probe 측)을 가져와 각 행을 버전화된 테이블(오른쪽 입력/build 측)의 해당 행의 관련 버전과 연관시킵니다. Flink는 SQL:2011 표준의 FOR SYSTEM_TIME AS OF 구문을 사용해 이 연산을 수행합니다. Temporal join의 구문은 다음과 같습니다.
SELECT [column_list]
FROM table1 [AS <alias1>]
[LEFT] JOIN table2 FOR SYSTEM_TIME AS OF table1.{ proctime | rowtime } [AS <alias2>]
ON table1.column-name1 = table2.column-name1
이벤트 시간 속성(즉 rowtime 속성)으로 과거의 어느 시점에 있던 키의 값을 검색할 수 있습니다. 이는 두 테이블을 공통 시점에 조인할 수 있게 합니다. 버전화된 테이블은 마지막 워터마크 이후의 모든 버전을 시간으로 식별하여 저장합니다.
예를 들어 각각 다른 통화로 된 가격이 있는 주문 테이블이 있다고 가정합니다. 이 테이블을 USD 같은 단일 통화로 적절히 정규화하려면 각 주문을 주문이 접수된 시점의 적절한 통화 환율과 조인해야 합니다.
-- Create a table of orders. This is a standard
-- append-only dynamic table.
CREATE TABLE orders (
order_id STRING,
price DECIMAL(32,2),
currency STRING,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '15' SECOND
) WITH (/* ... */);
-- Define a versioned table of currency rates.
-- This could be from a change-data-capture
-- such as Debezium, a compacted Kafka topic, or any other
-- way of defining a versioned table.
CREATE TABLE currency_rates (
currency STRING,
conversion_rate DECIMAL(32, 2),
update_time TIMESTAMP(3) METADATA FROM `values.source.timestamp` VIRTUAL,
WATERMARK FOR update_time AS update_time - INTERVAL '15' SECOND,
PRIMARY KEY(currency) NOT ENFORCED
) WITH (
'connector' = 'kafka',
'value.format' = 'debezium-json',
/* ... */
);
SELECT
order_id,
price,
orders.currency,
conversion_rate,
order_time
FROM orders
LEFT JOIN currency_rates FOR SYSTEM_TIME AS OF orders.order_time
ON orders.currency = currency_rates.currency;
order_id price currency conversion_rate order_time
======== ===== ======== =============== =========
o_001 11.11 EUR 1.14 12:00:00
o_002 12.51 EUR 1.10 12:06:00
참고: 이벤트 시간 temporal join은 왼쪽과 오른쪽의 워터마크에 의해 트리거됩니다. INTERVAL 시간 뺄셈은 지연 이벤트를 기다리기 위해 사용되어 조인이 기대에 부합하도록 합니다. 조인 양쪽 모두 워터마크가 올바르게 설정되었는지 확인하세요.
참고: 이벤트 시간 temporal join은 temporal join 조건의 동등 조건에 포함된 기본 키가 필요합니다. 예: currency_rates 테이블의 기본 키 currency_rates.currency가 조건 orders.currency = currency_rates.currency에 제약되어야 합니다.
일반 조인과 대조적으로, build 측의 변경에도 불구하고 이전 temporal table 결과는 영향을 받지 않습니다. interval joins과 비교해 temporal table join은 레코드가 조인될 시간 윈도우를 정의하지 않습니다. Probe 측의 레코드는 항상 시간 속성이 지정한 시점의 build 측 버전과 조인됩니다. 따라서 build 측의 행은 임의로 오래될 수 있습니다. 시간이 지나면 (주어진 기본 키에 대해) 더 이상 필요하지 않은 레코드 버전이 상태에서 제거됩니다.
처리 시간 Temporal Join
처리 시간 temporal table join은 처리 시간 속성을 사용하여 행을 외부 버전화된 테이블의 키의 최신 버전과 연관시킵니다.
정의상 처리 시간 속성을 사용하면 조인은 항상 주어진 키에 대한 가장 최신 값을 반환합니다. Lookup 테이블을 build 측의 모든 레코드를 저장하는 단순 HashMap<K, V>로 생각할 수 있습니다. 이 조인의 힘은 Flink 내에서 테이블을 동적 테이블로 구체화하는 것이 실현 가능하지 않을 때 Flink가 외부 시스템에 직접 작동할 수 있게 해준다는 것입니다.
다음 처리 시간 temporal table join 예시는 append-only 테이블 orders가 테이블 LatestRates와 조인되어야 함을 보여줍니다. LatestRates는 최신 환율로 구체화된 차원 테이블(예: HBase 테이블)입니다. 시간 10:15, 10:30, 10:52에 LatestRates의 내용은 다음과 같습니다.
10:15> SELECT * FROM LatestRates;
currency rate
======== ======
US Dollar 102
Euro 114
Yen 1
10:30> SELECT * FROM LatestRates;
currency rate
======== ======
US Dollar 102
Euro 114
Yen 1
10:52> SELECT * FROM LatestRates;
currency rate
======== ======
US Dollar 102
Euro 116 <==== changed from 114 to 116
Yen 1
시간 10:15와 10:30의 LastestRates 내용은 동일합니다. Euro 환율은 10:52에 114에서 116으로 변경되었습니다.
Orders는 주어진 amount와 currency에 대한 지불을 나타내는 append-only 테이블입니다. 예를 들어 10:15에 2 Euro 금액의 주문이 있었습니다.
SELECT * FROM Orders;
amount currency
====== =========
2 Euro <== arrived at time 10:15
1 US Dollar <== arrived at time 10:30
2 Euro <== arrived at time 10:52
이 테이블들이 주어졌을 때 공통 통화로 변환된 모든 Orders를 계산하려고 합니다.
amount currency rate amount*rate
====== ========= ======= ============
2 Euro 114 228 <== arrived at time 10:15
1 US Dollar 102 102 <== arrived at time 10:30
2 Euro 116 232 <== arrived at time 10:52
현재 임의의 view/테이블의 최신 버전과 함께 temporal join에 사용되는 FOR SYSTEM_TIME AS OF 구문은 아직 지원되지 않습니다. 다음과 같이 temporal table function 구문을 사용할 수 있습니다.
SELECT
o_amount, r_rate
FROM
Orders,
LATERAL TABLE (Rates(o_proctime))
WHERE
r_currency = o_currency
참고 임의의 테이블/view의 최신 버전과 함께 temporal join에 사용되는 FOR SYSTEM_TIME AS OF 구문이 지원되지 않는 이유는 의미론적 고려 때문입니다. 왼쪽 스트림의 조인 처리가 temporal table의 완전한 스냅샷을 기다리지 않기 때문에 프로덕션 환경에서 사용자를 오도할 수 있습니다. temporal table function에 의한 처리 시간 temporal join도 같은 의미론적 문제가 있지만 오래 생존해 왔으므로 호환성 관점에서 지원합니다.
처리 시간의 결과는 결정적이지 않습니다. 처리 시간 temporal join은 대부분 외부 테이블(즉 차원 테이블)로 스트림을 강화하는 데 사용됩니다.
일반 조인과 대조적으로, build 측의 변경에도 불구하고 이전 temporal table 결과는 영향을 받지 않습니다. interval joins과 비교해 temporal table join은 레코드가 조인되는 시간 윈도우를 정의하지 않습니다. 즉 오래된 행이 상태에 저장되지 않습니다.
Temporal Table Function Join
테이블 함수(temporal table function)와 테이블을 조인하는 구문은 Table Function과의 조인과 동일합니다.
참고: 현재 temporal table과의 inner join과 left outer join만 지원됩니다.
Rates가 temporal table function이라고 가정하면 조인은 SQL로 다음과 같이 표현할 수 있습니다.
SELECT
o_amount, r_rate
FROM
Orders,
LATERAL TABLE (Rates(o_proctime))
WHERE
r_currency = o_currency
위 Temporal Table DDL과 Temporal Table Function의 주요 차이점은 다음과 같습니다.
- Temporal table DDL은 SQL에서 정의할 수 있지만 temporal table function은 정의할 수 없습니다.
- 둘 다 temporal join 버전화된 테이블을 지원하지만, temporal table function만 임의의 테이블/view의 최신 버전을 temporal join할 수 있습니다.
Lookup Join
Lookup join은 일반적으로 외부 시스템에서 쿼리한 데이터로 테이블을 강화하는 데 사용됩니다. 조인은 한 테이블이 처리 시간 속성을 갖고 다른 테이블이 lookup 소스 커넥터로 지원되어야 합니다.
Lookup join은 오른쪽 테이블이 lookup 소스 커넥터로 지원되는 위 Processing Time Temporal Join 구문을 사용합니다.
다음 예시는 lookup join을 지정하는 구문을 보여줍니다.
-- Customers is backed by the JDBC connector and can be used for lookup joins
CREATE TEMPORARY TABLE Customers (
id INT,
name STRING,
country STRING,
zip STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://mysqlhost:3306/customerdb',
'table-name' = 'customers'
);
-- enrich each order with customer information
SELECT o.order_id, o.total, c.country, c.zip
FROM Orders AS o
JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.customer_id = c.id;
위 예시에서 Orders 테이블은 MySQL 데이터베이스에 있는 Customers 테이블의 데이터로 강화됩니다. FOR SYSTEM_TIME AS OF 절과 그에 따른 처리 시간 속성은 Orders 테이블의 각 행이 조인 연산자에 의해 처리되는 시점에 조인 술어와 일치하는 Customers 행과 조인되도록 보장합니다. 또한 조인된 Customer 행이 나중에 갱신되어도 조인 결과가 갱신되는 것을 방지합니다. Lookup join은 또한 필수 동등 조인 술어가 필요합니다. 위 예시에서는 o.customer_id = c.id입니다.
Array, Multiset 및 Map 평면화
Unnest는 주어진 array, multiset 또는 map의 각 요소에 대해 새 행을 반환합니다. CROSS JOIN과 LEFT JOIN을 모두 지원합니다.
-- Returns a new row for each element in a constant array
SELECT * FROM (VALUES('order_1')), UNNEST(ARRAY['shirt', 'pants', 'hat'])
id product_name
======= ============
order_1 shirt
order_1 pants
order_1 hat
-- Returns a new row for each element in the array
-- assuming a Orders table with an array column `product_names`
SELECT order_id, product_name
FROM Orders
CROSS JOIN UNNEST(product_names) AS t(product_name)
WITH ORDINALITY를 사용한 unnest도 지원됩니다. 현재 WITH ORDINALITY는 CROSS JOIN만 지원하고 LEFT JOIN은 지원하지 않습니다.
-- Returns a new row for each element in a constant array and its position in the array
SELECT *
FROM (VALUES ('order_1'), ('order_2'))
CROSS JOIN UNNEST(ARRAY['shirt', 'pants', 'hat'])
WITH ORDINALITY AS t(product_name, index)
id product_name index
======= ============ =====
order_1 shirt 1
order_1 pants 2
order_1 hat 3
order_2 shirt 1
order_2 pants 2
order_2 hat 3
-- Returns a new row for each element and its position in the array
-- assuming a Orders table with an array column `product_names`
SELECT order_id, product_name, product_index
FROM Orders
CROSS JOIN UNNEST(product_names)
WITH ORDINALITY AS t(product_name, product_index)
ordinality가 있는 unnest는 각 요소와 데이터 구조에서 요소의 위치(1부터 시작)를 반환합니다. 배열의 요소 순서는 보장됩니다. map과 multiset은 순서가 없으므로 요소의 순서는 보장되지 않습니다.
-- Returns a new row each key/value pair in the map.
SELECT *
FROM
(VALUES('order_1'))
CROSS JOIN UNNEST(MAP['shirt', 2, 'pants', 1, 'hat', 1]) WITH ORDINALITY
id product_name amount index
======= ============ ===== =====
order_1 shirt 2 1
order_1 pants 1 2
order_1 hat 1 3
-- Returns a new row for each instance of a element in a multiset
-- If an element has been seen twice (multiplicity is 2), it will be returned twice
WITH ProductMultiset AS
(SELECT COLLECT(product_name) AS product_multiset
FROM (
VALUES ('shirt'), ('pants'), ('hat'), ('shirt'), ('hat')
) AS t(product_name)) -- produces { 'shirt': 2, 'pants': 1, 'hat': 2 }
SELECT id, product_name, ordinality
FROM
(VALUES ('order_1'), ('order_2')) AS t(id),
ProductMultiset
CROSS JOIN UNNEST(product_multiset) WITH
ORDINALITY AS u(product_name, ordinality);
id product_name index
======= ============ =====
order_1 shirt 1
order_1 shirt 2
order_1 pants 3
order_1 hat 4
order_1 hat 5
order_2 shirt 1
order_2 shirt 2
order_2 pants 3
order_2 hat 4
order_1 hat 5
Table Function
테이블 함수의 결과와 함께 테이블을 조인합니다. 왼쪽(외부) 테이블의 각 행은 테이블 함수의 해당 호출이 생성한 모든 행과 조인됩니다. 사용자 정의 테이블 함수는 사용 전에 등록되어야 합니다.
INNER JOIN
테이블 함수 호출이 빈 결과를 반환하면 왼쪽(외부) 테이블의 행은 버려집니다.
SELECT order_id, res
FROM Orders,
LATERAL TABLE(table_func(order_id)) t(res)
LEFT OUTER JOIN
테이블 함수 호출이 빈 결과를 반환하면 해당 외부 행은 보존되고 결과는 null 값으로 채워집니다. 현재 lateral table에 대한 left outer join은 ON 절에 TRUE 리터럴이 필요합니다.
SELECT order_id, res
FROM Orders
LEFT OUTER JOIN LATERAL TABLE(table_func(order_id)) t(res)
ON TRUE