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_start와 window_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은 가까운 시일 내에 지원될 예정이에요.