Apache Flink — 스트림을 엔진으로 돌린다 · 클러스터 한 벌 · 实验
클러스터를 띄우고 잡 하나를 끝까지 따라간다
목표
Flink 로컬 클러스터를 띄우고 REST API 로 구조를 확인한 뒤, 배치 잡 하나를 제출부터 종료까지 따라간다. 병렬도·슬롯·실패 기록을 직접 바꿔 보고 그 결과를 보고서로 정리한다.
왜 중요한가
운영에서 받는 질문 — 잡이 왜 멈췄나, 병렬도를 올렸는데 왜 그대로인가, 왜 재시작되지 않았나 — 은 SQL 이 아니라 엔진 구조에 대한 것이다. JobManager 는 잡마다 JobMaster 를 만들고 슬롯을 나눠 주며, 모든 기록을 REST 로 내놓는다. 이 실습의 채점기는 클러스터에 묻지 않는다. 여러분이 파일로 저장한 REST 응답과 sql-client 출력만 읽고, 집계 기대값은 원본 CSV 에서 직접 계산해 대조한다. 그래서 응답은 손으로 고치지 말고 curl 출력 그대로 저장한다.
단계
1. flink-up 으로 클러스터를 띄우고 /overview 응답을 /root/flink/cluster/overview.json 에 저장하세요(TaskManager 1개 · 슬롯 2개가 보여야 합니다).
2. /taskmanagers 응답을 /root/flink/cluster/taskmanagers.json 에 저장하고, 그 memoryConfiguration 을 MiB 로 바꿔 반올림한 값 일곱 개를 /root/flink/cluster/memory.json 에 적으세요.
3. /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 에 저장하세요.
4. 잡이 끝난 뒤 /jobs/overview 응답을 /root/flink/cluster/jobs.json 에 저장하세요.
5. 같은 집계를 병렬도 2 · 잡 이름 flk-first-p2 로 돌리는 /root/flink/cluster/p2.sql 을 만들어 출력을 /root/flink/cluster/p2.out 에, 그 잡의 /jobs/<jid> 응답을 /root/flink/cluster/p2-job.json 에 저장하세요. 정점 병렬도에 2 가 보여야 합니다.
6. /opt/flink/conf/config.yaml 의 taskmanager.numberOfTaskSlots 를 4 로 바꾸고 클러스터를 다시 띄운 뒤 /overview 를 /root/flink/cluster/overview-4.json 에 저장하세요.
7. 잡 이름 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 에 저장하세요.
8. /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](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/flink-architecture/) · [REST API](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/rest_api/) · [TaskManager 메모리](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/deployment/memory/mem_setup_tm/) · [Adaptive Batch](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/deployment/adaptive_batch/) · [Task Failure Recovery](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/state/task_failure_recovery/)
8个步骤
- 클러스터를 띄우고 개요를 받는다
- TaskManager 메모리 예산을 숫자로 읽는다
- 배치 잡 하나를 제출한다
- 끝난 잡을 목록에서 찾는다
- 병렬도 2 로 다시 돌린다
- 슬롯을 4개로 바꿔 다시 띄운다
- 일부러 실패시키고 기록을 받는다
- 보고서 — 원인과 재시작 전략을 가른다