Apache Flink — ストリームを本物のエンジンで動かす
実行計画から状態と TTL を読む
한국어 원문으로 표시합니다.
목표
연산자마다 어떤 상태를 쥐고 TTL 이 얼마인지 COMPILE PLAN JSON 으로 읽고, 잡 전체 TTL 과 연산자별 힌트가 어디에 박히는지, 창 집계와 상태 백엔드가 무엇을 바꾸는지 확인한다.
왜 중요한가
끝없는 GROUP BY 와 일반 조인은 기본으로 상태를 영원히 쥐고, TTL 로 줄이면 결과가 틀릴 수 있다. TTL 은 "마지막 갱신 뒤 이만큼 지나면 언젠가 지운다" 라서 결과로는 판정할 수 없지만, 계획은 결정적으로 적혀 나온다. 운영에서 상태 문제를 다루는 첫 도구가 바로 이 계획이다. 이 실습의 채점기는 클러스터에 묻지 않고, 여러분이 저장한 계획 JSON · REST 응답 · sql-client 출력을 읽고 결과는 원본 CSV 로 직접 계산해 대조한다.
단계
flink-up뒤 /root/flink/state/ddl.sql 에 orders · payments · users 와 blackhole 싱크kv_sink(k STRING, v BIGINT)·win_sink(window_start, window_end, v BIGINT)를 정의하고, 거르기만 하는 INSERT 의 계획을 /root/flink/state/calc.sql 로 /root/flink/state/calc-plan.json 에 뽑으세요.- 사용자별
SUM(amount)를 내는 끝없는 GROUP BY 의 계획을 /root/flink/state/agg.sql 로 /root/flink/state/agg-plan.json 에 뽑으세요(TTL 설정 없이). SET 'table.exec.state.ttl' = '30 min';을 넣은 같은 계획을 /root/flink/state/agg-ttl.sql 로 /root/flink/state/agg-ttl-plan.json 에 뽑으세요.- 잡 기본값 30 min 을 둔 채 orders o · payments p 조인에
STATE_TTL('o' = '1d', 'p' = '2h')를 준 계획을 /root/flink/state/join-hint.sql 로 /root/flink/state/join-hint-plan.json 에 뽑으세요. - 잡 기본값 30 min 에 orders o ⋈ payments p ⋈ users u 조인과
STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d')를 준 계획을 /root/flink/state/cascade.sql 로 /root/flink/state/cascade-plan.json 에 뽑으세요. - 잡 기본값 30 min 을 둔 채 10분 TUMBLE 창 집계의 계획을 /root/flink/state/window-plan.json 에 뽑고, 같은 창의
window_start·window_end·orders·amount를 스트리밍으로 돌려 출력을 /root/flink/state/window.out 에 저장하세요(둘 다 /root/flink/state/window.sql 한 파일로). SET 'state.backend.type' = 'rocksdb'· 잡 이름flk-state-rocks로 사용자별orders·total을 내는 /root/flink/state/rocks.sql 의 출력을 /root/flink/state/rocks.out 에, 그 잡의/jobs/<jid>를 /root/flink/state/rocks-job.json,/jobs/<jid>/checkpoints/config를 /root/flink/state/rocks-ckpt.json 에 저장하세요.- /root/flink/state/report.json 에
default_ttl·join_left_ttl·join_right_ttl·cascade_second_left_ttl·window_state_entries·state_backend·rocks_groups를 적으세요.
참고
- 원본(머리글 없는 CSV, 시각은 초 단위, ts 오름차순):
state_orders.csv=order_id, user_id, amount, ts·state_payments.csv=pay_id, order_id, pay_method, ts·state_users.csv=user_id, region. 모두/opt/lab/fixtures/data/에 있습니다. orders · payments 의ts에 워터마크를 거세요. - 계획 뽑기:
COMPILE PLAN 'file:///root/flink/state/이름.json' FOR INSERT INTO 싱크 SELECT ...;— 잡을 돌리지 않습니다. 같은 경로에 파일이 있으면 덮어쓰지 않고 오류를 내므로 다시 뽑을 때는 먼저 지웁니다. - 계획 훑기:
jq -c '.nodes[] | {id, type, state}' 파일.json— 노드 id 순서가 아래(원천)에서 위(싱크)로의 순서입니다. - 흔한 실수: 문서대로 표에 별칭을 붙였다면 힌트 열쇠도 별칭이어야 합니다.
method는 SQL 예약어라 열 이름으로 쓰면 CREATE 가 실패합니다(pay_method). - 공식 문서: Hints — STATE_TTL · Configuration — table.exec.state.ttl · State Backends · Group Aggregation · Determinism
거르기만 하는 쿼리에는 상태가 없다
flink-up 뒤 /root/flink/state/ddl.sql 에 orders · payments · users 와 blackhole 싱크 kv_sink·win_sink 를 정의하세요. /root/flink/state/calc.sql 에 amount > 1000 인 주문의 order_id 와 amount * 2 를 kv_sink 에 넣는 INSERT 를 COMPILE PLAN 'file:///root/flink/state/calc-plan.json' 으로 뽑는 문장을 쓰고 sql-client.sh -i ddl.sql -f calc.sql 로 돌리세요.
COMPILE PLAN 은 INSERT 문을 받으므로 받을 싱크가 필요합니다. blackhole 커넥터는 받은 것을 버립니다. 만든 JSON 을 jq 로 열어 nodes 의 type 과 state 를 보세요 — 행 하나를 보고 바로 내보내는 연산자에는 기억할 것이 없습니다.
끝없는 GROUP BY — 기본은 영원히
/root/flink/state/agg.sql 에 사용자별 SUM(amount) 를 kv_sink 에 넣는 INSERT 의 계획을 /root/flink/state/agg-plan.json 에 뽑으세요. TTL 은 설정하지 않습니다.
집계 노드의 state 목록에 무엇이 있고 ttl 이 얼마인지 보세요. 0 ms 는 지우지 않는다는 뜻입니다. SUM 의 결과 타입이 싱크 열(BIGINT)과 다르면 CAST 로 맞춥니다.
잡 전체 TTL 을 건다
/root/flink/state/agg-ttl.sql 맨 앞에 SET 'table.exec.state.ttl' = '30 min'; 을 두고 2단계와 같은 GROUP BY 의 계획을 /root/flink/state/agg-ttl-plan.json 에 뽑으세요.
SET 은 그 뒤 문장부터 적용됩니다. 계획의 groupAggregateState ttl 이 어떻게 바뀌는지 보세요. 이 값은 '마지막 갱신 뒤 최소 이만큼은 보존' 이라는 뜻이라 결과로는 언제 지워졌는지 알 수 없습니다.
조인 양쪽에 다른 TTL — STATE_TTL 힌트
/root/flink/state/join-hint.sql 에 SET 'table.exec.state.ttl' = '30 min'; 을 둔 채, orders o 와 payments p 를 order_id 로 일반 조인해 kv_sink 에 넣는 INSERT 에 /*+ STATE_TTL('o' = '1d', 'p' = '2h') */ 를 주고 계획을 /root/flink/state/join-hint-plan.json 에 뽑으세요.
힌트는 SELECT 바로 뒤에 씁니다. 표에 별칭을 붙였다면 힌트의 열쇠도 별칭이어야 합니다. 계획의 조인 노드에서 leftState · rightState 의 ttl 이 잡 기본값과 다르게 나오는지 보세요.
이어진 조인 — 힌트가 닿지 않는 자리
/root/flink/state/cascade.sql 에 잡 기본값 30 min 을 두고, orders o ⋈ payments p(order_id) ⋈ users u(user_id) 를 kv_sink 에 넣는 INSERT 에 /*+ STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d') */ 를 준 계획을 /root/flink/state/cascade-plan.json 에 뽑으세요.
조인 노드가 둘 나옵니다. id 가 작은 쪽이 먼저(아래) 조인입니다. 힌트 세 개가 네 자리(첫 조인 왼쪽·오른쪽, 둘째 조인 왼쪽·오른쪽) 중 어디에 붙고, 남은 한 자리는 무엇을 받는지 보세요.
창 집계 — TTL 없이 스스로 비운다
/root/flink/state/window.sql 한 파일에 SET 'table.exec.state.ttl' = '30 min'; 을 두고 (1) 10분 TUMBLE 창의 SUM(amount) 를 win_sink 에 넣는 계획을 /root/flink/state/window-plan.json 에 뽑고, (2) 같은 창의 window_start·window_end·orders(COUNT)·amount(SUM)를 스트리밍 SELECT 로 돌리세요. 출력은 /root/flink/state/window.out 에 저장합니다.
창 TVF 는 FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)) 이고 GROUP BY window_start, window_end 로 묶습니다. TTL 을 걸었는데도 창 집계 노드에 state 항목이 있는지 보세요. 결과의 op 가 모두 +I 인 것도 같은 이유입니다 — 창은 워터마크가 끝을 지날 때 한 번 내고 비웁니다.
RocksDB 백엔드로 같은 집계를 돌린다
/root/flink/state/rocks.sql 에 SET 'state.backend.type' = 'rocksdb'; · SET 'pipeline.name' = 'flk-state-rocks'; 를 두고 사용자별 orders(COUNT)·total(SUM(amount))을 스트리밍으로 내세요. 출력은 /root/flink/state/rocks.out, 그 잡의 /jobs/<jid> 는 /root/flink/state/rocks-job.json, /jobs/<jid>/checkpoints/config 는 /root/flink/state/rocks-ckpt.json 에 저장합니다.
잡 id 는 /jobs/overview 에서 이름으로 찾습니다. checkpoints/config 응답의 state_backend 칸에 이 잡이 쓴 백엔드 클래스 이름이 적힙니다. 기본 백엔드로 돌린 잡이면 다른 이름이 나옵니다. 결과는 백엔드와 무관하게 같아야 합니다.
보고서 — 상태 지도를 숫자로
/root/flink/state/report.json 에 default_ttl(agg-plan 의 groupAggregateState ttl), join_left_ttl·join_right_ttl(join-hint-plan), cascade_second_left_ttl(cascade-plan 둘째 조인의 leftState ttl), window_state_entries(window-plan 의 state 항목 수, 정수), state_backend(rocks-ckpt 의 값), rocks_groups(rocks.out 최종 결과의 행 수 = 상태에 남은 키 수, 정수)를 적으세요.
ttl 값은 계획에 적힌 문자열(예: "2 h")을 그대로 옮기면 됩니다. jq 로 type 이 stream-exec-join 으로 시작하는 노드를 id 순으로 정렬하면 두 조인을 가를 수 있습니다. state 가 없는 노드는 그 칸이 null 입니다.