LabHub
배우기 러닝패스 코스

Apache Flink — Running Streams on a Real Engine

Joining orders three ways

LabHub 에서 이어서 보기

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

목표

같은 주문 흐름에 사용자·배송·환율을 일반 조인·구간 조인·이벤트 시간 temporal join 으로 붙여 보고, 각 조인이 무엇을 기억하고 결과를 고치는지 출력과 실행 계획으로 확인한다.

왜 중요한가

스트림 조인은 짝이 언제 올지 모르기 때문에 반대편 행을 상태에 쌓아 둔다. 시간 조건이 없으면 영원히, 시간 범위를 주면 워터마크가 지날 때까지, 버전 테이블이면 필요한 버전만 쥔다. 조인 종류를 잘못 고르면 상태가 끝없이 자라거나 지난 결과가 조용히 바뀐다. 이 실습의 채점기는 클러스터에 묻지 않고, 여러분이 저장한 sql-client 출력과 계획 JSON 을 원본 CSV 에서 직접 계산한 값과 대조한다.

단계

  1. flink-up 으로 클러스터를 띄우고 /root/flink/joins/ddl.sql 에 네 원천(users · orders · shipments · rates)을 정의하세요. orders·shipments·rates 의 시간 열에는 워터마크를 겁니다. 배치로 네 표의 행 수를 세는 /root/flink/joins/count.sql 의 출력을 /root/flink/joins/count.out 에 저장하세요(열 이름 tbl·n).
  2. /root/flink/joins/regular.sql 에 orders 와 users 를 user_id 로 붙여 등급별 orders(건수)·amount(합)를 내는 스트리밍 쿼리를 쓰고, 출력을 /root/flink/joins/regular.out 에 저장하세요.
  3. /root/flink/joins/interval.sql 에 주문 시각부터 2시간 안(양끝 포함)에 나간 배송을 붙이는 구간 조인을 쓰고 order_id·ship_id·delay_s(초)를 /root/flink/joins/interval.out 에 저장하세요.
  4. /root/flink/joins/unshipped.sql 에 2시간 안에 배송되지 않은 주문의 order_id·order_time 을 내는 바깥 구간 조인을 쓰고 출력을 /root/flink/joins/unshipped.out 에 저장하세요.
  5. rates 로 버전 뷰 rates_v 를 만들고, 주문 시각의 환율을 붙이는 temporal join 으로 order_id·currency·rate·amount_krw(amount × rate)를 내는 /root/flink/joins/temporal.sql 의 출력을 /root/flink/joins/temporal.out 에 저장하세요.
  6. 같은 rates_vFOR SYSTEM_TIME AS OF 없이 붙여 통화별 orders·total_krw 를 내는 /root/flink/joins/latest.sql 의 출력을 /root/flink/joins/latest.out 에 저장하세요.
  7. 세 조인(2·3·5단계)의 실행 계획을 COMPILE PLAN 으로 /root/flink/joins/regular-plan.json·/root/flink/joins/interval-plan.json·/root/flink/joins/temporal-plan.json 에 뽑으세요.
  8. /root/flink/joins/report.jsonmatched_orders·interval_rows·unshipped_orders·temporal_rows·orders_without_rate·usd_gap_krw 를 적으세요.

참고

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

flink-up/root/flink/joins/ddl.sql 에 users · orders · shipments · rates 를 정의하세요(order_time · ship_time · update_time 에 워터마크). 배치 모드로 네 표의 행 수를 tbl·n 열로 내는 /root/flink/joins/count.sqlsql-client.sh -i ddl.sql -f count.sql 로 돌려 /root/flink/joins/count.out 에 저장하세요.

워터마크는 WATERMARK FOR 시간열 AS 시간열 처럼 씁니다. 파일마다 자기 시간 순서로 정렬돼 있어 지연을 주지 않아도 늦은 행이 없습니다. 네 개의 SELECT 를 UNION ALL 로 이으면 표 하나로 나옵니다.

일반 조인 — 시간을 보지 않고 양쪽을 기억한다

/root/flink/joins/regular.sql 에 스트리밍 모드로 orders 와 users 를 user_id 로 안쪽 조인해 tierorders(건수)·amount(amount 합)를 내는 쿼리를 쓰고, 출력을 /root/flink/joins/regular.out 에 저장하세요.

일반 조인은 시간 조건 없이 등호만 씁니다. u41–u44 는 users 에 없는 사용자라 안쪽 조인에서 빠집니다. 오후에 가입한 사용자(signup_time 이 주문보다 늦음)의 주문도 붙는지 결과로 확인해 보세요 — 일반 조인은 도착 순서와 무관하게 과거·미래의 모든 짝을 찾습니다.

