LabHub
배우기 러닝패스 코스

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

重複を取り除き、順位を付ける

LabHub 에서 이어서 보기

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

목표

같은 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. 마지막 행 유지 위에서 statusorders(주문 수)·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.jsonn_rows·n_orders·update_pairs·shipped_now·topn_log_rows·norank_log_rows·top_books 를 적으세요.

참고

원천 둘을 정의하고 행 수를 센다

flink-up/root/flink/dedup/ddl.sql 에 events · sales 를 정의하세요(events 의 tsWATERMARK FOR ts AS ts). 배치 모드로 n_rows·n_orders 를 내는 /root/flink/dedup/count.sqlsql-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 statusorders(주문 수)·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.jsonn_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단계의 두 숫자로도 맞는지 확인해 보세요.