Apache Flink — 스트림을 엔진으로 돌린다 · 스트림 조인 · 실습
세 가지 조인으로 주문을 붙인다
목표
같은 주문 흐름에 사용자·배송·환율을 일반 조인·구간 조인·이벤트 시간 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_v 를 FOR 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.json 에 matched_orders·interval_rows·unshipped_orders·temporal_rows·orders_without_rate·usd_gap_krw 를 적으세요.
참고
- 원본(머리글 없는 CSV, 시각은 초 단위, 파일마다 자기 시간 순서로 정렬):
- 정의를 한 번만 쓰려면
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1—-i파일의 CREATE 문이 먼저 실행됩니다. - 스트리밍 결과는 맨 앞에
op열(+I · -U · +U · -D)이 붙습니다. 잡은 파일 끝에서 끝나고, 그때 워터마크가 끝까지 전진해 남은 구간·버전이 모두 처리됩니다. COMPILE PLAN은INSERT INTO문을 받습니다.'connector' = 'blackhole'싱크를 하나 만들어 쓰세요. 같은 경로에 파일이 있으면 덮어쓰지 않고 오류를 내므로 다시 뽑을 때는 먼저 지웁니다.- 흔한 실수: 바깥 구간 조인의 시간 조건을
WHERE에 두면 null 행이 걸러집니다. temporal join 의 오른쪽은 기본 키가 있어야 하는데, append-only 원천은 중복 제거 뷰로 만들어야 합니다. - 공식 문서: [Joins](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/joins/) · [Versioned Tables](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/sql-table-concepts/versioned_tables/) · [Deduplication](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/deduplication/) · [SQL Client](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/interfaces/sql-client/)
joins_users.csv = user_id, tier, signup_time · joins_orders.csv = order_id, user_id, currency, amount, order_time · joins_shipments.csv = ship_id, order_id, ship_time · joins_rates.csv = currency, rate, update_time (rate 는 원화, DECIMAL(10, 4) 로 읽으세요). 모두 /opt/lab/fixtures/data/ 에 있습니다.
8단계
- 네 원천을 정의하고 행 수를 센다
- 일반 조인 — 시간을 보지 않고 양쪽을 기억한다
- 구간 조인 — 2시간 안의 배송만 붙인다
- 바깥 구간 조인 — 제때 짝이 없던 주문
- temporal join — 주문 시각의 환율
- 같은 뷰를 일반 조인으로 — 지난 주문이 다시 계산된다
- 세 조인의 실행 계획에서 상태를 읽는다
- 보고서 — 조인마다 무엇이 붙고 빠졌나