윈도우 조인

윈도우 조인 (Window Join)

윈도우 조인은 조인 기준 자체에 시간 차원을 추가해요. 공통 키를 공유하고 같은 윈도우에 있는 두 스트림의 요소를 조인해요. 윈도우 조인의 의미론은 DataStream 윈도우 조인과 동일해요.

출처: 문서

본문

배치 | 스트리밍

스트리밍 쿼리의 경우, 연속 테이블의 다른 조인과 달리 윈도우 조인은 중간 결과를 방출하지 않고 윈도우 끝에서만 최종 결과를 방출해요. 또한 윈도우 조인은 더 이상 필요하지 않을 때 모든 중간 상태를 제거해요.

보통 Window Join은 Windowing TVF와 함께 사용돼요. 게다가 Window Join은 Window Aggregation, Window TopN, Window Join 같은 Windowing TVF 기반의 다른 연산 뒤에 올 수 있어요. 현재 Window Join은 조인 조건에 입력 테이블들의 윈도우 시작 동등성(window starts equality)과 입력 테이블들의 윈도우 끝 동등성(window ends equality)이 포함되어야 해요.

Window Join은 INNER/LEFT/RIGHT/FULL OUTER/ANTI/SEMI JOIN을 지원해요.

INNER/LEFT/RIGHT/FULL OUTER

다음은 INNER/LEFT/RIGHT/FULL OUTER Window Join 문의 구문이에요.

SELECT ...
FROM L [LEFT|RIGHT|FULL OUTER] JOIN R -- L and R are relations applied windowing TVF
ON L.window_start = R.window_start AND L.window_end = R.window_end AND ...

INNER/LEFT/RIGHT/FULL OUTER WINDOW JOIN의 구문은 서로 매우 유사해요. 여기서는 FULL OUTER JOIN의 예시만 제공해요.

윈도우 조인을 수행할 때 공통 키와 공통 tumbling 윈도우를 가진 모든 요소가 함께 조인돼요. Tumble Window TVF에서 동작하는 Window Join의 예시만 제공해요.

조인의 시간 영역을 고정된 5분 간격으로 범위를 한정함으로써, 데이터셋을 [12:00, 12:05)[12:05, 12:10)의 두 개의 별개 시간 윈도우로 잘랐어요. L2와 R2 행은 별개의 윈도우에 속하므로 서로 조인될 수 없어요.

Flink SQL> desc LeftTable;
+----------+------------------------+------+-----+--------+----------------------------------+
|     name |                   type | null | key | extras |                        watermark |
+----------+------------------------+------+-----+--------+----------------------------------+
| row_time | TIMESTAMP(3) *ROWTIME* | true |     |        | `row_time` - INTERVAL '1' SECOND |
|      num |                    INT | true |     |        |                                  |
|       id |                 STRING | true |     |        |                                  |
+----------+------------------------+------+-----+--------+----------------------------------+

Flink SQL> SELECT * FROM LeftTable;
+------------------+-----+----+
|         row_time | num | id |
+------------------+-----+----+
| 2020-04-15 12:02 |   1 | L1 |
| 2020-04-15 12:06 |   2 | L2 |
| 2020-04-15 12:03 |   3 | L3 |
+------------------+-----+----+

Flink SQL> desc RightTable;
+----------+------------------------+------+-----+--------+----------------------------------+
|     name |                   type | null | key | extras |                        watermark |
+----------+------------------------+------+-----+--------+----------------------------------+
| row_time | TIMESTAMP(3) *ROWTIME* | true |     |        | `row_time` - INTERVAL '1' SECOND |
|      num |                    INT | true |     |        |                                  |
|       id |                 STRING | true |     |        |                                  |
+----------+------------------------+------+-----+--------+----------------------------------+

Flink SQL> SELECT * FROM RightTable;
+------------------+-----+----+
|         row_time | num | id |
+------------------+-----+----+
| 2020-04-15 12:01 |   2 | R2 |
| 2020-04-15 12:04 |   3 | R3 |
| 2020-04-15 12:05 |   4 | R4 |
+------------------+-----+----+

Flink SQL> SELECT L.num as L_Num, L.id as L_Id, R.num as R_Num, R.id as R_Id,
           COALESCE(L.window_start, R.window_start) as window_start,
           COALESCE(L.window_end, R.window_end) as window_end
           FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
           ) L
           FULL JOIN (
               SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
           ) R
           ON L.num = R.num AND L.window_start = R.window_start AND L.window_end = R.window_end;
+-------+------+-------+------+------------------+------------------+
| L_Num | L_Id | R_Num | R_Id |     window_start |       window_end |
+-------+------+-------+------+------------------+------------------+
|     1 |   L1 |  null | null | 2020-04-15 12:00 | 2020-04-15 12:05 |
|  null | null |     2 |   R2 | 2020-04-15 12:00 | 2020-04-15 12:05 |
|     3 |   L3 |     3 |   R3 | 2020-04-15 12:00 | 2020-04-15 12:05 |
|     2 |   L2 |  null | null | 2020-04-15 12:05 | 2020-04-15 12:10 |
|  null | null |     4 |   R4 | 2020-04-15 12:05 | 2020-04-15 12:10 |
+-------+------+-------+------+------------------+------------------+

