조인

Flink SQL은 동적 테이블에 대한 복잡하고 유연한 조인 연산을 지원합니다. 다양한 의미의 질의를 처리하기 위해 여러 유형의 조인이 있습니다: Regular Join, Interval Join, Temporal Join, Lookup Join, 배열/멀티셋/맵 확장, Table Function 조인 등.

출처: 문서

본문

배치/스트리밍

Flink SQL은 동적 테이블에 대한 복잡하고 유연한 조인 연산을 지원합니다. 질의가 요구할 수 있는 다양한 의미를 설명하기 위해 여러 가지 유형의 조인이 있습니다.

기본적으로 조인의 순서는 최적화되지 않습니다. 테이블은 FROM 절에 지정된 순서대로 조인됩니다. 갱신 빈도가 가장 낮은 테이블을 먼저, 갱신 빈도가 가장 높은 테이블을 마지막에 배치하여 조인 질의 성능을 조정할 수 있습니다. 지원되지 않고 질의 실패를 일으키는 크로스 조인(카테시안 곱)이 되지 않는 순서로 테이블을 지정해야 합니다.

일반 조인 (Regular Joins)

일반 조인은 가장 일반적인 유형의 조인으로, 새 레코드 또는 조인 양쪽의 변경이 모두 보이고 전체 조인 결과에 영향을 미칩니다. 예를 들어 왼쪽에 새 레코드가 있으면 product id가 같을 때 오른쪽의 이전 및 이후 레코드 모두와 조인됩니다.

SELECT * FROM Orders
INNER JOIN Product
ON Orders.productId = Product.id

스트리밍 질의의 경우 일반 조인의 문법이 가장 유연하며 모든 종류의 갱신(insert, update, delete) 입력 테이블을 허용합니다. 그러나 이 연산은 중요한 운영상 의미를 가집니다: 조인 입력의 양쪽을 Flink 상태에 영원히 유지해야 합니다. 따라서 질의 결과를 계산하는 데 필요한 상태가 모든 입력 테이블의 고유 입력 행 수와 중간 조인 결과에 따라 무한히 커질 수 있습니다. 과도한 상태 크기를 방지하기 위해 적절한 state time-to-live(TTL)로 질의 구성을 제공할 수 있습니다. 이는 질의 결과의 정확성에 영향을 줄 수 있음에 유의하세요. 자세한 내용은 query configuration 참고.

{{< query_state_warning >}}

INNER Equi-JOIN

조인 조건에 의해 제한된 단순 카테시안 곱을 반환합니다. 현재 equi-join만 지원됩니다. 즉, 같음 술어(equality predicate)가 있는 적어도 하나의 접속(conjunctive) 조건이 있는 조인입니다. 임의의 크로스 또는 theta 조인은 지원되지 않습니다.

SELECT *
FROM Orders
INNER JOIN Product
ON Orders.product_id = Product.id

OUTER Equi-JOIN

적격 카테시안 곱에서 모든 행(즉, 조인 조건을 통과하는 모든 결합 행)을 반환하며, 조인 조건이 다른 테이블의 어떤 행과도 일치하지 않은 외부 테이블의 각 행에 대해 한 복사본을 추가로 반환합니다. Flink는 LEFT, RIGHT, FULL 외부 조인을 지원합니다. 현재 equi-join만 지원됩니다. 즉, 같음 술어가 있는 적어도 하나의 접속 조건이 있는 조인입니다. 임의의 크로스 또는 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

다중 일반 조인 (Multiple Regular Joins)

여러 체인 조인의 성능 문제가 있고 이들이 많은 상태를 생성한다면 MultiJoin 연산자 사용을 고려하세요. 자세한 내용은 tuning multiple regular joins 참고.

간격 조인 (Interval Joins)

조인 조건과 시간 제약에 의해 제한된 단순 카테시안 곱을 반환합니다. 간격 조인은 적어도 하나의 equi-join 술어와 양쪽의 시간을 제한하는 조인 조건을 요구합니다. 두 개의 적절한 범위 술어(<, <=, >=, >), BETWEEN 술어, 또는 두 입력 테이블의 동일한 유형(프로세싱 타임 또는 이벤트 타임)의 시간 속성을 비교하는 단일 같음 술어가 그러한 조건을 정의할 수 있습니다.

