Apache Flink — Running Streams on a Real Engine
Removing duplicates and ranking
한국어 원문으로 표시합니다.
목표
같은 ROW_NUMBER 패턴으로 첫 행 유지·마지막 행 유지·Top-N 을 만들어 보고, 각각이 내는 변경 로그의 모양과 양을 출력 개수와 실행 계획으로 확인한다.
왜 중요한가
재전송된 사건을 걷어 낼 때와 지금 상태를 남길 때는 SQL 로는 ASC 와 DESC 한 단어 차이지만, 앞의 것은 append-only 결과를, 뒤의 것은 철회·갱신 로그를 낸다. 이 차이가 어느 싱크를 쓸 수 있는지, 하류 집계가 무엇을 받는지를 정한다. Top-N 은 순번을 내보내느냐에 따라 같은 결과에 로그 양이 크게 달라진다. 이 실습의 채점기는 클러스터에 묻지 않고, 여러분이 저장한 sql-client 출력을 원본 CSV 에서 직접 계산한 값과 대조한다.
단계
flink-up으로 클러스터를 띄우고 /root/flink/dedup/ddl.sql 에 events · sales 를 정의하세요(events 의ts에 워터마크). 배치로n_rows(전체 행)·n_orders(서로 다른 주문 수)를 내는 /root/flink/dedup/count.sql 의 출력을 /root/flink/dedup/count.out 에 저장하세요.- 주문마다 가장 이른 이벤트를 남기는 첫 행 유지 쿼리 /root/flink/dedup/first.sql(
order_id·status·amount)의 출력을 /root/flink/dedup/first.out 에 저장하세요. - 주문마다 가장 늦은 이벤트를 남기는 마지막 행 유지 쿼리 /root/flink/dedup/last.sql 의 출력을 /root/flink/dedup/last.out 에 저장하세요.
- 2·3단계 쿼리를
EXPLAIN CHANGELOG_MODE하는 /root/flink/dedup/explain.sql 의 출력을 /root/flink/dedup/explain.out 에 저장하세요. - 마지막 행 유지 위에서
status별orders(주문 수)·amount(합)를 세는 /root/flink/dedup/board.sql 의 출력을 /root/flink/dedup/board.out 에 저장하세요. - sales 로 범주별 누적
qty합계 상위 3개를 순번과 함께 내는 /root/flink/dedup/topn.sql(category·product·total·rn)의 출력을 /root/flink/dedup/topn.out 에 저장하세요. - 같은 Top-3 에서 바깥 SELECT 의
rn만 뺀 /root/flink/dedup/topn-norank.sql 의 출력을 /root/flink/dedup/topn-norank.out 에 저장하세요. - /root/flink/dedup/report.json 에
n_rows·n_orders·update_pairs·shipped_now·topn_log_rows·norank_log_rows·top_books를 적으세요.
참고
- 원본(머리글 없는 CSV, 시각은 초 단위, 행 순서가 곧 도착 순서이고 ts 오름차순):
dedup_events.csv=order_id, status, amount, ts— 한 주문에 1–4건, 일부는 같은 행이 바로 뒤에 한 번 더 옵니다(재전송).dedup_sales.csv=sale_id, category, product, qty, ts. 모두/opt/lab/fixtures/data/에 있습니다. - 실행:
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1. 스트리밍 결과는 맨 앞에op열(+I · -U · +U · -D)이 붙습니다. - 중복 제거는 문서의 모양을 정확히 따라야 합니다: 안쪽 SELECT 에
ROW_NUMBER() OVER (PARTITION BY 키 ORDER BY 시간속성 ASC|DESC) AS rn, 바깥에WHERE rn = 1. - 흔한 실수: 처리 시간(
PROCTIME())으로 정렬하면 결과가 도착 시각에 따라 달라질 수 있습니다. 이 실습은 이벤트 시간ts로 정렬합니다. - 공식 문서: Deduplication · Top-N · EXPLAIN · Dynamic Tables
원천 둘을 정의하고 행 수를 센다
flink-up 뒤 /root/flink/dedup/ddl.sql 에 events · sales 를 정의하세요(events 의 ts 에 WATERMARK FOR ts AS ts). 배치 모드로 n_rows·n_orders 를 내는 /root/flink/dedup/count.sql 을 sql-client.sh -i ddl.sql -f count.sql 로 돌려 /root/flink/dedup/count.out 에 저장하세요.
주문 수는 COUNT(DISTINCT order_id) 입니다. 두 값의 차이가 곧 같은 주문에 대해 두 번째 이후로 온 이벤트 수입니다 — 3단계에서 이 숫자를 다시 만납니다.
첫 행 유지 — 한 번 내면 끝
/root/flink/dedup/first.sql 에 스트리밍 모드로 주문(order_id)마다 ts 가 가장 이른 이벤트 하나를 남기는 중복 제거를 쓰고 order_id·status·amount 를 내세요. 출력은 /root/flink/dedup/first.out 에 저장합니다.
안쪽 SELECT 에서 ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY ts ASC) 로 순번을 매기고, 바깥에서 순번이 1 인 것만 고릅니다. 출력의 op 가 전부 +I 인지 보세요 — 한 번 낸 첫 행은 바뀔 일이 없습니다.
마지막 행 유지 — 철회하고 고친다
/root/flink/dedup/last.sql 에 주문마다 ts 가 가장 늦은 이벤트를 남기는 중복 제거를 쓰고(order_id·status·amount), 출력을 /root/flink/dedup/last.out 에 저장하세요. 로그 개수(+I · -U · +U)가 1단계의 두 숫자와 어떤 관계인지 보세요.
정렬 방향만 바꾸면 됩니다. 같은 주문의 이벤트가 새로 올 때마다 앞서 낸 행을 -U 로 거두고 새 행을 +U 로 냅니다. 재전송으로 완전히 같은 행이 와도 쌍이 나옵니다. 처리 시간으로 정렬하면 이 개수가 흔들릴 수 있으니 ts 를 쓰세요.
계획에서 두 중복 제거를 가른다
/root/flink/dedup/explain.sql 에 2단계와 3단계 쿼리 앞에 각각 EXPLAIN CHANGELOG_MODE 를 붙인 두 문장을 쓰고, 출력을 /root/flink/dedup/explain.out 에 저장하세요.
실행 계획(Optimized Execution Plan)에 Deduplicate(keep=[...]) 노드가, 물리 계획(Optimized Physical Plan)의 줄마다 changelogMode=[...] 가 붙습니다. 두 쿼리에서 keep 과 outputInsertOnly, changelogMode 가 어떻게 다른지 비교하세요.
지금 상태로 센다 — 중복 제거 위의 집계
/root/flink/dedup/board.sql 에 3단계의 마지막 행 유지를 안쪽에 두고 바깥에서 GROUP BY status 로 orders(주문 수)·amount(amount 합)를 내는 쿼리를 쓰고, 출력을 /root/flink/dedup/board.out 에 저장하세요.
주문 하나가 created 에서 paid 로 바뀌면 중복 제거가 -U created · +U paid 를 내고, 집계는 created 칸에서 빼고 paid 칸에 더합니다. 최종 orders 의 합은 주문 수와 같아야 합니다. 중복 제거 없이 events 를 바로 세면 이 합이 행 수만큼 커집니다.
Top-3 — 순위가 바뀔 때마다 고친다
/root/flink/dedup/topn.sql 에 sales 를 (category, product)별 SUM(qty) AS total 로 모은 뒤, 범주마다 total 내림차순 상위 3개를 순번 rn 과 함께 내는 쿼리를 쓰세요(category·product·total·rn). 출력은 /root/flink/dedup/topn.out 에 저장합니다.
GROUP BY 합계를 서브쿼리로 두고, 그 위에 ROW_NUMBER() OVER (PARTITION BY category ORDER BY total DESC) 를 매긴 뒤 바깥에서 rn <= 3 으로 고릅니다. 합계가 바뀔 때마다 순위가 움직이므로 출력에 -U · -D 가 섞여 나옵니다. 최종 상태만 보면 범주마다 세 줄입니다.
순번을 빼면 로그가 준다
6단계 쿼리에서 바깥 SELECT 의 rn 만 뺀 /root/flink/dedup/topn-norank.sql(category·product·total)을 돌려 출력을 /root/flink/dedup/topn-norank.out 에 저장하세요. 최종 상위 3개는 같고 로그 행 수는 6단계보다 적어야 합니다.
순번 칸은 결과의 유일 키 일부라서, 내보내면 한 상품의 순위가 오를 때 그 아래 순위의 행이 모두 다시 나갑니다. 순번을 빼면 바뀐 상품의 행만 보내면 됩니다. 두 파일의 로그 행 수를 grep -c 로 세어 비교해 보세요.
보고서 — 로그의 양을 숫자로
/root/flink/dedup/report.json 에 n_rows·n_orders(1단계), update_pairs(last.out 의 -U 개수), shipped_now(board.out 최종 상태의 shipped 주문 수), topn_log_rows·norank_log_rows(두 Top-3 출력의 로그 행 수), top_books(books 범주 1위 상품)를 적으세요.
모두 앞에서 저장한 출력에서 옮깁니다. 로그 행은 '| +I |' 처럼 op 로 시작하는 줄입니다. 변경 로그가 있는 표의 최종 값은 그 키의 마지막 +I · +U 입니다. update_pairs 는 1단계의 두 숫자로도 맞는지 확인해 보세요.