LabHub
배우기 러닝패스 코스

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

実行計画から状態と TTL を読む

LabHub 에서 이어서 보기

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

목표

연산자마다 어떤 상태를 쥐고 TTL 이 얼마인지 COMPILE PLAN JSON 으로 읽고, 잡 전체 TTL 과 연산자별 힌트가 어디에 박히는지, 창 집계와 상태 백엔드가 무엇을 바꾸는지 확인한다.

왜 중요한가

끝없는 GROUP BY 와 일반 조인은 기본으로 상태를 영원히 쥐고, TTL 로 줄이면 결과가 틀릴 수 있다. TTL 은 "마지막 갱신 뒤 이만큼 지나면 언젠가 지운다" 라서 결과로는 판정할 수 없지만, 계획은 결정적으로 적혀 나온다. 운영에서 상태 문제를 다루는 첫 도구가 바로 이 계획이다. 이 실습의 채점기는 클러스터에 묻지 않고, 여러분이 저장한 계획 JSON · REST 응답 · sql-client 출력을 읽고 결과는 원본 CSV 로 직접 계산해 대조한다.

단계

  1. 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 에 뽑으세요.
  2. 사용자별 SUM(amount) 를 내는 끝없는 GROUP BY 의 계획을 /root/flink/state/agg.sql/root/flink/state/agg-plan.json 에 뽑으세요(TTL 설정 없이).
  3. SET 'table.exec.state.ttl' = '30 min'; 을 넣은 같은 계획을 /root/flink/state/agg-ttl.sql/root/flink/state/agg-ttl-plan.json 에 뽑으세요.
  4. 잡 기본값 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 에 뽑으세요.
  5. 잡 기본값 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 에 뽑으세요.
  6. 잡 기본값 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 한 파일로).
  7. 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 에 저장하세요.
  8. /root/flink/state/report.jsondefault_ttl·join_left_ttl·join_right_ttl·cascade_second_left_ttl·window_state_entries·state_backend·rocks_groups 를 적으세요.

참고

거르기만 하는 쿼리에는 상태가 없다

flink-up/root/flink/state/ddl.sql 에 orders · payments · users 와 blackhole 싱크 kv_sink·win_sink 를 정의하세요. /root/flink/state/calc.sqlamount > 1000 인 주문의 order_idamount * 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.sqlSET 'table.exec.state.ttl' = '30 min'; 을 둔 채, orders o 와 payments porder_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.sqlSET '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.jsondefault_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 입니다.