LabHub
배우기 러닝패스 코스

Apache Flink — ストリームを本物のエンジンで動かす

4種類のウィンドウ TVF で同じ注文を切り分ける

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

같은 주문 스트림에 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_timeorder_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.jsonorders·hop_assignments·cumulate_windows·sessions·max_session_orders 를 적으세요.

참고

창 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.sqlHOP(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.sqlCUMULATE(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.sqlSESSION(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.jsonorders(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 배입니다.