윈도우 테이블 함수
윈도우 테이블 함수 (Windowing TVFs)
Windows 는 무한 스트림 처리의 핵심입니다. Windows 는 스트림을 계산을 적용할 수 있는 유한 크기의 "버킷(buckets)" 으로 분할합니다. 이 문서는 Flink SQL 에서 윈도우잉이 수행되는 방식과 프로그래머가 그 제공 기능을 최대한 활용할 수 있는 방법에 초점을 맞춥니다.
출처: 문서
본문
이 절은 Batch 와 Streaming 모드에서 모두 지원됩니다.
Windows 는 무한 스트림 처리의 핵심입니다. Windows 는 스트림을 계산을 적용할 수 있는 유한 크기의 "버킷(buckets)" 으로 분할합니다. 이 문서는 Flink SQL 에서 윈도우잉이 수행되는 방식과 프로그래머가 그 제공 기능을 최대한 활용할 수 있는 방법에 초점을 맞춥니다.
Apache Flink 는 테이블의 요소를 window 로 나누는 여러 윈도우 테이블 함수(TVF) 를 제공합니다:
- Tumble Windows
- Hop Windows
- Cumulate Windows
- Session Windows (현재 스트리밍 모드에서만 지원)
사용하는 윈도우 테이블 함수에 따라 각 요소는 논리적으로 둘 이상의 window 에 속할 수 있습니다. 예를 들어 HOP 윈도우잉은 단일 요소가 여러 window 에 할당될 수 있는 겹치는 window 를 만듭니다.
윈도우 TVF 는 Flink 가 정의한 다형성 테이블 함수(Polymorphic Table Functions, PTF) 입니다. PTF 는 SQL 2016 표준의 일부로, 테이블을 매개변수로 가질 수 있는 특별한 테이블 함수입니다. PTF 는 테이블의 형태를 바꾸는 강력한 기능입니다. PTF 는 의미상 테이블처럼 사용되므로 호출은 SELECT 문의 FROM 절에서 발생합니다.
윈도우 TVF 는 레거시 Grouped Window Functions 의 대체입니다. 윈도우 TVF 는 SQL 표준에 더 부합하고 Window TopN, Window Join 같은 복잡한 window 기반 계산을 지원하는 데 더 강력합니다. 그러나 Grouped Window Functions 은 Window Aggregation 만 지원할 수 있습니다.
윈도우 TVF 를 기반으로 추가 계산을 적용하는 방법은 다음을 참고하세요:
윈도우 함수 (Window Functions)
Apache Flink 는 4개의 내장 윈도우 TVF 를 제공합니다: TUMBLE, HOP, CUMULATE, SESSION. 윈도우 TVF 의 반환 값은 원래 관계의 모든 열과 할당된 window 를 나타내는 "window_start", "window_end", "window_time" 이라는 추가 3개 열을 포함하는 새 관계입니다. 스트리밍 모드에서 "window_time" 필드는 window 의 시간 속성(time attributes) 입니다. 배치 모드에서 "window_time" 필드는 입력 시간 필드 유형에 따라 TIMESTAMP 또는 TIMESTAMP_LTZ 유형의 속성입니다. "window_time" 필드는 이후의 시간 기반 연산(예: 다른 윈도우 TVF, interval joins, over aggregations) 에서 사용할 수 있습니다. window_time 의 값은 항상 window_end - 1ms 와 같습니다.
TUMBLE
TUMBLE 함수는 각 요소를 지정된 window size 의 window 에 할당합니다. Tumbling windows 는 고정된 크기를 가지며 겹치지 않습니다. 예를 들어 크기가 5분인 tumbling window 를 지정한다고 가정해 보세요. 이 경우 Flink 는 현재 window 를 평가하고 5분마다 새 window 를 시작합니다.
TUMBLE 함수는 시간 속성 필드를 기반으로 관계의 각 행에 window 를 할당합니다. 스트리밍 모드에서 시간 속성 필드는 이벤트 또는 프로세싱 시간 속성 중 하나여야 합니다. 배치 모드에서 윈도우 테이블 함수의 시간 속성 필드는 TIMESTAMP 또는 TIMESTAMP_LTZ 유형의 속성이어야 합니다. TUMBLE 의 반환 값은 원래 관계의 모든 열과 할당된 window 를 나타내는 "window_start", "window_end", "window_time" 이라는 추가 3개 열을 포함하는 새 관계입니다. 원래 시간 속성 "timecol" 은 윈도우 TVF 후 일반 타임스탬프 열이 됩니다.
TUMBLE 함수는 세 개의 필수 매개변수와 하나의 선택 매개변수를 받습니다:
TUMBLE(TABLE data, DESCRIPTOR(timecol), size [, offset ])
data: 시간 속성 열이 있는 모든 관계가 될 수 있는 테이블 매개변수.timecol: 데이터의 어떤 시간 속성 열이 tumbling windows 에 매핑될지 나타내는 열 설명자.size: tumbling windows 의 너비를 지정하는 기간(duration).offset: window 시작이 이동될 오프셋을 지정하는 선택 매개변수.
Bid 테이블에 대한 호출 예:
-- 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 TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES);
-- or with the named params
-- note: the DATA param must be the first
Flink SQL> SELECT * FROM
TUMBLE(
DATA => TABLE Bid,
TIMECOL => DESCRIPTOR(bidtime),
SIZE => INTERVAL '10' MINUTES);
+------------------+-------+------+------------------+------------------+-------------------------+
| bidtime | price | item | window_start | window_end | window_time |
+------------------+-------+------+------------------+------------------+-------------------------+
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:07 | 2.00 | A | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:09 | 5.00 | D | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:17 | 6.00 | F | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
+------------------+-------+------+------------------+------------------+-------------------------+
-- apply aggregation on the tumbling windowed table
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 |
+------------------+------------------+-------------+
참고: 윈도우잉의 동작을 더 잘 이해하기 위해 타임스탬프 값의 표시를 단순화하여 뒤의 0을 표시하지 않습니다. 예:
2020-04-15 08:05는 유형이TIMESTAMP(3)이면 Flink SQL Client 에서2020-04-15 08:05:00.000으로 표시되어야 합니다.
HOP
HOP 함수는 요소를 고정 길이의 window 에 할당합니다. TUMBLE 윈도우 함수와 마찬가지로 window 의 크기는 window size 매개변수로 구성됩니다. 추가 윈도우 slide 매개변수는 hopping window 가 얼마나 자주 시작되는지 제어합니다. 따라서 slide 가 window size 보다 작으면 hopping windows 는 겹칠 수 있습니다. 이 경우 요소는 여러 window 에 할당됩니다. Hopping windows 는 "sliding windows" 라고도 합니다.
예를 들어 10분 크기이고 5분씩 슬라이드하는 window 를 가질 수 있습니다. 이를 통해 5분마다 지난 10분 동안 도착한 이벤트를 포함하는 window 를 얻습니다.
HOP 함수는 시간 속성 필드를 기반으로 size 의 간격 내의 행을 덮고 매 slide 마다 이동하는 window 를 할당합니다. 스트리밍 모드에서 시간 속성 필드는 이벤트 또는 프로세싱 시간 속성 중 하나여야 합니다. 배치 모드에서 윈도우 테이블 함수의 시간 속성 필드는 TIMESTAMP 또는 TIMESTAMP_LTZ 유형의 속성이어야 합니다. HOP 의 반환 값은 원래 관계의 모든 열과 할당된 window 를 나타내는 "window_start", "window_end", "window_time" 이라는 추가 3개 열을 포함하는 새 관계입니다. 원래 시간 속성 "timecol" 은 윈도우 TVF 후 일반 타임스탬프 열이 됩니다.
HOP 은 네 개의 필수 매개변수와 하나의 선택 매개변수를 받습니다:
HOP(TABLE data, DESCRIPTOR(timecol), slide, size [, offset ])
data: 시간 속성 열이 있는 모든 관계가 될 수 있는 테이블 매개변수.timecol: 데이터의 어떤 시간 속성 열이 hopping windows 에 매핑될지 나타내는 열 설명자.slide: 순차 hopping windows 시작 사이의 기간을 지정하는 기간.size: hopping windows 의 너비를 지정하는 기간.offset: window 시작이 이동될 오프셋을 지정하는 선택 매개변수.
Bid 테이블에 대한 호출 예:
> SELECT * FROM HOP(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES, INTERVAL '10' MINUTES);
-- or with the named params
-- note: the DATA param must be the first
> SELECT * FROM
HOP(
DATA => TABLE Bid,
TIMECOL => DESCRIPTOR(bidtime),
SLIDE => INTERVAL '5' MINUTES,
SIZE => INTERVAL '10' MINUTES);
+------------------+-------+------+------------------+------------------+-------------------------+
| bidtime | price | item | window_start | window_end | window_time |
+------------------+-------+------+------------------+------------------+-------------------------+
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:05 | 2020-04-15 08:15 | 2020-04-15 08:14:59.999 |
| 2020-04-15 08:07 | 2.00 | A | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:07 | 2.00 | A | 2020-04-15 08:05 | 2020-04-15 08:15 | 2020-04-15 08:14:59.999 |
| 2020-04-15 08:09 | 5.00 | D | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:09 | 5.00 | D | 2020-04-15 08:05 | 2020-04-15 08:15 | 2020-04-15 08:14:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:05 | 2020-04-15 08:15 | 2020-04-15 08:14:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:05 | 2020-04-15 08:15 | 2020-04-15 08:14:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:17 | 6.00 | F | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:17 | 6.00 | F | 2020-04-15 08:15 | 2020-04-15 08:25 | 2020-04-15 08:24:59.999 |
+------------------+-------+------+------------------+------------------+-------------------------+
-- apply aggregation on the hopping windowed table
> 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 |
+------------------+------------------+-------------+
CUMULATE
누적(Cumulating) window 는 고정 window 간격에서 조기 발화(early firing)가 있는 tumbling window 같은 일부 시나리오에서 매우 유용합니다. 예를 들어 일일 대시보드가 00:00 부터 매분까지 누적 UV 를 그리는 경우, 10:00 의 UV 는 00:00 부터 10:00 까지의 총 UV 수를 나타냅니다. 이는 CUMULATE 윈도우잉으로 쉽고 효율적으로 구현할 수 있습니다.
CUMULATE 함수는 step 크기의 초기 간격 내의 행을 덮는 window 에 요소를 할당하고, 최대 window size 까지 매 step 마다 step size 하나를 더 확장합니다(window 시작은 고정). CUMULATE 함수를 먼저 최대 window size 로 TUMBLE 윈도우잉을 적용하고, 각 tumbling window 를 같은 window 시작과 step-size 차이의 window 끝을 가진 여러 window 로 분할하는 것으로 생각할 수 있습니다. 따라서 누적 window 는 겹치며 고정된 크기가 없습니다.
예를 들어 1시간 step 이고 1일 최대 크기인 누적 window 를 가질 수 있으며, 매일 [00:00, 01:00), [00:00, 02:00), [00:00, 03:00), …, [00:00, 24:00) window 를 얻습니다.
CUMULATE 함수는 시간 속성 열을 기반으로 window 를 할당합니다. 스트리밍 모드에서 시간 속성 필드는 이벤트 또는 프로세싱 시간 속성 중 하나여야 합니다. 배치 모드에서 윈도우 테이블 함수의 시간 속성 필드는 TIMESTAMP 또는 TIMESTAMP_LTZ 유형의 속성이어야 합니다. CUMULATE 의 반환 값은 원래 관계의 모든 열과 할당된 window 를 나타내는 "window_start", "window_end", "window_time" 이라는 추가 3개 열을 포함하는 새 관계입니다. 원래 시간 속성 "timecol" 은 윈도우 TVF 후 일반 타임스탬프 열이 됩니다.
CUMULATE 는 네 개의 필수 매개변수와 하나의 선택 매개변수를 받습니다:
CUMULATE(TABLE data, DESCRIPTOR(timecol), step, size)
data: 시간 속성 열이 있는 모든 관계가 될 수 있는 테이블 매개변수.timecol: 데이터의 어떤 시간 속성 열이 cumulating windows 에 매핑될지 나타내는 열 설명자.step: 순차 cumulating windows 의 끝 사이의 증가된 window 크기를 지정하는 기간.size: cumulating windows 의 최대 너비를 지정하는 기간.size는step의 정수 배수여야 합니다.offset: window 시작이 이동될 오프셋을 지정하는 선택 매개변수.
Bid 테이블에 대한 호출 예:
> SELECT * FROM
CUMULATE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '2' MINUTES, INTERVAL '10' MINUTES);
-- or with the named params
-- note: the DATA param must be the first
> SELECT * FROM
CUMULATE(
DATA => TABLE Bid,
TIMECOL => DESCRIPTOR(bidtime),
STEP => INTERVAL '2' MINUTES,
SIZE => INTERVAL '10' MINUTES);
+------------------+-------+------+------------------+------------------+-------------------------+
| bidtime | price | item | window_start | window_end | window_time |
+------------------+-------+------+------------------+------------------+-------------------------+
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:00 | 2020-04-15 08:06 | 2020-04-15 08:05:59.999 |
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:00 | 2020-04-15 08:08 | 2020-04-15 08:07:59.999 |
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:07 | 2.00 | A | 2020-04-15 08:00 | 2020-04-15 08:08 | 2020-04-15 08:07:59.999 |
| 2020-04-15 08:07 | 2.00 | A | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:09 | 5.00 | D | 2020-04-15 08:00 | 2020-04-15 08:10 | 2020-04-15 08:09:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:12 | 2020-04-15 08:11:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:16 | 2020-04-15 08:15:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:18 | 2020-04-15 08:17:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:10 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:10 | 2020-04-15 08:16 | 2020-04-15 08:15:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:10 | 2020-04-15 08:18 | 2020-04-15 08:17:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
| 2020-04-15 08:17 | 6.00 | F | 2020-04-15 08:10 | 2020-04-15 08:18 | 2020-04-15 08:17:59.999 |
| 2020-04-15 08:17 | 6.00 | F | 2020-04-15 08:10 | 2020-04-15 08:20 | 2020-04-15 08:19:59.999 |
+------------------+-------+------+------------------+------------------+-------------------------+
-- apply aggregation on the cumulating windowed table
> 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
참고:
- Session Window TVF 는 현재 배치 모드에서 지원되지 않습니다.
- Session Window Aggregation 은 현재 Performance Tuning 의 어떤 최적화도 지원하지 않습니다.
- Session Window Join, Session Window TopN, Session Window Deduplication 은 개념적으로 지원되며 베타 모드입니다. 문제는 JIRA 에 보고할 수 있습니다.
SESSION 함수는 활동의 세션(session) 으로 요소를 그룹화합니다. TUMBLE window 및 HOP window 와 달리 session windows 는 겹치지 않으며 고정된 시작과 끝 시간이 없습니다. 대신 session window 는 일정 기간 동안 요소를 받지 않을 때, 즉 비활동의 간격이 발생했을 때 닫힙니다. session window 는 비활동 기간이 얼마나 되는지 정의하는 정적 session gap 으로 구성되어야 합니다. 이 기간이 만료되면 현재 session 이 닫히고 후속 요소가 새 session window 에 할당됩니다.
예를 들어 gap 이 10분인 window 를 가질 수 있습니다. 이를 통해 같은 사용자의 두 이벤트 사이의 간격이 10분 미만이면 이 이벤트들은 같은 session window 로 그룹화됩니다. 최신 이벤트 후 10분 동안 데이터가 없으면 이 session window 는 닫히고 다운스트림으로 전송됩니다. 후속 이벤트는 새 session window 에 할당됩니다.
SESSION 함수는 날짜/시간에 기반해 행을 덮는 window 를 할당합니다. 스트리밍 모드에서 시간 속성 필드는 이벤트 또는 프로세싱 시간 속성 중 하나여야 합니다. SESSION 의 반환 값은 원래 관계의 모든 열과 할당된 window 를 나타내는 "window_start", "window_end", "window_time" 이라는 추가 3개 열을 포함하는 새 관계입니다. 원래 시간 속성 "timecol" 은 윈도우 TVF 후 일반 타임스탬프 열이 됩니다.
SESSION 은 세 개의 필수 매개변수와 하나의 선택 매개변수를 받습니다:
SESSION(TABLE data [PARTITION BY(keycols, ...)], DESCRIPTOR(timecol), gap)
data: 시간 속성 열이 있는 모든 관계가 될 수 있는 테이블 매개변수.keycols: session windows 전에 데이터를 파티셔닝하는 데 사용할 열을 나타내는 열 설명자.timecol: 데이터의 어떤 시간 속성 열이 session windows 에 매핑될지 나타내는 열 설명자.gap: 두 이벤트가 같은 session window 의 일부로 간주되기 위한 최대 타임스탬프 간격.
Bid 테이블에 대한 호출 예:
-- 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:07 | 4.00 | A |
| 2020-04-15 08:06 | 2.00 | A |
| 2020-04-15 08:09 | 5.00 | B |
| 2020-04-15 08:08 | 3.00 | A |
| 2020-04-15 08:17 | 1.00 | B |
+------------------+-------+------+
-- session window with partition keys
> SELECT * FROM SESSION(TABLE Bid PARTITION BY item, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES);
-- or with the named params
-- note: the DATA param must be the first
> SELECT * FROM
SESSION(
DATA => TABLE Bid PARTITION BY item,
TIMECOL => DESCRIPTOR(bidtime),
GAP => INTERVAL '5' MINUTES);
+------------------+-------+------+------------------+------------------+-------------------------+
| bidtime | price | item | window_start | window_end | window_time |
+------------------+-------+------+------------------+------------------+-------------------------+
| 2020-04-15 08:07 | 4.00 | A | 2020-04-15 08:06 | 2020-04-15 08:13 | 2020-04-15 08:12:59.999 |
| 2020-04-15 08:06 | 2.00 | A | 2020-04-15 08:06 | 2020-04-15 08:13 | 2020-04-15 08:12:59.999 |
| 2020-04-15 08:08 | 3.00 | A | 2020-04-15 08:06 | 2020-04-15 08:13 | 2020-04-15 08:12:59.999 |
| 2020-04-15 08:09 | 5.00 | B | 2020-04-15 08:09 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:17 | 1.00 | B | 2020-04-15 08:17 | 2020-04-15 08:22 | 2020-04-15 08:21:59.999 |
+------------------+-------+------+------------------+------------------+-------------------------+
-- apply aggregation on the session windowed table with partition keys
> SELECT window_start, window_end, item, SUM(price) AS total_price
FROM SESSION(TABLE Bid PARTITION BY item, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)
GROUP BY item, window_start, window_end;
+------------------+------------------+------+-------------+
| window_start | window_end | item | total_price |
+------------------+------------------+------+-------------+
| 2020-04-15 08:06 | 2020-04-15 08:13 | A | 9.00 |
| 2020-04-15 08:09 | 2020-04-15 08:14 | B | 5.00 |
| 2020-04-15 08:17 | 2020-04-15 08:22 | B | 1.00 |
+------------------+------------------+------+-------------+
-- session window without partition keys
> SELECT * FROM SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES);
-- or with the named params
-- note: the DATA param must be the first
> SELECT * FROM
SESSION(
DATA => TABLE Bid,
TIMECOL => DESCRIPTOR(bidtime),
GAP => INTERVAL '5' MINUTES);
+------------------+-------+------+------------------+------------------+-------------------------+
| bidtime | price | item | window_start | window_end | window_time |
+------------------+-------+------+------------------+------------------+-------------------------+
| 2020-04-15 08:07 | 4.00 | A | 2020-04-15 08:06 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:06 | 2.00 | A | 2020-04-15 08:06 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:08 | 3.00 | A | 2020-04-15 08:06 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:09 | 5.00 | B | 2020-04-15 08:06 | 2020-04-15 08:14 | 2020-04-15 08:13:59.999 |
| 2020-04-15 08:17 | 1.00 | B | 2020-04-15 08:17 | 2020-04-15 08:22 | 2020-04-15 08:21:59.999 |
+------------------+-------+------+------------------+------------------+-------------------------+
-- apply aggregation on the session windowed table without partition keys
> SELECT window_start, window_end, SUM(price) AS total_price
FROM SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)
GROUP BY window_start, window_end;
+------------------+------------------+-------------+
| window_start | window_end | total_price |
+------------------+------------------+-------------+
| 2020-04-15 08:06 | 2020-04-15 08:14 | 14.00 |
| 2020-04-15 08:17 | 2020-04-15 08:22 | 1.00 |
+------------------+------------------+-------------+
Window Offset
Offset 은 window 할당을 변경하는 데 사용할 수 있는 선택 매개변수입니다. 양수 기간 또는 음수 기간일 수 있습니다. window offset 의 기본값은 0 입니다. 다른 offset 값을 설정하면 같은 레코드가 다른 window 에 할당될 수 있습니다.
예를 들어 타임스탬프 2021-06-30 00:00:04 의 레코드가 size 가 10 MINUTE 인 Tumble window 에서 어떤 window 에 할당될까요?
offset값이-16 MINUTE이면 레코드는 window [2021-06-29 23:54:00,2021-06-30 00:04:00) 에 할당됩니다.offset값이-6 MINUTE이면 레코드는 window [2021-06-29 23:54:00,2021-06-30 00:04:00) 에 할당됩니다.offset이-4 MINUTE이면 레코드는 window [2021-06-29 23:56:00,2021-06-30 00:06:00) 에 할당됩니다.offset이0이면 레코드는 window [2021-06-30 00:00:00,2021-06-30 00:10:00) 에 할당됩니다.offset이4 MINUTE이면 레코드는 window [2021-06-29 23:54:00,2021-06-30 00:04:00) 에 할당됩니다.offset이6 MINUTE이면 레코드는 window [2021-06-29 23:56:00,2021-06-30 00:06:00) 에 할당됩니다.offset이16 MINUTE이면 레코드는 window [2021-06-29 23:56:00,2021-06-30 00:06:00) 에 할당됩니다.
일부 window offset 매개변수는 window 할당에 같은 효과를 가질 수 있음을 알 수 있습니다. 위의 경우 size 가 10 MINUTE 인 Tumble window 에 -16 MINUTE, -6 MINUTE, 4 MINUTE 는 같은 효과를 가집니다.
참고: window offset 의 효과는 window 할당을 업데이트하기 위한 것일 뿐이며 Watermark 에는 영향을 주지 않습니다.
다음 SQL 에서 Tumble window 에서 offset 을 사용하는 방법을 설명하는 예를 보여줍니다.
-- NOTE: Currently Flink doesn't support evaluating individual window table-valued function,
-- window table-valued function should be used with aggregate operation,
-- this example is just used for explaining the syntax and the data produced by table-valued function.
Flink SQL> SELECT * FROM
TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES, INTERVAL '1' MINUTES);
-- or with the named params
-- note: the DATA param must be the first and `OFFSET` should be wrapped with double quotes
Flink SQL> SELECT * FROM
TUMBLE(
DATA => TABLE Bid,
TIMECOL => DESCRIPTOR(bidtime),
SIZE => INTERVAL '10' MINUTES,
`OFFSET` => INTERVAL '1' MINUTES);
+------------------+-------+------+------------------+------------------+-------------------------+
| bidtime | price | item | window_start | window_end | window_time |
+------------------+-------+------+------------------+------------------+-------------------------+
| 2020-04-15 08:05 | 4.00 | C | 2020-04-15 08:01 | 2020-04-15 08:11 | 2020-04-15 08:10:59.999 |
| 2020-04-15 08:07 | 2.00 | A | 2020-04-15 08:01 | 2020-04-15 08:11 | 2020-04-15 08:10:59.999 |
| 2020-04-15 08:09 | 5.00 | D | 2020-04-15 08:01 | 2020-04-15 08:11 | 2020-04-15 08:10:59.999 |
| 2020-04-15 08:11 | 3.00 | B | 2020-04-15 08:11 | 2020-04-15 08:21 | 2020-04-15 08:20:59.999 |
| 2020-04-15 08:13 | 1.00 | E | 2020-04-15 08:11 | 2020-04-15 08:21 | 2020-04-15 08:20:59.999 |
| 2020-04-15 08:17 | 6.00 | F | 2020-04-15 08:11 | 2020-04-15 08:21 | 2020-04-15 08:20:59.999 |
+------------------+-------+------+------------------+------------------+-------------------------+
-- apply aggregation on the tumbling windowed table
Flink SQL> SELECT window_start, window_end, SUM(price) AS total_price
FROM TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES, INTERVAL '1' MINUTES)
GROUP BY window_start, window_end;
+------------------+------------------+-------------+
| window_start | window_end | total_price |
+------------------+------------------+-------------+
| 2020-04-15 08:01 | 2020-04-15 08:11 | 11.00 |
| 2020-04-15 08:11 | 2020-04-15 08:21 | 10.00 |
+------------------+------------------+-------------+
참고: 윈도우잉의 동작을 더 잘 이해하기 위해 타임스탬프 값의 표시를 단순화하여 뒤의 0을 표시하지 않습니다. 예:
2020-04-15 08:05는 유형이TIMESTAMP(3)이면 Flink SQL Client 에서2020-04-15 08:05:00.000으로 표시되어야 합니다.