Apache Flink — ストリームを本物のエンジンで動かす
クラスターを起動し、ジョブを一つ最後まで追う
한국어 원문으로 표시합니다.
목표
Flink 로컬 클러스터를 띄우고 REST API 로 구조를 확인한 뒤, 배치 잡 하나를 제출부터 종료까지 따라간다. 병렬도·슬롯·실패 기록을 직접 바꿔 보고 그 결과를 보고서로 정리한다.
왜 중요한가
운영에서 받는 질문 — 잡이 왜 멈췄나, 병렬도를 올렸는데 왜 그대로인가, 왜 재시작되지 않았나 — 은 SQL 이 아니라 엔진 구조에 대한 것이다. JobManager 는 잡마다 JobMaster 를 만들고 슬롯을 나눠 주며, 모든 기록을 REST 로 내놓는다. 이 실습의 채점기는 클러스터에 묻지 않는다. 여러분이 파일로 저장한 REST 응답과 sql-client 출력만 읽고, 집계 기대값은 원본 CSV 에서 직접 계산해 대조한다. 그래서 응답은 손으로 고치지 말고 curl 출력 그대로 저장한다.
단계
flink-up으로 클러스터를 띄우고/overview응답을 /root/flink/cluster/overview.json 에 저장하세요(TaskManager 1개 · 슬롯 2개가 보여야 합니다)./taskmanagers응답을 /root/flink/cluster/taskmanagers.json 에 저장하고, 그memoryConfiguration을 MiB 로 바꿔 반올림한 값 일곱 개를 /root/flink/cluster/memory.json 에 적으세요.- /root/flink/cluster/first.sql 에 배치 모드 · 잡 이름
flk-first-orders로/opt/lab/fixtures/data/cluster_orders.csv를 읽어 status 별orders(건수)·revenue(amount 합)를 내는 SQL 을 쓰고,sql-client.sh -f의 출력을 /root/flink/cluster/first.out 에 저장하세요. - 잡이 끝난 뒤
/jobs/overview응답을 /root/flink/cluster/jobs.json 에 저장하세요. - 같은 집계를 병렬도 2 · 잡 이름
flk-first-p2로 돌리는 /root/flink/cluster/p2.sql 을 만들어 출력을 /root/flink/cluster/p2.out 에, 그 잡의/jobs/<jid>응답을 /root/flink/cluster/p2-job.json 에 저장하세요. 정점 병렬도에 2 가 보여야 합니다. /opt/flink/conf/config.yaml의taskmanager.numberOfTaskSlots를 4 로 바꾸고 클러스터를 다시 띄운 뒤/overview를 /root/flink/cluster/overview-4.json 에 저장하세요.- 잡 이름
flk-bad-cast로status를INT로 CAST 하다 죽는 잡을 /root/flink/cluster/fail.sql 로 돌리고, 그 잡의/jobs/<jid>를 /root/flink/cluster/fail-job.json,/jobs/<jid>/exceptions를 /root/flink/cluster/fail-exceptions.json 에 저장하세요. - /root/flink/cluster/report.json 에
slots_total·first_job_id·p2_job_id·failed_job_id·restart_strategy·root_cause를 적으세요.
참고
- 원본 열:
order_id BIGINT, user_id STRING, status STRING, amount INT, order_time TIMESTAMP(3)(머리글 없는 CSV). - SQL 파일 실행:
sql-client.sh -f 파일.sql > 파일.out 2>&1. 오류가 난 문장에서 멈추고, 출력 파일에[ERROR]가 남습니다. - 결과는 표(tableau) 모드로 찍히게 설정돼 있습니다. 배치로 돌면
op열이 없고, 스트리밍으로 돌면 맨 앞에op열(+I · -U · +U)이 붙습니다. - 잡 id 찾기:
curl -s localhost:8081/jobs/overview | jq -r '.jobs[] | select(.name=="잡이름") | .jid' - 설정 파일에서 슬롯 수는
taskmanager:아래 들여쓰기 된numberOfTaskSlots:줄입니다. 파일만 고쳐서는 반영되지 않습니다(flink-down뒤flink-up). - 흔한 실수: 배치 잡은 적응형 스케줄러가 자료 크기로 병렬도를 다시 정합니다. 작은 파일이면 1 이 나옵니다.
- 이 파드에는 인터넷이 없습니다. 메모리 한도는 2Gi 라 클러스터와 SQL 클라이언트가 함께 1.5GB 안팎을 씁니다.
- 공식 문서: Flink Architecture · REST API · TaskManager 메모리 · Adaptive Batch · Task Failure Recovery
클러스터를 띄우고 개요를 받는다
flink-up 으로 클러스터를 띄운 뒤 curl -s localhost:8081/overview 응답을 /root/flink/cluster/overview.json 에 그대로 저장하세요.
flink-up 은 start-cluster.sh 를 부르고 TaskManager 가 JobManager 에 붙을 때까지 기다립니다. 붙기 전에 받은 응답에는 taskmanagers 가 0 으로 찍힙니다. 저장할 디렉터리부터 만드세요.
TaskManager 메모리 예산을 숫자로 읽는다
/taskmanagers 응답을 /root/flink/cluster/taskmanagers.json 에 저장하고, memoryConfiguration 의 값을 MiB 로 바꿔 반올림한 total_process_mb·total_flink_mb·jvm_metaspace_mb·jvm_overhead_mb·task_heap_mb·managed_mb·network_mb 를 /root/flink/cluster/memory.json 에 정수로 적으세요.
값은 바이트 단위입니다. 1048576 으로 나눠 반올림합니다. 프로세스 전체는 Flink 메모리 + 메타스페이스 + 오버헤드와 같아야 합니다 — 이 항등식이 맞지 않으면 어딘가 잘못 옮긴 것입니다. jq 의 round 를 쓰면 한 번에 만들 수 있습니다.
배치 잡 하나를 제출한다
/root/flink/cluster/first.sql 에 SET 'execution.runtime-mode' = 'batch'; · SET 'pipeline.name' = 'flk-first-orders'; · 원본 CSV 를 읽는 CREATE TABLE · status 별 orders(건수)와 revenue(amount 합)를 내는 SELECT 를 쓰고, sql-client.sh -f first.sql > first.out 2>&1 로 /root/flink/cluster/first.out 을 만드세요.
filesystem 커넥터에 path 는 file:///opt/lab/fixtures/data/cluster_orders.csv, format 은 csv 입니다. 열 이름은 지시문의 원본 열 그대로 쓰고, 결과 열에는 AS orders · AS revenue 로 별칭을 붙입니다. 출력에 op 열이 보이면 배치가 아니라 스트리밍으로 돈 것입니다.
끝난 잡을 목록에서 찾는다
curl -s localhost:8081/jobs/overview 응답을 /root/flink/cluster/jobs.json 에 저장하세요. 이름이 flk-first-orders 인 잡이 FINISHED 로 보여야 합니다.
JobManager 는 끝난 잡을 한동안 기억합니다(클러스터를 내리면 사라집니다). 이름이 안 보이면 first.sql 의 pipeline.name 이 SELECT 보다 앞에 있는지 보세요 — SET 은 그 뒤에 제출되는 잡부터 적용됩니다.
병렬도 2 로 다시 돌린다
같은 집계를 SET 'parallelism.default' = '2'; · 잡 이름 flk-first-p2 로 돌리는 /root/flink/cluster/p2.sql 을 만들어 출력을 /root/flink/cluster/p2.out 에 저장하고, 그 잡의 /jobs/<jid> 응답을 /root/flink/cluster/p2-job.json 에 저장하세요. 정점(vertices) 병렬도의 최댓값이 2 여야 하고, 집계 결과는 병렬도 1 일 때와 같아야 합니다.
병렬도를 2 로 줬는데 정점이 전부 1 이라면 배치 잡의 적응형 스케줄러가 자료 크기(작은 파일)를 보고 병렬도를 다시 정한 것입니다. 그 자동 결정을 끄는 설정이 execution.batch.adaptive.auto-parallelism 아래에 있습니다. 결과를 모으는 싱크는 1 로 남는 것이 정상입니다.
슬롯을 4개로 바꿔 다시 띄운다
/opt/flink/conf/config.yaml 의 taskmanager.numberOfTaskSlots 를 4 로 바꾸고 클러스터를 내렸다 다시 띄운 뒤 /overview 응답을 /root/flink/cluster/overview-4.json 에 저장하세요.
슬롯 수는 TaskManager 프로세스가 뜰 때 한 번 읽는 값입니다. 파일만 고치고 개요를 다시 받으면 여전히 2 가 나옵니다. flink-down 으로 내리고 flink-up 으로 올리세요. 같은 이름의 줄이 다른 블록에 있지 않은지도 보세요.
일부러 실패시키고 기록을 받는다
잡 이름을 flk-bad-cast 로 두고 status 를 INT 로 CAST 하는 SELECT 를 /root/flink/cluster/fail.sql 로 돌리세요. 그 잡의 /jobs/<jid> 를 /root/flink/cluster/fail-job.json, /jobs/<jid>/exceptions 를 /root/flink/cluster/fail-exceptions.json 에 저장하세요.
문자열 'paid' 를 정수로 바꾸는 순간 태스크가 예외를 던지고, 체크포인트가 꺼진 잡은 재시작하지 않고 바로 FAILED 가 됩니다. sql-client 출력에도 오류가 찍히지만, 채점기는 JobManager 가 남긴 예외 기록(exceptionHistory)을 봅니다.
보고서 — 원인과 재시작 전략을 가른다
/root/flink/cluster/report.json 에 slots_total(지금 클러스터의 슬롯 수), first_job_id·p2_job_id·failed_job_id(각 잡의 jid), restart_strategy(맨 위 예외 문장에서 재시작을 막은 전략 이름), root_cause(스택의 마지막 Caused by: 예외 클래스, 패키지 포함)를 적으세요.
맨 위 예외는 'Recovery is suppressed by …' 로 시작하는 JobException 입니다. 그것은 원인이 아니라 '재시작하지 않았다' 는 결과입니다. 진짜 원인은 스택을 따라 내려가 마지막 Caused by 에 있습니다. jid 는 앞에서 저장한 JSON 들에서 옮기면 됩니다.