Apache Flink — Running Streams on a Real Engine
Slice the Same Orders with Four Window TVFs
한국어 원문으로 표시합니다.
목표
같은 주문 스트림에 TUMBLE·HOP·CUMULATE·SESSION 창을 걸어 창 열이 어떻게 붙고 한 행이 몇 개 창에 들어가는지 결과로 확인한다. 창 집계 위의 창 Top-N 과 창 TVF 바로 위의 창 Top-N 을 만든다.
왜 중요한가
창의 종류를 잘못 고르면 숫자가 조용히 틀린다 — HOP 의 건수를 합하면 겹친 만큼 부풀고, GROUP BY 에서 창 열을 빼면 창 집계가 아니라 무한 집계가 되어 갱신 로그가 섞인다. 창 TVF 는 행에 창 열을 붙이는 함수일 뿐이라, 그 열이 어떻게 붙는지만 정확히 알면 집계든 순위든 평범한 SQL 로 올릴 수 있다. 이 실습의 채점기는 클러스터에 묻지 않는다 — 여러분이 저장한 sql-client 출력을 읽고, 창마다의 기대값을 원본 CSV 에서 파이썬으로 직접 계산해 대조한다.
단계
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 에 저장하세요.- 10분
TUMBLE창·가게(shop)별cnt(건수)·revenue(amount 합)를 내는 /root/flink/windows/tumble.sql 의 출력을 /root/flink/windows/tumble.out 에 저장하세요. HOP(slide 5분, size 10분) 창마다cnt를 내는 /root/flink/windows/hop.sql 의 출력을 /root/flink/windows/hop.out 에 저장하세요.CUMULATE(step 10분, size 1시간) 창마다revenue를 내는 /root/flink/windows/cumulate.sql 의 출력을 /root/flink/windows/cumulate.out 에 저장하세요.SESSION(PARTITION BY shop, gap 5분) 창·가게별cnt를 내는 /root/flink/windows/session.sql 의 출력을 /root/flink/windows/session.out 에 저장하세요.- 10분 TUMBLE 창마다 가게별 매출 상위 2곳(
window_start, window_end, shop, revenue, rownum)을 내는 /root/flink/windows/top-shops.sql 의 출력을 /root/flink/windows/top-shops.out 에 저장하세요. - 30분 TUMBLE 창마다 금액이 큰 주문 3건(
order_id, shop, amount, window_start, window_end, rownum)을 집계 없이 내는 /root/flink/windows/top-orders.sql 의 출력을 /root/flink/windows/top-orders.out 에 저장하세요. - /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 · Window Aggregation · Window Top-N · Time Attributes
창 TVF 가 붙이는 열 세 개
flink-up 으로 클러스터를 띄우고, /root/flink/windows/assign.sql 에 원본 표 orders(참고의 열과 워터마크)와 SELECT order_id, ts, window_start, window_end, window_time FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE) WHERE order_id <= 20 을 스트리밍으로 돌려 출력을 /root/flink/windows/assign.out 에 저장하세요.
창 TVF 는 원래 열을 그대로 두고 창 열 세 개를 붙여 돌려줍니다. 창은 [시작, 끝) 반열린 구간이라 경계 시각에 딱 찍힌 주문은 그 시각에 시작하는 창으로 갑니다. window_time 과 window_end 의 차이를 보세요.
가게별 10분 집계
/root/flink/windows/tumble.sql 에 10분 TUMBLE 창과 shop 으로 묶어 window_start, window_end, shop, COUNT(*) AS cnt, SUM(amount) AS revenue 를 내는 창 집계를 쓰고, 출력을 /root/flink/windows/tumble.out 에 저장하세요.
창 집계는 GROUP BY 에 window_start 와 window_end 를 넣습니다. 창이 겹치지 않으므로 cnt 를 모두 더하면 원래 주문 수와 같습니다. 결과는 창이 닫힐 때 +I 로 한 번만 나옵니다.
HOP — 한 주문이 두 창에 들어간다
/root/flink/windows/hop.sql 에 HOP(TABLE orders, DESCRIPTOR(ts), INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) 창마다 window_start, window_end, COUNT(*) AS cnt 를 내는 집계를 쓰고, 출력을 /root/flink/windows/hop.out 에 저장하세요.
HOP 의 셋째 인자가 slide(창이 시작하는 간격), 넷째가 size(창 길이)입니다. 5분마다 시작하는 10분 창이면 창이 절반씩 겹쳐 한 주문이 창 두 개에 들어갑니다. cnt 합을 원래 주문 수와 비교해 보세요. 맨 앞 창은 첫 주문보다 5분 이른 시각에 시작할 수 있습니다.
CUMULATE — 시작이 고정된 누적 창
/root/flink/windows/cumulate.sql 에 CUMULATE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE, INTERVAL '1' HOUR) 창마다 window_start, window_end, SUM(amount) AS revenue 를 내는 집계를 쓰고, 출력을 /root/flink/windows/cumulate.out 에 저장하세요.
CUMULATE 는 size(1시간)로 TUMBLE 한 뒤 그 안을 step(10분)마다 끝이 늘어나는 창으로 나눈 것입니다. 시작이 같은 창의 revenue 는 끝이 늦을수록 커지거나 같아야 합니다. 한 시간마다 창이 몇 개 나오는지 세어 보세요.
SESSION — 가게마다 길이가 다른 창
/root/flink/windows/session.sql 에 SESSION(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), INTERVAL '5' MINUTE) 창·가게별 window_start, window_end, shop, COUNT(*) AS cnt 를 내는 집계를 쓰고, 출력을 /root/flink/windows/session.out 에 저장하세요.
세션은 같은 가게의 이웃한 주문 간격이 5분 이하면 이어지고, 넘으면 새 세션이 시작됩니다. 세션의 시작은 첫 주문 시각, 끝은 마지막 주문 + 5분입니다. PARTITION BY 를 빼면 가게를 섞어 세션을 나눕니다.
창 집계 위의 창 Top-N
/root/flink/windows/top-shops.sql 에 10분 TUMBLE 창·가게별 매출(SUM(amount) AS revenue)을 구한 뒤, ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY revenue DESC) AS rownum 으로 창마다 상위 2곳만 남겨 window_start, window_end, shop, revenue, rownum 을 내세요. 출력은 /root/flink/windows/top-shops.out 에 저장합니다.
창 집계를 서브쿼리로 두고 그 위에서 ROW_NUMBER 를 매긴 뒤 바깥에서 rownum <= 2 로 거릅니다. PARTITION BY 에 창 열 두 개가 있어야 창 Top-N 이 되어 창이 닫힐 때 한 번만 결과를 냅니다. 가게가 한 곳뿐인 창은 한 줄만 나옵니다.
창 TVF 바로 위의 창 Top-N
/root/flink/windows/top-orders.sql 에 집계 없이 30분 TUMBLE 창 TVF 바로 위에서 ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY amount DESC) 로 순위를 매겨 창마다 금액이 큰 주문 3건의 order_id, shop, amount, window_start, window_end, rownum 을 내세요. 출력은 /root/flink/windows/top-orders.out 에 저장합니다.
창 Top-N 은 창 집계 없이도 창 TVF 의 결과 위에 바로 올릴 수 있습니다. 이때 순위의 대상은 집계 행이 아니라 주문 행 자체입니다. GROUP BY 는 쓰지 않습니다. 금액에 동점이 없게 만든 자료라 순위가 흔들리지 않습니다.
보고서 — 창마다 몇 번 세어지나
/root/flink/windows/report.json 에 orders(tumble.out 의 cnt 합), hop_assignments(hop.out 의 cnt 합), cumulate_windows(cumulate.out 의 줄 수), sessions(session.out 의 줄 수), max_session_orders(session.out 의 cnt 최댓값)를 정수로 적으세요.
결과 행은 '| +I |' 로 시작합니다. awk -F'|' 로 나누면 1번 칸은 빈 칸, 2번 칸이 op 이고 그 뒤로 SELECT 의 열 순서대로 이어집니다. hop_assignments 를 orders 와 비교해 보세요 — size/slide 배입니다.