Apache Flink — 스트림을 엔진으로 돌린다 · 창 TVF 와 창 Top-N · 이론
창 TVF — 행에 창 열 세 개를 붙이는 함수
한 줄 요약
Flink SQL 의 창은 표를 받아 표를 돌려주는 함수(창 TVF) 다. TUMBLE·HOP·CUMULATE·SESSION 은 원래 행에 window_start·window_end·window_time 세 열을 붙여 돌려줄 뿐이고, 집계와 Top-N 은 그 열로 묶는 평범한 SQL 이다. 창의 종류가 정하는 것은 단 하나 — 한 행이 몇 개의, 어떤 창에 들어가는가 다.
왜 이게 필요했나
앞 모듈에서 본 대로 GROUP BY user_id 같은 무한 집계는 결과를 끝없이 고쳐 쓰고, 키마다 상태를 영원히 쥔다. 대시보드가 원하는 것은 대개 그게 아니다. "10분마다 가게별 매출", "5분마다 갱신되는 최근 10분 주문 수", "오늘 0시부터 지금까지의 누적 매출", "한 번 들어와서 나갈 때까지의 세션" — 모두 시간으로 자른 구간 위의 집계다. 구간이 닫히면 결과를 한 번 내고 상태를 버리면 된다.
예전 Flink SQL 은 이것을 GROUP BY TUMBLE(ts, ...) 같은 특수 문법(Grouped Window Functions)으로 했다. 집계에는 쓸 수 있었지만, 창마다 순위를 매기거나 두 스트림을 같은 창으로 조인하는 일에는 쓸 수 없었다. 공식 문서(Windowing TVF)는 창 TVF 가 그 문법을 대체한다고 적는다. 창을 "행에 열을 붙이는 함수" 로 바꾸자 창 위에 무엇이든 올릴 수 있게 됐다 — 창 집계, 창 Top-N, 창 조인, 창 중복 제거.
어떻게 동작하나
창 TVF 는 FROM 자리에 쓴다. 첫 인자는 표, 둘째는 시간 속성 열, 나머지는 크기다.
SELECT window_start, window_end, shop, COUNT(*) AS cntFROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)GROUP BY window_start, window_end, shop;| TVF | 인자 | 한 행이 들어가는 창 수 | 겹침 |
| --- | --- | --- | --- |
| TUMBLE | size | 1 | 없음 |
| HOP | slide, size | size / slide | 있음 |
| CUMULATE | step, size | 시작이 같은 창 중 끝이 그 행보다 뒤인 것 전부 | 있음 |
| SESSION | (PARTITION BY 키) gap | 1 | 없음, 크기가 제각각 |
몇 가지 규칙이 결과를 정한다.
창은 반열린 구간이다. [window_start, window_end) — 09:10:00 에 딱 찍힌 주문은 [09:00, 09:10) 이 아니라 [09:10, 09:20) 에 들어간다. window_time 은 문서대로 늘 window_end − 1ms 이고, 스트리밍에서는 이 열이 시간 속성으로 남아 다음 창 연산에 쓸 수 있다. 반대로 window_start·window_end 는 평범한 타임스탬프가 되어 시간 속성이 아니다.
창의 시작은 에포크에 맞춰진다. 10분 창은 정각 기준 00·10·20분에 시작하고, 첫 행이 09:01:35 에 와도 창은 09:00 에 시작한다. 이것을 옮기는 것이 선택 인자 offset 이다.
HOP 과 CUMULATE 는 인자 순서가 함정이다. 둘 다 작은 값(slide, step)이 먼저, 큰 값(size)이 나중이다. 반대로 쓰면 잡이 뜨기 전에 거절된다 — 실습 환경에서 HOP 은 'size must be an integral multiple of slide', CUMULATE 는 'maxSize must be an integral multiple of step' 오류를 냈다. 즉 size 는 slide(step)의 정수배여야 한다. HOP(5분, 10분)에서 한 행은 창 두 개에 들어가므로 창별 건수를 다 더하면 원래 행 수의 두 배가 된다. 이것은 버그가 아니라 정의다. CUMULATE 는 문서 표현대로 "size 로 TUMBLE 한 뒤 그 안을 step 마다 끝이 늘어나는 창으로 나눈 것" 이라, 한 시간 창에 step 10분이면 시작이 같은 창이 여섯 개 나온다.
SESSION 은 크기가 없다. 같은 키의 이웃한 행 간격이 gap 이하면 한 세션으로 이어지고, 세션의 끝은 마지막 행 + gap 이다. 그래서 세션마다 길이가 다르다. 문서 기준으로 SESSION TVF 는 아직 배치 모드를 지원하지 않는다.
창 연산은 끝에 한 번만 낸다. 문서(Window Aggregation)는 창 집계가 중간 결과를 내지 않고 창이 끝날 때 최종 결과만 내며, 필요 없어진 상태를 모두 지운다고 적는다. 그래서 결과 로그에는 +I 만 있다. 창이 "끝났다" 는 판단은 앞 모듈의 워터마크가 한다.
창 Top-N 은 창 열로 나눈다. ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY ...) 처럼 PARTITION BY 에 창 열 두 개가 있어야 옵티마이저가 창 Top-N 으로 번역한다. 그러면 일반 Top-N 과 달리 순위가 바뀔 때마다 -U/+U 를 내지 않고 창이 닫힐 때 상위 N 개만 한 번 낸다. 창 집계 위에 올릴 수도, 창 TVF 바로 위에 올려 행 자체의 순위를 매길 수도 있다. 문서 기준으로 TVF 바로 위의 창 Top-N 은 TUMBLE·HOP·CUMULATE 만 지원한다.
현장에서 만나는 모습
"5분마다 최근 1시간" 대시보드를 HOP 으로 만들면 한 행이 창 12개에 들어간다. 상태와 계산이 그만큼 늘고, 창별 건수를 합해 "총 주문 수" 로 쓰면 12배로 부풀어 있다. 누적 지표("오늘 지금까지")를 HOP 으로 흉내 내는 경우도 흔한데, 그 자리는 시작이 고정된 CUMULATE 가 맞다.
두 번째는 순위의 흔들림이다. 금액이 같은 주문 둘이 3위를 다투면 ROW_NUMBER 는 둘 중 하나를 고른다. 어느 쪽인지 보장하지 않으므로 재처리하면 달라질 수 있다. 정렬 키에 주문 번호 같은 보조 키를 더해 동점을 깨 두는 습관이 필요하다. 이 실습의 자료는 금액과 창별 매출에 동점이 없게 만들어 두었다.
세 번째는 시간대다. 이 실습은 TIMESTAMP(3) 라 창이 적힌 시각 그대로 잘린다. TIMESTAMP_LTZ 열로 하루 창을 만들면 세션 시간대에 따라 "하루" 의 경계가 바뀐다. 일별 창이 한국 시각 오전 9시에 잘려 있다면 이것부터 본다.
다음 실습에서 할 것
주문 119건에 10분 TUMBLE 을 걸어 창 열이 어떻게 붙는지 보고, 가게별 10분 집계를 낸다. HOP(5분, 10분)의 건수 합이 두 배가 되는 것, CUMULATE(10분, 1시간)가 시작이 같은 창 여섯 개를 내는 것, SESSION(간격 5분)이 가게마다 길이가 다른 세션을 만드는 것을 확인한다. 10분 창마다 매출 상위 두 가게와 30분 창마다 큰 주문 세 건을 창 Top-N 으로 뽑고, 숫자를 보고서로 정리한다.