윈도우 집계
윈도우 집계 (Window Aggregation)
본문
Window TVF 집계 (Window TVF Aggregation)
배치 / 스트리밍
윈도우 집계는 Windowing TVF가 적용된 릴레이션의 "window_start"와 "window_end" 컬럼을 포함하는 GROUP BY 절에 정의됩니다. 일반 GROUP BY 절이 있는 쿼리처럼, 윈도우 집계로 그룹화하는 쿼리는 그룹당 하나의 결과 행을 계산합니다.
SELECT ...
FROM <windowed_table> -- relation applied windowing TVF
GROUP BY window_start, window_end, ...
연속 테이블의 다른 집계와 달리 윈도우 집계는 중간 결과를 내보내지 않고 최종 결과, 즉 윈도우가 끝날 때의 전체 집계만 내보냅니다. 또한 윈도우 집계는 더 이상 필요하지 않을 때 모든 중간 상태를 제거합니다.
Windowing TVF
Flink는 TUMBLE, HOP, CUMULATE, SESSION 유형의 윈도우 집계를 지원합니다. 스트리밍 모드에서 윈도우 테이블 값 함수의 시간 속성 필드는 이벤트 시간 또는 처리 시간 속성 중 하나여야 합니다. 더 많은 윈도잉 함수 정보는 Windowing TVF를 참고하세요. 배치 모드에서 윈도우 테이블 값 함수의 시간 속성 필드는 TIMESTAMP 또는 TIMESTAMP_LTZ 타입의 속성이어야 합니다.
참고:
SESSION윈도우 집계는 현재 배치 모드에서 지원되지 않습니다.
다음은 TUMBLE, HOP, CUMULATE, SESSION 윈도우 집계의 몇 가지 예시입니다.
-- 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 | C | supplier1 |
| 2020-04-15 08:07 | 2.00 | A | supplier1 |
| 2020-04-15 08:09 | 5.00 | D | supplier2 |
| 2020-04-15 08:11 | 3.00 | B | supplier2 |
| 2020-04-15 08:13 | 1.00 | E | supplier1 |
| 2020-04-15 08:17 | 6.00 | F | supplier2 |
+------------------+-------+------+-------------+
-- tumbling window aggregation
Flink SQL> SELECT window_start, window_end, SUM(price) AS total_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
GROUP BY window_start, window_end;
+------------------+------------------+-------------+
| window_start | window_end | total_price |
+------------------+------------------+-------------+
| 2020-04-15 08:00 | 2020-04-15 08:10 | 11.00 |
| 2020-04-15 08:10 | 2020-04-15 08:20 | 10.00 |
+------------------+------------------+-------------+
-- hopping window aggregation
Flink SQL> SELECT window_start, window_end, SUM(price) AS total_price
FROM HOP(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES, INTERVAL '10' MINUTES)
GROUP BY window_start, window_end;
+------------------+------------------+-------------+
| window_start | window_end | total_price |
+------------------+------------------+-------------+
| 2020-04-15 08:00 | 2020-04-15 08:10 | 11.00 |
| 2020-04-15 08:05 | 2020-04-15 08:15 | 15.00 |
| 2020-04-15 08:10 | 2020-04-15 08:20 | 10.00 |
| 2020-04-15 08:15 | 2020-04-15 08:25 | 6.00 |
+------------------+------------------+-------------+
-- cumulative window aggregation
Flink SQL> SELECT window_start, window_end, SUM(price) AS total_price
FROM CUMULATE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '2' MINUTES, INTERVAL '10' MINUTES)
GROUP BY window_start, window_end;
+------------------+------------------+-------------+
| window_start | window_end | total_price |
+------------------+------------------+-------------+
| 2020-04-15 08:00 | 2020-04-15 08:06 | 4.00 |
| 2020-04-15 08:00 | 2020-04-15 08:08 | 6.00 |
| 2020-04-15 08:00 | 2020-04-15 08:10 | 11.00 |
| 2020-04-15 08:10 | 2020-04-15 08:12 | 3.00 |
| 2020-04-15 08:10 | 2020-04-15 08:14 | 4.00 |
| 2020-04-15 08:10 | 2020-04-15 08:16 | 4.00 |
| 2020-04-15 08:10 | 2020-04-15 08:18 | 10.00 |
| 2020-04-15 08:10 | 2020-04-15 08:20 | 10.00 |
+------------------+------------------+-------------+
-- session window aggregation with partition keys
Flink SQL> SELECT window_start, window_end, supplier_id, SUM(price) AS total_price
FROM SESSION(TABLE Bid PARTITION BY supplier_id, DESCRIPTOR(bidtime), INTERVAL '2' MINUTES)
GROUP BY window_start, window_end, supplier_id;
+------------------+------------------+-------------+-------------+
| window_start | window_end | supplier_id | total_price |
+------------------+------------------+-------------+-------------+
| 2020-04-15 08:05 | 2020-04-15 08:09 | supplier1 | 6.00 |
| 2020-04-15 08:09 | 2020-04-15 08:13 | supplier2 | 8.00 |
| 2020-04-15 08:13 | 2020-04-15 08:15 | supplier1 | 1.00 |
| 2020-04-15 08:17 | 2020-04-15 08:19 | supplier2 | 6.00 |
+------------------+------------------+-------------+-------------+
-- session window aggregation without partition keys
Flink SQL> SELECT window_start, window_end, SUM(price) AS total_price
FROM SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '2' MINUTES)
GROUP BY window_start, window_end;
+------------------+------------------+-------------+
| window_start | window_end | total_price |
+------------------+------------------+-------------+
| 2020-04-15 08:05 | 2020-04-15 08:15 | 15.00 |
| 2020-04-15 08:17 | 2020-04-15 08:19 | 6.00 |
+------------------+------------------+-------------+
참고: 윈도잉의 동작을 더 잘 이해하기 위해 타임스탬프 값 표시를 단순화해 뒤의 0을 보여주지 않습니다. 예를 들어
2020-04-15 08:05는 타입이TIMESTAMP(3)이면 Flink SQL Client에서2020-04-15 08:05:00.000으로 표시되어야 합니다.
GROUPING SETS
윈도우 집계는 GROUPING SETS 문법도 지원합니다. 그룹핑 세트는 표준 GROUP BY로 표현할 수 있는 것보다 더 복잡한 그룹핑 연산을 허용합니다. 행은 지정된 각 그룹핑 세트에 의해 개별적으로 그룹화되며, 단순 GROUP BY 절과 마찬가지로 각 그룹에 대해 집계가 계산됩니다.
GROUPING SETS가 있는 윈도우 집계는 window_start와 window_end 컬럼이 모두 GROUP BY 절에 있어야 하며, GROUPING SETS 절에는 있으면 안 됩니다.
Flink SQL> SELECT window_start, window_end, supplier_id, SUM(price) AS total_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
GROUP BY window_start, window_end, GROUPING SETS ((supplier_id), ());
+------------------+------------------+-------------+-------------+
| window_start | window_end | supplier_id | total_price |
+------------------+------------------+-------------+-------------+
| 2020-04-15 08:00 | 2020-04-15 08:10 | (NULL) | 11.00 |
| 2020-04-15 08:00 | 2020-04-15 08:10 | supplier2 | 5.00 |
| 2020-04-15 08:00 | 2020-04-15 08:10 | supplier1 | 6.00 |
| 2020-04-15 08:10 | 2020-04-15 08:20 | (NULL) | 10.00 |
| 2020-04-15 08:10 | 2020-04-15 08:20 | supplier2 | 9.00 |
| 2020-04-15 08:10 | 2020-04-15 08:20 | supplier1 | 1.00 |
+------------------+------------------+-------------+-------------+
GROUPING SETS의 각 하위 목록은 0개 이상의 컬럼이나 표현식을 지정할 수 있으며 GROUP BY 절에서 직접 사용된 것과 같은 방식으로 해석됩니다. 빈 그룹핑 세트는 모든 행이 단일 그룹으로 집계된다는 뜻이며, 이 그룹은 입력 행이 없어도 출력됩니다.
그룹핑 컬럼이나 표현식이 나타나지 않는 그룹핑 세트의 결과 행에서는 그룹핑 컬럼이나 표현식에 대한 참조가 null 값으로 대체됩니다.
ROLLUP
ROLLUP은 공통 유형의 그룹핑 세트를 지정하기 위한 축약 표기입니다. 주어진 표현식 목록과 목록의 모든 접두사(빈 목록 포함)를 나타냅니다.
ROLLUP이 있는 윈도우 집계는 window_start와 window_end 컬럼이 모두 GROUP BY 절에 있어야 하며, ROLLUP 절에는 있으면 안 됩니다.
예를 들어 다음 쿼리는 위의 것과 동등합니다.
SELECT window_start, window_end, supplier_id, SUM(price) AS total_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
GROUP BY window_start, window_end, ROLLUP (supplier_id);
CUBE
CUBE는 공통 유형의 그룹핑 세트를 지정하기 위한 축약 표기입니다. 주어진 목록과 그 가능한 모든 부분 집합, 즉 멱집합(power set)을 나타냅니다.
CUBE가 있는 윈도우 집계는 window_start와 window_end 컬럼이 모두 GROUP BY 절에 있어야 하며, CUBE 절에는 있으면 안 됩니다.
예를 들어 다음 두 쿼리는 동등합니다.
SELECT window_start, window_end, item, supplier_id, SUM(price) AS total_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
GROUP BY window_start, window_end, CUBE (supplier_id, item);
SELECT window_start, window_end, item, supplier_id, SUM(price) AS total_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
GROUP BY window_start, window_end, GROUPING SETS (
(supplier_id, item),
(supplier_id ),
( item),
( )
)
그룹 윈도우 시작·끝 타임스탬프 선택 (Selecting Group Window Start and End Timestamps)
그룹 윈도우의 시작과 끝 타임스탬프는 그룹화된 window_start와 window_end 컬럼으로 선택할 수 있습니다.
캐스케이딩 윈도우 집계 (Cascading Window Aggregation)
window_start와 window_end 컬럼은 일반 타임스탬프 컬럼이지 시간 속성이 아닙니다. 따라서 후속 시간 기반 연산에서 시간 속성으로 사용할 수 없습니다. 시간 속성을 전파하려면 GROUP BY 절에 window_time 컬럼을 추가로 포함해야 합니다. window_time은 Windowing TVF가 생성하는 세 번째 컬럼으로, 할당된 윈도우의 시간 속성입니다. GROUP BY 절에 window_time을 추가하면 window_time도 선택할 수 있는 그룹 키가 됩니다. 그러면 후속 쿼리가 이 컬럼을 캐스케이딩 윈도우 집계와 Window TopN 같은 후속 시간 기반 연산에 사용할 수 있습니다.
다음은 첫 번째 윈도우 집계가 두 번째 윈도우 집계를 위해 시간 속성을 전파하는 캐스케이딩 윈도우 집계를 보여줍니다.
-- tumbling 5 minutes for each supplier_id
CREATE VIEW window1 AS
-- Note: The window start and window end fields of inner Window TVF are optional in the select clause. However, if they appear in the clause, they need to be aliased to prevent name conflicting with the window start and window end of the outer Window TVF.
SELECT window_start AS window_5mintumble_start, window_end AS window_5mintumble_end, window_time AS rowtime, SUM(price) AS partial_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)
GROUP BY supplier_id, window_start, window_end, window_time;
-- tumbling 10 minutes on the first window
SELECT window_start, window_end, SUM(partial_price) AS total_price
FROM TUMBLE(TABLE window1, DESCRIPTOR(rowtime), INTERVAL '10' MINUTES)
GROUP BY window_start, window_end;
그룹 윈도우 집계 (Group Window Aggregation)
배치 / 스트리밍
경고: 그룹 윈도우 집계는 더 이상 사용되지 않습니다(deprecated). 더 강력하고 효과적인 Window TVF 집계를 사용하는 것이 권장됩니다.
Group Window Aggregation과 비교해 Window TVF Aggregation은 다음과 같은 많은 장점이 있습니다:
- Performance Tuning에서 언급된 모든 성능 최적화를 갖습니다.
- 표준
GROUPING SETS문법을 지원합니다.- 윈도우 집계 결과 후 Window TopN을 적용할 수 있습니다.
- 그 외 여러 가지.
그룹 윈도우 집계는 SQL 쿼리의 GROUP BY 절에 정의됩니다. 일반 GROUP BY 절이 있는 쿼리처럼, 그룹 윈도우 함수를 포함하는 GROUP BY 절이 있는 쿼리는 그룹당 하나의 결과 행을 계산합니다. 다음 그룹 윈도우 함수들이 배치 및 스트리밍 테이블의 SQL에 지원됩니다.
그룹 윈도우 함수 (Group Window Functions)
| 그룹 윈도우 함수 | 설명 |
|---|---|
TUMBLE(time_attr, interval) |
텀블링 시간 윈도우를 정의합니다. 텀블링 시간 윈도우는 고정된 기간(interval)을 가진 겹치지 않는 연속 윈도우에 행을 할당합니다. 예를 들어 5분 텀블링 윈도우는 5분 간격으로 행을 그룹화합니다. 텀블링 윈도우는 이벤트 시간(스트림 + 배치) 또는 처리 시간(스트림)에 정의할 수 있습니다. |
HOP(time_attr, interval, interval) |
호핑 시간 윈도우(Table API에서는 슬라이딩 윈도우라고 함)를 정의합니다. 호핑 시간 윈도우는 고정된 기간(두 번째 interval 파라미터)을 가지며 지정된 호프 간격(첫 번째 interval 파라미터)으로 호프합니다. 호프 간격이 윈도우 크기보다 작으면 호핑 윈도우는 겹칩니다. 따라서 행은 여러 윈도우에 할당될 수 있습니다. 예를 들어 크기 15분, 호프 간격 5분의 호핑 윈도우는 각 행을 5분 간격으로 평가되는 3개의 서로 다른 15분 크기 윈도우에 할당합니다. 호핑 윈도우는 이벤트 시간(스트림 + 배치) 또는 처리 시간(스트림)에 정의할 수 있습니다. |
SESSION(time_attr, interval) |
세션 시간 윈도우를 정의합니다. 세션 시간 윈도우는 고정된 기간을 가지지 않지만 그 경계가 비활동(inactivity)의 시간 interval로 정의됩니다. 즉, 정의된 갭 기간 동안 이벤트가 나타나지 않으면 세션 윈도우가 닫힙니다. 예를 들어 30분 갭이 있는 세션 윈도우는 30분 비활동 후 행이 관찰되면 시작되고(그렇지 않으면 그 행은 기존 윈도우에 추가됨), 30분 내에 행이 추가되지 않으면 닫힙니다. 세션 윈도우는 이벤트 시간(스트림 + 배치) 또는 처리 시간(스트림)에서 동작할 수 있습니다. |
시간 속성 (Time Attributes)
스트리밍 모드에서 그룹 윈도우 함수의 time_attr 인자는 행의 처리 시간 또는 이벤트 시간을 지정하는 유효한 시간 속성을 참조해야 합니다. 시간 속성을 정의하는 방법은 시간 속성 문서를 참고하세요.
배치 모드에서 그룹 윈도우 함수의 time_attr 인자는 TIMESTAMP 타입의 속성이어야 합니다.
그룹 윈도우 시작·끝 타임스탬프 선택 (Selecting Group Window Start and End Timestamps)
그룹 윈도우의 시작과 끝 타임스탬프 그리고 시간 속성은 다음 보조 함수로 선택할 수 있습니다:
| 보조 함수 | 설명 |
|---|---|
TUMBLE_START(time_attr, interval), HOP_START(time_attr, interval, interval), SESSION_START(time_attr, interval) |
해당 텀블링, 호핑 또는 세션 윈도우의 포함(inclusive) 하한 타임스탬프를 반환합니다. |
TUMBLE_END(time_attr, interval), HOP_END(time_attr, interval, interval), SESSION_END(time_attr, interval) |
해당 텀블링, 호핑 또는 세션 윈도우의 배타(exclusive) 상한 타임스탬프를 반환합니다. 참고: 배타적 상한 타임스탬프는 구간 조인, 그룹 윈도우 또는 over 윈도우 집계 같은 후속 시간 기반 연산에서 rowtime 속성으로 사용할 수 없습니다. |
TUMBLE_ROWTIME(time_attr, interval), HOP_ROWTIME(time_attr, interval, interval), SESSION_ROWTIME(time_attr, interval) |
해당 텀블링, 호핑 또는 세션 윈도우의 포함 상한 타임스탬프를 반환합니다. 결과 속성은 구간 조인, 그룹 윈도우 또는 over 윈도우 집계 같은 후속 시간 기반 연산에서 사용할 수 있는 rowtime 속성입니다. |
TUMBLE_PROCTIME(time_attr, interval), HOP_PROCTIME(time_attr, interval, interval), SESSION_PROCTIME(time_attr, interval) |
구간 조인, 그룹 윈도우 또는 over 윈도우 집계 같은 후속 시간 기반 연산에서 사용할 수 있는 proctime 속성을 반환합니다. |
참고: 보조 함수는
GROUP BY절의 그룹 윈도우 함수와 정확히 같은 인자로 호출되어야 합니다.
다음 예시들은 스트리밍 테이블에서 그룹 윈도우를 가진 SQL 쿼리를 지정하는 방법을 보여줍니다.
CREATE TABLE Orders (
user BIGINT,
product STRING,
amount INT,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '1' MINUTE
) WITH (...);
SELECT
user,
TUMBLE_START(order_time, INTERVAL '1' DAY) AS wStart,
SUM(amount) FROM Orders
GROUP BY
TUMBLE(order_time, INTERVAL '1' DAY),
user