예를 들어 이 질의는 주문이 수신된 지 4시간 후에 배송된 경우 모든 주문을 해당 배송과 조인합니다.

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

다음 술어는 유효한 간격 조인 조건의 예입니다:

  • ltime = rtime
  • ltime >= rtime AND ltime < rtime + INTERVAL '10' MINUTE
  • ltime BETWEEN rtime - INTERVAL '10' SECOND AND rtime + INTERVAL '5' SECOND

스트리밍 질의의 경우 일반 조인과 비교해 간격 조인은 시간 속성을 가진 append-only 테이블만 지원합니다. 시간 속성은 준단조 증가(quasi-monotonic increasing)하므로 Flink는 결과의 정확성에 영향을 주지 않고 상태에서 오래된 값을 제거할 수 있습니다.

시간 조인 (Temporal Joins)

시간 테이블(Temporal table)은 시간에 따라 진화하는 테이블입니다 - Flink에서는 동적 테이블(dynamic table)로 알려져 있습니다. 시간 테이블의 행은 하나 이상의 시간 기간과 연관되며 모든 Flink 테이블은 시간적(동적)입니다. 시간 테이블은 하나 이상의 버전 테이블 스냅샷을 포함하며, 변경 사항을 추적하는 변경 히스토리 테이블(예: 데이터베이스 changelog, 모든 스냅샷 포함) 또는 변경을 머티리얼라이즈하는 변경 디멘션 테이블(예: 최신 스냅샷을 포함하는 데이터베이스 테이블)일 수 있습니다.

이벤트 타임 시간 조인 (Event Time Temporal Join)

이벤트 타임 시간 조인은 versioned table에 대하여 조인할 수 있게 합니다. 즉, 테이블을 변경하는 메타데이터로 보강하고 특정 시점의 값을 검색할 수 있습니다.

시간 조인은 임의 테이블(왼쪽 입력/프로브 측)을 가져와 각 행을 버전 테이블(오른쪽 입력/빌드 측)의 해당 행 관련 버전과 연관시킵니다. Flink는 SQL:2011 표준의 FOR SYSTEM_TIME AS OF 문법을 사용하여 이 연산을 수행합니다. 시간 조인의 문법은 다음과 같습니다:

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

참고: 이벤트 타임 시간 조인은 왼쪽과 오른쪽의 워터마크에 의해 트리거됩니다. INTERVAL 시간 빼기는 늦은 이벤트를 기다려 조인이 기대치를 충족하도록 하는 데 사용됩니다. 조인의 양쪽이 워터마크를 올바르게 설정했는지 확인하세요.

참고: 이벤트 타임 시간 조인은 시간 조인 조건의 등가 조건에 기본 키가 포함되어 있어야 합니다. 예: 테이블 currency_rates의 기본 키 currency_rates.currency가 조건 orders.currency = currency_rates.currency에 제한되어야 합니다.

일반 조인과 대조적으로 빌드 측의 변경에도 이전 시간 테이블 결과는 영향을 받지 않습니다. 간격 조인과 비교하면 시간 테이블 조인은 레코드가 조인될 시간 윈도우를 정의하지 않습니다. 프로브 측의 레코드는 항상 시간 속성이 지정한 시점의 빌드 측 버전과 조인됩니다. 따라서 빌드 측의 행은 임의로 오래될 수 있습니다. 시간이 지남에 따라 (주어진 기본 키에 대한) 더 이상 필요하지 않은 레코드 버전은 상태에서 제거됩니다.

프로세싱 타임 시간 조인 (Processing Time Temporal Join)

프로세싱 타임 시간 테이블 조인은 프로세싱 타임 속성을 사용하여 외부 버전 테이블의 키 최신 버전에 행을 연관시킵니다.

