Apache Flink — 스트림을 엔진으로 돌린다 · 창 TVF 와 창 Top-N · 실습
창 TVF 네 가지로 같은 주문을 자른다
목표
같은 주문 스트림에 TUMBLE·HOP·CUMULATE·SESSION 창을 걸어 창 열이 어떻게 붙고 한 행이 몇 개 창에 들어가는지 결과로 확인한다. 창 집계 위의 창 Top-N 과 창 TVF 바로 위의 창 Top-N 을 만든다.
왜 중요한가
창의 종류를 잘못 고르면 숫자가 조용히 틀린다 — HOP 의 건수를 합하면 겹친 만큼 부풀고, GROUP BY 에서 창 열을 빼면 창 집계가 아니라 무한 집계가 되어 갱신 로그가 섞인다. 창 TVF 는 행에 창 열을 붙이는 함수일 뿐이라, 그 열이 어떻게 붙는지만 정확히 알면 집계든 순위든 평범한 SQL 로 올릴 수 있다. 이 실습의 채점기는 클러스터에 묻지 않는다 — 여러분이 저장한 sql-client 출력을 읽고, 창마다의 기대값을 원본 CSV 에서 파이썬으로 직접 계산해 대조한다.
단계
1. flink-up 으로 클러스터를 띄우고, 10분 TUMBLE 이 붙인 order_id, ts, window_start, window_end, window_time 을 order_id <= 20 인 주문만 뽑는 /root/flink/windows/assign.sql 을 돌려 출력을 /root/flink/windows/assign.out 에 저장하세요.
2. 10분 TUMBLE 창·가게(shop)별 cnt(건수)·revenue(amount 합)를 내는 /root/flink/windows/tumble.sql 의 출력을 /root/flink/windows/tumble.out 에 저장하세요.
3. HOP(slide 5분, size 10분) 창마다 cnt 를 내는 /root/flink/windows/hop.sql 의 출력을 /root/flink/windows/hop.out 에 저장하세요.
4. CUMULATE(step 10분, size 1시간) 창마다 revenue 를 내는 /root/flink/windows/cumulate.sql 의 출력을 /root/flink/windows/cumulate.out 에 저장하세요.
5. SESSION(PARTITION BY shop, gap 5분) 창·가게별 cnt 를 내는 /root/flink/windows/session.sql 의 출력을 /root/flink/windows/session.out 에 저장하세요.
6. 10분 TUMBLE 창마다 가게별 매출 상위 2곳(window_start, window_end, shop, revenue, rownum)을 내는 /root/flink/windows/top-shops.sql 의 출력을 /root/flink/windows/top-shops.out 에 저장하세요.
7. 30분 TUMBLE 창마다 금액이 큰 주문 3건(order_id, shop, amount, window_start, window_end, rownum)을 집계 없이 내는 /root/flink/windows/top-orders.sql 의 출력을 /root/flink/windows/top-orders.out 에 저장하세요.
8. /root/flink/windows/report.json 에 orders·hop_assignments·cumulate_windows·sessions·max_session_orders 를 적으세요.
참고
- 원본:
/opt/lab/fixtures/data/windows_orders.csv, 열order_id BIGINT, shop STRING, amount INT, ts TIMESTAMP(3)(머리글 없는 CSV, ts 오름차순). 표에는WATERMARK FOR ts AS ts - INTERVAL '1' SECOND를 두고 스트리밍 모드로 돌립니다. - 모양:
FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE). HOP 은(TABLE, DESCRIPTOR, slide, size), CUMULATE 는(TABLE, DESCRIPTOR, step, size), SESSION 은(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), gap). - 창 집계는
GROUP BY window_start, window_end, ...로 묶습니다. 창 열을 빼면 무한 집계가 되어 -U/+U 가 섞입니다. - 흔한 실수: HOP·CUMULATE 의 인자 순서(작은 값이 먼저)를 바꾸는 것. size 가 slide(step)의 정수배가 아니라는 오류로 거절됩니다.
- 흔한 실수: 창 Top-N 의 PARTITION BY 에 window_start, window_end 를 빼는 것 — 그러면 창 Top-N 이 아니라 일반 Top-N 이 되어 순위가 바뀔 때마다 로그가 나옵니다.
- 공식 문서: [Windowing TVF](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/window-tvf/) · [Window Aggregation](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/window-agg/) · [Window Top-N](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/window-topn/) · [Time Attributes](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/sql-table-concepts/time_attributes/)
8단계
- 창 TVF 가 붙이는 열 세 개
- 가게별 10분 집계
- HOP — 한 주문이 두 창에 들어간다
- CUMULATE — 시작이 고정된 누적 창
- SESSION — 가게마다 길이가 다른 창
- 창 집계 위의 창 Top-N
- 창 TVF 바로 위의 창 Top-N
- 보고서 — 창마다 몇 번 세어지나