참고: 윈도우잉의 동작을 더 잘 이해하기 위해, 타임스탬프 값을 표시할 때 뒤따르는 0을 표시하지 않도록 단순화했어요. 예: 타입이 TIMESTAMP(3)이면 2020-04-15 08:05는 Flink SQL Client에서 2020-04-15 08:05:00.000으로 표시되어야 해요.

SEMI

Semi Window Join은 공통 윈도우 내에서 오른쪽에 일치하는 행이 적어도 하나 있으면 왼쪽 레코드에서 행 하나를 반환해요.

Flink SQL> SELECT *
           FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
           ) L WHERE L.num IN (
             SELECT num FROM (   
               SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
             ) R WHERE L.window_start = R.window_start AND L.window_end = R.window_end);
+------------------+-----+----+------------------+------------------+-------------------------+
|         row_time | num | id |     window_start |       window_end |            window_time  |
+------------------+-----+----+------------------+------------------+-------------------------+
| 2020-04-15 12:03 |   3 | L3 | 2020-04-15 12:00 | 2020-04-15 12:05 | 2020-04-15 12:04:59.999 |
+------------------+-----+----+------------------+------------------+-------------------------+

Flink SQL> SELECT *
           FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
           ) L WHERE EXISTS (
             SELECT * FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
             ) R WHERE L.num = R.num AND L.window_start = R.window_start AND L.window_end = R.window_end);
+------------------+-----+----+------------------+------------------+-------------------------+
|         row_time | num | id |     window_start |       window_end |            window_time  |
+------------------+-----+----+------------------+------------------+-------------------------+
| 2020-04-15 12:03 |   3 | L3 | 2020-04-15 12:00 | 2020-04-15 12:05 | 2020-04-15 12:04:59.999 |
+------------------+-----+----+------------------+------------------+-------------------------+

참고: 윈도우잉의 동작을 더 잘 이해하기 위해, 타임스탬프 값을 표시할 때 뒤따르는 0을 표시하지 않도록 단순화했어요. 예: 타입이 TIMESTAMP(3)이면 2020-04-15 08:05는 Flink SQL Client에서 2020-04-15 08:05:00.000으로 표시되어야 해요.

ANTI

Anti Window Join은 Inner Window Join의 반대예요. 각 공통 윈도우 내에서 조인되지 않은 모든 행을 포함해요.

Flink SQL> SELECT *
           FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
           ) L WHERE L.num NOT IN (
             SELECT num FROM (   
               SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
             ) R WHERE L.window_start = R.window_start AND L.window_end = R.window_end);
+------------------+-----+----+------------------+------------------+-------------------------+
|         row_time | num | id |     window_start |       window_end |            window_time  |
+------------------+-----+----+------------------+------------------+-------------------------+
| 2020-04-15 12:02 |   1 | L1 | 2020-04-15 12:00 | 2020-04-15 12:05 | 2020-04-15 12:04:59.999 |
| 2020-04-15 12:06 |   2 | L2 | 2020-04-15 12:05 | 2020-04-15 12:10 | 2020-04-15 12:09:59.999 |
+------------------+-----+----+------------------+------------------+-------------------------+

Flink SQL> SELECT *
           FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
           ) L WHERE NOT EXISTS (
             SELECT * FROM (
               SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES))
             ) R WHERE L.num = R.num AND L.window_start = R.window_start AND L.window_end = R.window_end);
+------------------+-----+----+------------------+------------------+-------------------------+
|         row_time | num | id |     window_start |       window_end |            window_time  |
+------------------+-----+----+------------------+------------------+-------------------------+
| 2020-04-15 12:02 |   1 | L1 | 2020-04-15 12:00 | 2020-04-15 12:05 | 2020-04-15 12:04:59.999 |
| 2020-04-15 12:06 |   2 | L2 | 2020-04-15 12:05 | 2020-04-15 12:10 | 2020-04-15 12:09:59.999 |
+------------------+-----+----+------------------+------------------+-------------------------+

참고: 윈도우잉의 동작을 더 잘 이해하기 위해, 타임스탬프 값을 표시할 때 뒤따르는 0을 표시하지 않도록 단순화했어요. 예: 타입이 TIMESTAMP(3)이면 2020-04-15 08:05는 Flink SQL Client에서 2020-04-15 08:05:00.000으로 표시되어야 해요.

제한 (Limitation)

Join 절에 대한 제한

현재 윈도우 조인은 조인 조건에 입력 테이블들의 윈도우 시작 동등성과 윈도우 끝 동등성이 포함되어야 해요. 나중에는 windwoing TVF가 TUMBLE 또는 HOP인 경우 조인 조건을 윈도우 시작 동등성만 포함하도록 단순화할 수도 있어요.

입력의 Windowing TVF에 대한 제한

현재 windwoing TVF는 좌우 입력이 동일해야 해요. 이는 나중에 확장될 수 있어요. 예를 들어 같은 윈도우 크기의 tumbling 윈도우가 sliding 윈도우와 조인하는 경우가 있어요.

Windowing TVF 뒤에 직접 오는 Window Join에 대한 제한

현재 Window Join이 Windowing TVF 뒤에 오면, Windowing TVF는 Session 윈도우가 아닌 Tumble Window, Hop Window 또는 Cumulate Window여야 해요.

더 알아보기 (Learn more)