윈도우 중복 제거

윈도우 중복 제거 (Window Deduplication)

Window Deduplication은 컬럼 집합에 대해 중복되는 행을 제거하고 각 창(window)과 파티션 키에 대해 첫 번째 또는 마지막 행을 유지하는 특수한 Deduplication입니다. 스트리밍 쿼리에서 일반 Deduplicate와 달리 Window Deduplication은 중간 결과를 내보내지 않고 창이 끝날 때 최종 결과만 내보냅니다.

출처: 문서

본문

Window Deduplication은 컬럼 집합에 대해 중복되는 행을 제거하고 각 창과 파티션 키에 대해 첫 번째 또는 마지막 행을 유지하는 특수한 Deduplication입니다.

스트리밍 쿼리의 경우, 연속 테이블에 대한 일반 Deduplicate와 달리 Window Deduplication은 중간 결과를 내보내지 않고 창이 끝날 때 최종 결과만 내보냅니다. 또한 Window Deduplication은 더 이상 필요하지 않을 때 모든 중간 상태를 정리합니다. 따라서 사용자가 레코드별로 갱신된 결과를 필요로 하지 않는다면 Window Deduplication 쿼리가 더 나은 성능을 제공합니다.

보통 Window Deduplication은 Windowing TVF와 직접 함께 사용됩니다. 그 외에도 Window Deduplication은 Window Aggregation, Window TopN, Window Join처럼 Windowing TVF를 기반으로 한 다른 연산과 함께 사용할 수 있습니다.

Window Deduplication은 일반 Deduplication과 같은 문법으로 정의할 수 있습니다. 자세한 내용은 Deduplication 문서를 참조하세요. 그 외에도 Window Deduplication은 PARTITION BY 절이 관계의 window_startwindow_end 컬럼을 포함해야 합니다. 그렇지 않으면 옵티마이저가 쿼리를 변환할 수 없습니다.

Flink는 Window Top-N 쿼리와 같은 방식으로 ROW_NUMBER()를 사용해 중복을 제거합니다. 이론적으로 Window Deduplication은 N이 1이고 처리 시간 또는 이벤트 시간으로 정렬하는 Window Top-N의 특수한 경우입니다.

다음은 Window Deduplication 문의 문법입니다.

SELECT [column_list]
FROM (
   SELECT [column_list],
     ROW_NUMBER() OVER (PARTITION BY window_start, window_end [, col_key1...]
       ORDER BY time_attr [asc|desc]) AS rownum
   FROM table_name) -- relation applied windowing TVF
WHERE (rownum = 1 | rownum <=1 | rownum < 2) [AND conditions]

파라미터 사양:

  • ROW_NUMBER(): 각 행에 1부터 시작하는 고유하고 순차적인 번호를 할당합니다.
  • PARTITION BY window_start, window_end [, col_key1...]: window_start, window_end와 다른 파티션 키를 포함하는 파티션 컬럼을 지정합니다.
  • ORDER BY time_attr [asc|desc]: 정렬 컬럼을 지정합니다. 이는 시간 속성(time attribute)이어야 합니다. 현재 Flink는 처리 시간 속성이벤트 시간 속성을 지원합니다. ASC로 정렬하면 첫 번째 행을 유지하고, DESC로 정렬하면 마지막 행을 유지합니다.
  • WHERE (rownum = 1 | rownum <=1 | rownum < 2): rownum = 1 | rownum <=1 | rownum < 2는 옵티마이저가 쿼리를 Window Deduplication으로 변환될 수 있다고 인식하도록 하기 위해 필요합니다.

참고: 위 패턴을 정확히 따라야 합니다. 그렇지 않으면 옵티마이저가 쿼리를 Window Deduplication으로 변환하지 않습니다.

예시 (Example)

다음 예시는 10분마다 텀블링 창에서 마지막 레코드를 유지하는 방법을 보여줍니다.

-- tables must have time attribute, e.g. `bidtime` in this table
Flink SQL> DESC Bid;
+-------------+------------------------+------+-----+--------+---------------------------------+
|        name |                   type | null | key | extras |                       watermark |
+-------------+------------------------+------+-----+--------+---------------------------------+
|     bidtime | TIMESTAMP(3) *ROWTIME* | true |     |        | `bidtime` - INTERVAL '1' SECOND |
|       price |         DECIMAL(10, 2) | true |     |        |                                 |
|        item |                 STRING | true |     |        |                                 |
+-------------+------------------------+------+-----+--------+---------------------------------+

Flink SQL> SELECT * FROM Bid;
+------------------+-------+------+
|          bidtime | price | item |
+------------------+-------+------+
| 2020-04-15 08:05 |  4.00 | C    |
| 2020-04-15 08:07 |  2.00 | A    |
| 2020-04-15 08:09 |  5.00 | D    |
| 2020-04-15 08:11 |  3.00 | B    |
| 2020-04-15 08:13 |  1.00 | E    |
| 2020-04-15 08:17 |  6.00 | F    |
+------------------+-------+------+

Flink SQL> SELECT *
  FROM (
    SELECT bidtime, price, item, supplier_id, window_start, window_end,
      ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY bidtime DESC) AS rownum
    FROM TABLE(
               TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES))
  ) WHERE rownum <= 1;
+------------------+-------+------+-------------+------------------+------------------+--------+
|          bidtime | price | item | supplier_id |     window_start |       window_end | rownum |
+------------------+-------+------+-------------+------------------+------------------+--------+
| 2020-04-15 08:09 |  5.00 |    D |   supplier4 | 2020-04-15 08:00 | 2020-04-15 08:10 |      1 |
| 2020-04-15 08:17 |  6.00 |    F |   supplier5 | 2020-04-15 08:10 | 2020-04-15 08:20 |      1 |
+------------------+-------+------+-------------+------------------+------------------+--------+

참고: 윈도우 동작을 더 잘 이해하기 위해 타임스탬프 표시를 단순화해 끝의 0을 표시하지 않았습니다. 예를 들어 타입이 TIMESTAMP(3)이라면 2020-04-15 08:05는 Flink SQL Client에서 2020-04-15 08:05:00.000으로 표시되어야 합니다.

제한 사항 (Limitation)

Windowing TVF 바로 뒤에 오는 Window Deduplication의 제한

현재 Windowing TVF 뒤에 Window Deduplication이 오는 경우, Windowing TVF는 Session 창이 아닌 Tumble 창, Hop 창, Cumulate 창이어야 합니다. Session 창은 가까운 미래에 지원될 예정입니다.

정렬 키 시간 속성의 제한

현재 Window Deduplication은 정렬 키가 처리 시간 속성이 아닌 이벤트 시간 속성이어야 합니다. 처리 시간으로 정렬하는 것은 가까운 미래에 지원될 예정입니다.

더 알아보기 (Learn more)