Window Top-N

Window Top-N (윈도우 탑-N)

Window Top-N은 각 윈도우와 기타 파티셔닝 키에 대해 N개의 가장 작거나 큰 값을 반환하는 특수한 [Top-N]({{< ref "docs/sql/reference/queries/topn" >}})이에요. 배치와 스트리밍 모두에서 지원돼요.

출처: Window Top-N

본문

{{< label Batch >}} {{< label Streaming >}}

Window Top-N은 각 윈도우와 기타 파티셔닝 키에 대해 N개의 가장 작거나 큰 값을 반환하는 특수한 [Top-N]({{< ref "docs/sql/reference/queries/topn" >}})이에요.

스트리밍 쿼리의 경우, 연속 테이블에 대한 일반 Top-N과 달리 Window Top-N은 중간 결과를 방출하지 않고 윈도우 끝에서 총 상위 N개 레코드라는 최종 결과만 방출해요. 또한 Window Top-N은 더 이상 필요하지 않을 때 모든 중간 상태를 제거해요. 따라서 사용자가 레코드별 업데이트 결과를 필요로 하지 않는다면 Window Top-N 쿼리가 더 나은 성능을 보여요. 보통 Window Top-N은 [Windowing TVF]({{< ref "docs/sql/reference/queries/window-tvf" >}})와 함께 직접 사용돼요. 게다가 Window Top-N은 [Windowing TVF]({{< ref "docs/sql/reference/queries/window-tvf" >}})에 기반한 다른 연산들과도 함께 사용될 수 있어요. 예를 들어 [Window Aggregation]({{< ref "docs/sql/reference/queries/window-agg" >}}), [Window TopN]({{< ref "docs/sql/reference/queries/window-topn">}}), 그리고 [Window Join]({{< ref "docs/sql/reference/queries/window-join">}})과 함께요.

{{< hint info >}} Note: SESSION Window Top-N은 현재 배치 모드에서 지원되지 않아요. {{< /hint >}}

Window Top-N은 일반 Top-N과 같은 구문으로 정의할 수 있어요. 자세한 내용은 [Top-N 문서]({{< ref "docs/sql/reference/queries/topn" >}})를 참조하세요. 그 외에 Window Top-N은 PARTITION BY 절이 [Windowing TVF]({{< ref "docs/sql/reference/queries/window-tvf" >}}) 또는 [Window Aggregation]({{< ref "docs/sql/reference/queries/window-agg" >}})을 적용한 관계의 window_startwindow_end 열을 포함해야 해요. 그렇지 않으면 옵티마이저가 쿼리를 변환할 수 없어요.

다음은 Window Top-N 문의 구문이에요:

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

Example (예제)

Window Aggregation 뒤의 Window Top-N

다음 예제는 10분마다의 텀블링 윈도우에서 판매액이 가장 높은 Top 3 공급업체를 계산하는 방법을 보여줘요.

-- 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 |     |        |                                 |
| supplier_id |                 STRING | true |     |        |                                 |
+-------------+------------------------+------+-----+--------+---------------------------------+

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

Flink SQL> SELECT *
  FROM (
    SELECT *, ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY price DESC) as rownum
    FROM (
      SELECT window_start, window_end, supplier_id, SUM(price) as price, COUNT(*) as cnt
      FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
      GROUP BY window_start, window_end, supplier_id
    )
  ) WHERE rownum <= 3;
+------------------+------------------+-------------+-------+-----+--------+
|     window_start |       window_end | supplier_id | price | cnt | rownum |
+------------------+------------------+-------------+-------+-----+--------+
| 2020-04-15 08:00 | 2020-04-15 08:10 |   supplier1 |  6.00 |   2 |      1 |
| 2020-04-15 08:00 | 2020-04-15 08:10 |   supplier4 |  5.00 |   1 |      2 |
| 2020-04-15 08:00 | 2020-04-15 08:10 |   supplier2 |  4.00 |   1 |      3 |
| 2020-04-15 08:10 | 2020-04-15 08:20 |   supplier5 |  6.00 |   1 |      1 |
| 2020-04-15 08:10 | 2020-04-15 08:20 |   supplier2 |  3.00 |   1 |      2 |
| 2020-04-15 08:10 | 2020-04-15 08:20 |   supplier3 |  2.00 |   1 |      3 |
+------------------+------------------+-------------+-------+-----+--------+

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

Windowing TVF 뒤의 Window Top-N

다음 예제는 10분마다의 텀블링 윈도우에서 가격이 가장 높은 Top 3 항목을 계산하는 방법을 보여줘요.

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 price DESC) as rownum
    FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
  ) WHERE rownum <= 3;
+------------------+-------+------+-------------+------------------+------------------+--------+
|          bidtime | price | item | supplier_id |     window_start |       window_end | rownum |
+------------------+-------+------+-------------+------------------+------------------+--------+
| 2020-04-15 08:05 |  4.00 |    A |   supplier1 | 2020-04-15 08:00 | 2020-04-15 08:10 |      2 |
| 2020-04-15 08:06 |  4.00 |    C |   supplier2 | 2020-04-15 08:00 | 2020-04-15 08:10 |      3 |
| 2020-04-15 08:09 |  5.00 |    D |   supplier4 | 2020-04-15 08:00 | 2020-04-15 08:10 |      1 |
| 2020-04-15 08:11 |  2.00 |    B |   supplier3 | 2020-04-15 08:10 | 2020-04-15 08:20 |      3 |
| 2020-04-15 08:15 |  3.00 |    H |   supplier2 | 2020-04-15 08:10 | 2020-04-15 08:20 |      2 |
| 2020-04-15 08:17 |  6.00 |    F |   supplier5 | 2020-04-15 08:10 | 2020-04-15 08:20 |      1 |
+------------------+-------+------+-------------+------------------+------------------+--------+

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

Limitation (제한 사항)

현재 Flink는 Tumble Windows, Hop Windows, Cumulate Windows가 있는 [Windowing TVF]({{< ref "docs/sql/reference/queries/window-tvf" >}}) 뒤의 Window Top-N만 지원해요. Session windows가 있는 [Windowing TVF]({{< ref "docs/sql/reference/queries/window-tvf" >}}) 뒤의 Window Top-N은 가까운 시일 내에 지원될 예정이에요.

더 알아보기 (Learn more)