정의상 프로세싱 타임 속성을 사용하면 조인이 주어진 키에 대해 항상 가장 최신의 값을 반환합니다. 룩업 테이블을 빌드 측의 모든 레코드를 저장하는 단순한 HashMap<K, V>로 생각할 수 있습니다. 이 조인의 강점은 Flink 내에서 테이블을 동적 테이블로 머티리얼라이즈하는 것이 불가능할 때 외부 시스템에 대해 직접 작업할 수 있게 한다는 것입니다.

다음 프로세싱 타임 시간 테이블 조인 예시는 테이블 LatestRates와 조인되어야 하는 append-only 테이블 orders를 보여줍니다. 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:1510:30LastestRates 내용은 같습니다. 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

현재 임의 뷰/테이블의 최신 버전과의 시간 조인에 사용되는 FOR SYSTEM_TIME AS OF 문법은 아직 지원되지 않습니다. 대신 다음과 같이 시간 테이블 함수 문법을 사용할 수 있습니다:

SELECT
  o_amount, r_rate
FROM
  Orders,
  LATERAL TABLE (Rates(o_proctime))
WHERE
  r_currency = o_currency

참고 임의 테이블/뷰의 최신 버전과의 시간 조인에 FOR SYSTEM_TIME AS OF 문법이 지원되지 않는 이유는 순전히 의미론적 고려 때문입니다. 왼쪽 스트림에 대한 조인 처리가 시간 테이블의 완전한 스냅샷을 기다리지 않기 때문이며, 이는 프로덕션 환경에서 사용자를 오도할 수 있습니다. 시간 테이블 함수에 의한 프로세싱 타임 시간 조인도 같은 의미론적 문제가 있지만 오랫동안 살아 있어 호환성 관점에서 지원합니다.

프로세싱 타임의 결과는 결정적이지 않습니다. 프로세싱 타임 시간 조인은 외부 테이블(즉, 디멘션 테이블)로 스트림을 보강하는 데 가장 자주 사용됩니다.

일반 조인과 대조적으로 빌드 측의 변경에도 이전 시간 테이블 결과는 영향을 받지 않습니다. 간격 조인과 비교하면 시간 테이블 조인은 레코드가 조인되는 시간 윈도우를 정의하지 않습니다. 즉, 오래된 행은 상태에 저장되지 않습니다.

시간 테이블 함수 조인 (Temporal Table Function Join)

테이블을 시간 테이블 함수(temporal table function)와 조인하는 문법은 Table Function과의 조인과 같습니다.

참고: 현재 시간 테이블과의 inner join과 left outer join만 지원됩니다.

Rates가 시간 테이블 함수라고 가정하면 조인은 SQL에서 다음과 같이 표현할 수 있습니다:

SELECT
  o_amount, r_rate
FROM
  Orders,
  LATERAL TABLE (Rates(o_proctime))
WHERE
  r_currency = o_currency

위의 Temporal Table DDL과 Temporal Table Function의 주요 차이는 다음과 같습니다:

  • 시간 테이블 DDL은 SQL로 정의할 수 있지만 시간 테이블 함수는 정의할 수 없습니다;
  • 시간 테이블 DDL과 시간 테이블 함수 모두 버전 테이블의 시간 조인을 지원하지만, 시간 테이블 함수만 임의 테이블/뷰의 최신 버전을 시간 조인할 수 있습니다.

룩업 조인 (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 테이블의 각 행이 조인 연산자가 Orders 행을 처리하는 시점에 조인 술어와 일치하는 Customers 행과 조인되도록 보장합니다. 또한 조인된 Customer 행이 나중에 갱신되어도 조인 결과가 갱신되지 않게 합니다. 룩업 조인은 위 예시의 o.customer_id = c.id 같은 필수 같음 조인 술어도 요구합니다.

배열, 멀티셋, 맵 확장 (Array, Multiset and Map Expansion)

Unnest는 주어진 배열, 멀티셋 또는 맵의 각 요소에 대해 새 행을 반환합니다. CROSS JOINLEFT 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로의 unnesting도 지원됩니다. 현재 WITH ORDINALITYCROSS 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부터 시작)를 반환합니다. 배열의 요소 순서는 보장됩니다. 맵과 멀티셋은 순서가 없으므로 요소의 순서는 보장되지 않습니다.

-- 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

더 알아보기 (Learn more)