Flink Operations Playground
Flink Operations Playground
이 플레이그라운드는 Flink 애플리케이션을 작성하는 것보다 운영에 초점을 맞춰요. 플랫폼 엔지니어, DevOps 팀, 또는 프로덕션에서 Flink를 실행하고 관리하는 방법을 이해하려는 개발자에게 이상적이에요.
출처: 문서
본문
배우게 될 내용 (What You'll Learn)
이 플레이그라운드에서 다음을 하게 돼요.
- Web UI를 사용해 Flink 애플리케이션 배포·모니터링
- Flink가 exactly-once 보장으로 실패에서 어떻게 복구하는지 관찰
- savepoint를 사용해 job 업그레이드 수행
- job 확장·축소 (rescaling)
- REST API로 job 메트릭 조회
- 운영 작업에 Flink CLI 사용
다양한 환경에서 Apache Flink를 배포·운영하는 방법은 많아요. 이런 다양성과 관계없이 Flink 클러스터의 기본 구성 요소는 동일하게 유지되며, 비슷한 운영 원칙이 적용돼요.
이 플레이그라운드의 구조 (Anatomy of this Playground)
이 플레이그라운드는 오래 지속되는(long living) Flink Session Cluster와 Kafka Cluster로 구성돼요.
Flink 클러스터는 항상 JobManager와 하나 이상의 Flink TaskManager로 구성돼요. JobManager는 Job 제출 처리, Job의 감독, 리소스 관리를 담당해요. Flink TaskManager는 worker 프로세스이며, Flink Job을 구성하는 실제 Tasks의 실행을 담당해요. 이 플레이그라운드에서는 단일 TaskManager로 시작하지만, 나중에 더 많은 TaskManager로 확장할 거예요.
또한 이 플레이그라운드에는 전용 client 컨테이너가 함께 제공되는데, 이를 사용해 처음에 Flink Job을 제출하고 나중에 다양한 운영 작업을 수행해요. client 컨테이너는 Flink 클러스터 자체에 필요하지 않고, 사용 편의를 위해 포함된 것이에요.
Kafka 클러스터는 Zookeeper 서버와 Kafka Broker로 구성돼요.
플레이그라운드가 시작되면 Flink Event Count라는 Flink Job이 JobManager에 제출돼요. 또한 input과 output 두 개의 Kafka Topic이 생성돼요.
Job은 각각 타임스탬프와 페이지를 가진 ClickEvents를 input topic에서 소비해요. 이벤트는 페이지별로 keyed되고 15초 윈도우로 집계돼요. 결과는 output topic에 기록돼요.
여섯 개의 서로 다른 페이지가 있고 페이지당 15초마다 1000개의 클릭 이벤트를 생성해요. 따라서 Flink job의 출력은 페이지·윈도우당 1000번의 조회를 보여줘야 해요.
플레이그라운드 시작하기 (Starting the Playground)
플레이그라운드 환경은 몇 단계만으로 설정돼요. 필요한 명령을 안내하고 모든 것이 올바르게 실행되고 있는지 검증하는 방법을 보여줄게요.
머신에 Docker (20.10+)와 docker compose (2.1+)가 설치되어 있다고 가정해요. 필요한 구성 파일은 flink-playgrounds 저장소에 있어요. 먼저 코드를 체크아웃하고 docker 이미지를 빌드해요.
git clone https://github.com/apache/flink-playgrounds.git
cd flink-playgrounds/operations-playground
docker compose build
그다음 플레이그라운드를 시작해요.
docker compose up -d
그 후 다음 명령으로 실행 중인 Docker 컨테이너를 검사할 수 있어요.
docker compose ps
Name Command State Ports
-----------------------------------------------------------------------------------------------------------------------------
operations-playground_clickevent-generator_1 /docker-entrypoint.sh java ... Up 6123/tcp, 8081/tcp
operations-playground_client_1 /docker-entrypoint.sh flin ... Exit 0
operations-playground_jobmanager_1 /docker-entrypoint.sh jobm ... Up 6123/tcp, 0.0.0.0:8081->8081/tcp
operations-playground_kafka_1 start-kafka.sh Up 0.0.0.0:9094->9094/tcp
operations-playground_taskmanager_1 /docker-entrypoint.sh task ... Up 6123/tcp, 8081/tcp
operations-playground_zookeeper_1 /bin/sh -c /usr/sbin/sshd ... Up 2181/tcp, 22/tcp, 2888/tcp, 3888/tcp
이는 client 컨테이너가 Flink Job을 성공적으로 제출했고(Exit 0) 모든 클러스터 구성 요소와 데이터 생성기가 실행 중(Up)임을 나타내요.
플레이그라운드 환경을 중지하려면 다음을 호출해요.
docker compose down -v
플레이그라운드 진입 (Entering the Playground)
이 플레이그라운드에서 시도하고 확인할 수 있는 것이 많아요. 다음 두 절에서 Flink 클러스터와 상호작용하는 방법을 보여주고 Flink의 주요 기능 중 일부를 시연할 거예요.
Flink WebUI
Flink 클러스터를 관찰하는 가장 자연스러운 시작점은 http://localhost:8081에 노출된 WebUI예요. 모든 것이 잘 진행되었다면 클러스터가 처음에는 하나의 TaskManager로 구성되어 Click Event Count라는 Job을 실행함을 볼 수 있을 거예요.
Flink WebUI에는 Flink 클러스터와 그 Jobs에 대한 유용하고 흥미로운 정보(JobGraph, Metrics, Checkpointing Statistics, TaskManager Status 등)가 많이 포함되어 있어요.
로그 (Logs)
JobManager의 로그는 docker compose로 tail 수 있어요.
docker compose logs -f jobmanager
초기 시작 후에는 주로 모든 체크포인트 완료에 대한 로그 메시지를 볼 수 있어요.
TaskManager 로그도 같은 방식으로 tail할 수 있어요.
docker compose logs -f taskmanager
초기 시작 후에는 주로 모든 체크포인트 완료에 대한 로그 메시지를 볼 수 있어요.
Flink CLI
Flink CLI는 client 컨테이너 안에서 사용할 수 있어요. 예를 들어 Flink CLI의 도움말 메시지를 출력하려면 다음을 실행해요.
docker compose run --no-deps client flink --help
Flink REST API
Flink REST API는 호스트의 localhost:8081 또는 client 컨테이너에서 jobmanager:8081로 노출돼요. 예를 들어 현재 실행 중인 모든 job을 나열하려면 다음을 실행해요.
curl localhost:8081/jobs
Kafka Topics
Kafka Topic에 기록된 레코드를 보려면 다음을 실행해요.
//input topic (1000 records/s)
docker compose exec kafka kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic input
//output topic (24 records/min)
docker compose exec kafka kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic output
놀 시간! (Time to Play!)
이제 Flink와 Docker 컨테이너와 상호작용하는 법을 배웠으니, 플레이그라운드에서 시도할 수 있는 몇 가지 흔한 운영 작업을 살펴보자. 이러한 작업은 모두 서로 독립적이며, 즉 어떤 순서로든 수행할 수 있어요. 대부분의 작업은 CLI와 REST API로 실행할 수 있어요.
실행 중인 Job 나열 (Listing Running Jobs)
명령:
docker compose run --no-deps client flink list
예상 출력. 요청:
curl localhost:8081/jobs
예상 응답 (보기 좋게 출력). JobID는 제출 시 Job에 할당되며 CLI 또는 REST API로 Job에 작업을 수행하는 데 필요해요.
실패 및 복구 관찰 (Observing Failure & Recovery)
Flink는 (부분) 실패에서 exactly-once 처리 보장을 제공해요. 이 플레이그라운드에서 이 동작을 관찰하고 어느 정도 검증할 수 있어요.
1단계: 출력 관찰 (Observing the Output)
위에서 설명했듯이 이 플레이그라운드의 이벤트는 각 윈도우가 정확히 1000개의 레코드를 포함하도록 생성돼요. 따라서 Flink가 데이터 손실이나 중복 없이 TaskManager 실패에서 성공적으로 복구하는지 검증하려면, output topic을 tail하여 복구 후 모든 윈도우가 존재하고 개수가 올바른지 확인하면 돼요. 이를 위해 output topic에서 읽기를 시작하고 이 명령을 복구(3단계) 후까지 계속 실행해 두세요.
docker compose exec kafka kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic output
2단계: 장애 주입 (Introducing a Fault)
부분 실패를 시뮬레이션하려면 TaskManager를 kill할 수 있어요. 프로덕션 구성에서 이는 TaskManager 프로세스 상실, TaskManager 머신 상실, 또는 단순히 프레임워크나 사용자 코드에서 (예: 외부 리소스의 일시적 사용 불가로 인해) 던져지는 일시적 예외에 해당할 수 있어요.
docker compose kill taskmanager
몇 초 후 JobManager는 TaskManager의 상실을 인지하고, 영향을 받은 Job을 취소한 뒤 즉시 복구를 위해 다시 제출해요. Job이 재시작되면 task는 SCHEDULED 상태로 유지되는데, 이는 보라색 사각형으로 표시돼요 (아래 스크린샷 참고).
이 시점에서 Job의 task는 SCHEDULED 상태에서 RUNNING으로 이동할 수 없는데, task를 실행할 리소스(TaskManager가 제공하는 TaskSlot)가 없기 때문이에요. 새 TaskManager가 사용 가능해질 때까지 Job은 취소와 재제출의 순환을 겪게 돼요.
그 동안 데이터 생성기는 ClickEvents를 input topic에 계속 밀어 넣어요. 이는 Job이 다운되어 있는 동안 데이터가 생성되는 실제 프로덕션 구성과 유사해요.
3단계: 복구 (Recovery)
TaskManager를 재시작하면 JobManager에 다시 연결돼요.
docker compose up -d taskmanager
JobManager가 새 TaskManager를 통지받으면 복구 중인 Job의 task를 새로 사용 가능한 TaskSlot에 스케줄해요. 재시작 시 task는 실패 전에 만든 마지막 성공적인 체크포인트에서 상태를 복구하고 RUNNING 상태로 전환해요.
Job은 중단 동안 축적된 입력 이벤트의 전체 백로그를 Kafka에서 빠르게 처리하고, 스트림의 헤드에 도달할 때까지 훨씬 더 높은 속도(> 24 records/minute)로 출력을 생성해요. 출력에서 모든 키(페이지)가 모든 시간 윈도우에 대해 존재하고 모든 개수가 정확히 1000임을 볼 수 있을 거예요. FlinkKafkaProducer를 "at-least-once" 모드로 사용하고 있으므로, 일부 중복 출력 레코드가 보일 가능성이 있어요.
참고: 대부분의 프로덕션 구성은 리소스 매니저(Kubernetes, Yarn)에 의존해 실패한 프로세스를 자동으로 재시작해요.
Job 업그레이드 및 리스케일링 (Upgrading & Rescaling a Job)
Flink Job 업그레이드는 항상 두 단계를 포함해요. 첫째, Flink Job은 Savepoint로 정상적으로(gracefully) 중지돼요. Savepoint는 잘 정의된 전역적으로 일관된 시점에서 완전한 애플리케이션 상태의 일관된 스냅샷이에요 (체크포인트와 유사). 둘째, 업그레이드된 Flink Job이 Savepoint에서 시작돼요. 여기서 "업그레이드"는 다음을 포함한 서로 다른 것을 의미할 수 있어요.
- 구성 업그레이드 (Job의 병렬도 포함)
- Job의 토폴로지 업그레이드 (추가/제거된 Operators)
- Job의 사용자 정의 함수 업그레이드
업그레이드를 시작하기 전에 output topic을 tail해서 업그레이드 과정에서 데이터가 손실되거나 손상되지 않는지 관찰하고 싶을 수 있어요.
docker compose exec kafka kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic output
1단계: Job 중지 (Stopping the Job)
Job을 정상적으로 중지하려면 CLI 또는 REST API의 "stop" 명령을 사용해야 해요. 이를 위해 모든 실행 중인 Job을 나열하거나 WebUI에서 얻을 수 있는 Job의 JobID가 필요해요. JobID로 Job 중지를 진행할 수 있어요.
2a단계: 변경 없이 Job 재시작 (Restart Job without Changes)
이제 이 Savepoint에서 업그레이드된 Job을 재시작할 수 있어요. 간단히 하기 위해 변경 없이 재시작하는 것으로 시작할 수 있어요.
# Uploading the JAR from the Client container
docker compose run --no-deps client curl -X POST -H "Expect:" \
-F "jarfile=@/opt/ClickCountJob.jar" http://jobmanager:8081/jars/upload