Apache Flink — 스트림을 엔진으로 돌린다 · 중복 제거와 Top-N · 실습
중복을 걷어 내고 순위를 매긴다
목표
같은 ROW_NUMBER 패턴으로 첫 행 유지·마지막 행 유지·Top-N 을 만들어 보고, 각각이 내는 변경 로그의 모양과 양을 출력 개수와 실행 계획으로 확인한다.
왜 중요한가
재전송된 사건을 걷어 낼 때와 지금 상태를 남길 때는 SQL 로는 ASC 와 DESC 한 단어 차이지만, 앞의 것은 append-only 결과를, 뒤의 것은 철회·갱신 로그를 낸다. 이 차이가 어느 싱크를 쓸 수 있는지, 하류 집계가 무엇을 받는지를 정한다. Top-N 은 순번을 내보내느냐에 따라 같은 결과에 로그 양이 크게 달라진다. 이 실습의 채점기는 클러스터에 묻지 않고, 여러분이 저장한 sql-client 출력을 원본 CSV 에서 직접 계산한 값과 대조한다.
단계
1. 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 에 저장하세요.
2. 주문마다 가장 이른 이벤트를 남기는 첫 행 유지 쿼리 /root/flink/dedup/first.sql(order_id·status·amount)의 출력을 /root/flink/dedup/first.out 에 저장하세요.
3. 주문마다 가장 늦은 이벤트를 남기는 마지막 행 유지 쿼리 /root/flink/dedup/last.sql 의 출력을 /root/flink/dedup/last.out 에 저장하세요.
4. 2·3단계 쿼리를 EXPLAIN CHANGELOG_MODE 하는 /root/flink/dedup/explain.sql 의 출력을 /root/flink/dedup/explain.out 에 저장하세요.
5. 마지막 행 유지 위에서 status 별 orders(주문 수)·amount(합)를 세는 /root/flink/dedup/board.sql 의 출력을 /root/flink/dedup/board.out 에 저장하세요.
6. sales 로 범주별 누적 qty 합계 상위 3개를 순번과 함께 내는 /root/flink/dedup/topn.sql(category·product·total·rn)의 출력을 /root/flink/dedup/topn.out 에 저장하세요.
7. 같은 Top-3 에서 바깥 SELECT 의 rn 만 뺀 /root/flink/dedup/topn-norank.sql 의 출력을 /root/flink/dedup/topn-norank.out 에 저장하세요.
8. /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](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/deduplication/) · [Top-N](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/topn/) · [EXPLAIN](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/utility/explain/) · [Dynamic Tables](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/sql-table-concepts/dynamic_tables/)
8단계
- 원천 둘을 정의하고 행 수를 센다
- 첫 행 유지 — 한 번 내면 끝
- 마지막 행 유지 — 철회하고 고친다
- 계획에서 두 중복 제거를 가른다
- 지금 상태로 센다 — 중복 제거 위의 집계
- Top-3 — 순위가 바뀔 때마다 고친다
- 순번을 빼면 로그가 준다
- 보고서 — 로그의 양을 숫자로