Apache Flink — 스트림을 엔진으로 돌린다 · 상태와 TTL · 실습
실행 계획에서 상태와 TTL 을 읽는다
목표
연산자마다 어떤 상태를 쥐고 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.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](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/hints/) · [Configuration — table.exec.state.ttl](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/dev/table/config/) · [State Backends](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/state/state_backends/) · [Group Aggregation](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/group-agg/) · [Determinism](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/sql-table-concepts/determinism/)
8단계
- 거르기만 하는 쿼리에는 상태가 없다
- 끝없는 GROUP BY — 기본은 영원히
- 잡 전체 TTL 을 건다
- 조인 양쪽에 다른 TTL — STATE_TTL 힌트
- 이어진 조인 — 힌트가 닿지 않는 자리
- 창 집계 — TTL 없이 스스로 비운다
- RocksDB 백엔드로 같은 집계를 돌린다
- 보고서 — 상태 지도를 숫자로