구간 조인 — 2시간 안의 배송만 붙인다

/root/flink/joins/interval.sql 에 orders 와 shipments 를 order_id 로 붙이되 ship_time 이 주문 시각부터 2시간 안(양끝 포함)인 것만 남기는 구간 조인을 쓰고, order_id·ship_id·delay_s(지연 초, TIMESTAMPDIFF(SECOND, ...))를 /root/flink/joins/interval.out 에 저장하세요.

구간 조인은 등호 하나와 양쪽 시간을 묶는 범위가 있어야 합니다. BETWEEN a AND b 는 양끝을 포함합니다. 지연이 정확히 0초·7200초인 배송이 있고, 7201초인 배송도 있습니다. 두 번에 나눠 배송된 주문은 두 줄이 됩니다.

바깥 구간 조인 — 제때 짝이 없던 주문

/root/flink/joins/unshipped.sql 에 orders LEFT JOIN shipments 로 2시간 안에 배송이 하나도 없던 주문의 order_id·order_time 을 내는 쿼리를 쓰고, 출력을 /root/flink/joins/unshipped.out 에 저장하세요.

시간 조건은 ON 절에 두고, 짝이 없어 null 로 채워진 행만 WHERE 로 고릅니다. 배송 기록이 아예 없는 주문뿐 아니라 2시간을 넘겨 나간 주문도 들어가야 합니다. null 행은 워터마크가 (주문 시각 + 2시간)을 지난 뒤에 한 번 나옵니다.

temporal join — 주문 시각의 환율

/root/flink/joins/temporal.sql 에서 rates 로 통화별 최신 행만 남기는 버전 뷰 rates_v 를 만들고, orders 를 FOR SYSTEM_TIME AS OF o.order_time 으로 붙여 order_id·currency·rate·amount_krw(amount × rate)를 내세요. 출력은 /root/flink/joins/temporal.out 에 저장합니다.

rates 는 append-only 라 기본 키를 걸 수 없습니다. ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) 가 1 인 행만 남기는 뷰를 만들면 currency 가 기본 키, update_time 이 이벤트 시간인 버전 뷰가 됩니다. 안쪽 조인이면 첫 환율보다 이른 주문은 빠집니다.

같은 뷰를 일반 조인으로 — 지난 주문이 다시 계산된다

/root/flink/joins/latest.sql 에서 5단계의 rates_vFOR SYSTEM_TIME AS OF 없이 orders 와 currency 로 붙여 통화별 orders(건수)·total_krw(amount × rate 의 합)를 내고, 출력을 /root/flink/joins/latest.out 에 저장하세요.

버전 뷰는 환율이 바뀔 때마다 갱신을 냅니다. 일반 조인은 그 갱신을 받아 지난 주문까지 다시 계산하므로 출력에 -U 가 보입니다. 최종 합계는 temporal join 의 합계와 다릅니다 — 어느 환율로 계산된 것인지 생각해 보세요.

세 조인의 실행 계획에서 상태를 읽는다

2·3·5단계 조인을 blackhole 싱크에 넣는 INSERT 로 바꿔 COMPILE PLAN 으로 /root/flink/joins/regular-plan.json·/root/flink/joins/interval-plan.json·/root/flink/joins/temporal-plan.json 을 만드세요. (일반 조인은 orders⋈users, 구간 조인은 orders⋈shipments, temporal join 은 orders⋈rates_v)

COMPILE PLAN 'file:///경로.json' FOR INSERT INTO 싱크 SELECT ... 모양입니다. 파일이 이미 있으면 오류가 나므로 다시 뽑을 때는 지우고 돌립니다. 만든 뒤 jq 로 nodes 의 type 과 state 를 훑어보세요 — 어느 조인 노드에 leftState · rightState 가 있고 어느 노드에는 없는지가 요점입니다.

보고서 — 조인마다 무엇이 붙고 빠졌나

/root/flink/joins/report.jsonmatched_orders(일반 조인으로 사용자와 붙은 주문 수), interval_rows(구간 조인 결과 행 수), unshipped_orders(4단계 행 수), temporal_rows(temporal join 결과 행 수), orders_without_rate(전체 주문 − temporal_rows), usd_gap_krw(latest.out 의 USD total_krw − temporal.out 의 USD amount_krw 합, 소수)를 적으세요.

모두 앞에서 저장한 출력에서 옮길 수 있습니다. 변경 로그가 있는 표는 마지막 +I/+U 가 최종 값입니다. grep 으로 '| +I |' 같은 행만 골라 awk -F'|' 로 칸을 자르면 셀 수 있습니다. 정수 칸은 정수로, usd_gap_krw 는 숫자로 적습니다.