YARN
YARN
이 문서는 YARN 위에 완전히 기능하는 Flink 클러스터를 설정하는 방법을 안내합니다. Apache Hadoop YARN은 많은 데이터 처리 프레임워크에서 인기 있는 리소스 제공자이며, Flink는 Application Mode와 Session Mode로 YARN에 배포됩니다.
출처: 문서
본문
시작하기 (Getting Started)
이 Getting Started 섹션은 YARN 위에 완전히 기능하는 Flink 클러스터를 설정하는 방법을 안내합니다.
소개 (Introduction)
Apache Hadoop YARN은 많은 데이터 처리 프레임워크에서 인기 있는 리소스 제공자입니다. Flink 서비스는 YARN의 ResourceManager에 제출되며, ResourceManager는 YARN NodeManager가 관리하는 머신에 컨테이너를 생성합니다. Flink는 그러한 컨테이너 안에 JobManager와 TaskManager 인스턴스를 배포합니다.
Flink는 JobManager에서 실행 중인 작업이 요구하는 처리 슬롯 수에 따라 TaskManager 리소스를 동적으로 할당하고 해제할 수 있습니다.
준비 (Preparation)
이 Getting Started 섹션은 버전 2.10.2부터 시작하는 기능적인 YARN 환경을 가정합니다. YARN 환경은 Amazon EMR, Google Cloud DataProc 같은 서비스 또는 Cloudera 같은 제품을 통해 가장 편리하게 제공됩니다. 로컬에서 YARN 환경을 수동으로 설정하거나 클러스터에서 설정하는 것은 이 Getting Started 튜토리얼을 진행하는 데 권장되지 않습니다.
yarn top을 실행하여 YARN 클러스터가 Flink 애플리케이션을 수락할 준비가 되었는지 확인하세요. 오류 메시지가 없어야 합니다.- 다운로드 페이지에서 최신 Flink 배포판을 다운로드하고 압축을 풀어주세요.
- 중요
HADOOP_CLASSPATH환경 변수가 설정되어 있는지 확인하세요(echo $HADOOP_CLASSPATH로 확인 가능). 없으면 다음과 같이 설정하세요:
export HADOOP_CLASSPATH=`hadoop classpath`
YARN에서 Flink 세션 시작 (Starting a Flink Session on YARN)
HADOOP_CLASSPATH 환경 변수가 설정되었는지 확인한 후 Flink on YARN 세션을 시작하고 예제 작업을 제출할 수 있습니다:
# we assume to be in the root directory of
# the unzipped Flink distribution
# (0) export HADOOP_CLASSPATH
export HADOOP_CLASSPATH=`hadoop classpath`
# (1) Start YARN Session
./bin/yarn-session.sh --detached
# (2) You can now access the Flink Web Interface through the
# URL printed in the last lines of the command output, or through
# the YARN ResourceManager web UI.
# (3) Submit example job
./bin/flink run ./examples/streaming/TopSpeedWindowing.jar
# (4) Stop YARN session (replace the application id based
# on the output of the yarn-session.sh command)
echo "stop" | ./bin/yarn-session.sh -id application_XXXXX_XXX
축하합니다! Flink를 YARN에 배포하여 Flink 애플리케이션을 성공적으로 실행했습니다.
Flink on YARN이 지원하는 배포 모드 (Deployment Modes Supported by Flink on YARN)
프로덕션 사용 시에는 애플리케이션 간 격리가 더 좋기 때문에 Application Mode로 Flink 애플리케이션을 배포하는 것을 권장합니다.
Application Mode
Application Mode 뒤의 고수준 직관에 대해서는 deployment mode overview를 참고하세요.
Application Mode는 YARN에서 Flink 클러스터를 시작하며, 애플리케이션 jar의 main() 메서드가 YARN의 JobManager에서 실행됩니다. 애플리케이션이 끝나면 즉시 클러스터가 종료됩니다. yarn application -kill <ApplicationId>로 클러스터를 수동으로 중지하거나 Flink 작업을 취소할 수 있습니다.
./bin/flink run -t yarn-application ./examples/streaming/TopSpeedWindowing.jar
Application Mode 클러스터가 배포되면 취소나 savepoint 수행 같은 작업을 위해 상호작용할 수 있습니다.
# List running job on the cluster
./bin/flink list -t yarn-application -Dyarn.application.id=application_XXXX_YY
# Cancel running job
./bin/flink cancel -t yarn-application -Dyarn.application.id=application_XXXX_YY <jobId>
Application Cluster에서 작업을 취소하면 클러스터가 중지됨에 유의하세요.
Application mode의 전체 잠재력을 활용하려면 yarn.provided.lib.dirs 구성 옵션과 함께 사용하고 애플리케이션 jar를 클러스터의 모든 노드가 접근할 수 있는 위치에 사전 업로드하는 것을 고려하세요. 이 경우 명령은 다음과 같습니다:
./bin/flink run -t yarn-application \
-Dyarn.provided.lib.dirs="hdfs://myhdfs/my-remote-flink-dist-dir" \
hdfs://myhdfs/jars/my-application.jar
위 명령은 필요한 Flink jar와 애플리케이션 jar를 클라이언트가 클러스터로 보내는 대신 지정된 원격 위치에서 가져오므로 작업 제출을 매우 가볍게 만들어 줍니다.
Session Mode
Session mode 뒤의 고수준 직관에 대해서는 deployment mode overview를 참고하세요.
Session Mode 배포는 페이지 상단의 Getting Started 가이드에서 설명합니다.
Session Mode에는 두 가지 운영 모드가 있습니다:
- attached mode(기본값):
yarn-session.sh클라이언트가 Flink 클러스터를 YARN에 제출하지만, 클라이언트는 실행을 계속하며 클러스터 상태를 추적합니다. 클러스터가 실패하면 클라이언트가 오류를 표시합니다. 클라이언트가 종료되면 클러스터도 종료하도록 신호를 보냅니다. - detached mode(
-d또는--detached):yarn-session.sh클라이언트가 Flink 클러스터를 YARN에 제출한 뒤 클라이언트가 반환합니다. Flink 클러스터를 중지하려면 클라이언트를 다시 호출하거나 YARN 도구를 사용해야 합니다.
Session mode는 /tmp/.yarn-properties-<username>에 숨겨진 YARN 속성 파일을 만듭니다. 이 파일은 작업 제출 시 명령 줄 인터페이스가 클러스터를 발견하는 데 사용합니다.
Flink 작업 제출 시 명령 줄 인터페이스에서 대상 YARN 클러스터를 수동으로 지정할 수도 있습니다. 예시는 다음과 같습니다:
./bin/flink run -t yarn-session \
-Dyarn.application.id=application_XXXX_YY \
./examples/streaming/TopSpeedWindowing.jar
다음 명령으로 YARN 세션에 다시 연결할 수 있습니다:
./bin/yarn-session.sh -id application_XXXX_YY
Flink 구성 파일을 통한 구성 전달 외에도, 제출 시 -Dkey=value 인자로 ./bin/yarn-session.sh 클라이언트에 어떤 구성이든 전달할 수 있습니다.
YARN 세션 클라이언트는 자주 쓰이는 설정을 위한 몇 가지 "단축 인자(shortcut arguments)"도 있습니다. ./bin/yarn-session.sh -h로 나열할 수 있습니다.
Flink on YARN 참조 (Flink on YARN Reference)
Flink on YARN 구성 (Configuring Flink on YARN)
YARN 특정 구성은 구성 페이지에 나열되어 있습니다.
다음 구성 파라미터는 Flink on YARN이 관리합니다. 런타임에 프레임워크가 덮어쓸 수 있기 때문입니다:
jobmanager.rpc.address(Flink on YARN이 JobManager 컨테이너의 주소로 동적으로 설정)io.tmp.dirs(설정되지 않으면 Flink가 YARN이 정의한 임시 디렉터리를 설정)high-availability.cluster-id(HA 서비스에서 여러 클러스터를 구분하기 위해 자동 생성된 ID)
추가 Hadoop 구성 파일을 Flink에 전달해야 한다면, Hadoop 구성 파일을 포함하는 디렉터리 이름을 받는 HADOOP_CONF_DIR 환경 변수로 전달할 수 있습니다. 기본적으로 필요한 모든 Hadoop 구성 파일은 HADOOP_CLASSPATH 환경 변수의 클래스패스에서 로드됩니다.
리소스 할당 동작 (Resource Allocation Behavior)
YARN에서 실행되는 JobManager는 기존 리소스로 제출된 모든 작업을 실행할 수 없으면 추가 TaskManager를 요청합니다. 특히 Session Mode에서 실행할 때 JobManager는 필요하면 추가 작업이 제출됨에 따라 추가 TaskManager를 할당합니다. 사용되지 않는 TaskManager는 타임아웃 후 다시 해제됩니다.
JobManager와 TaskManager 프로세스의 메모리 구성은 YARN 구현이 존중합니다. 보고되는 VCores 수는 기본적으로 TaskManager당 구성된 슬롯 수와 같습니다. yarn.containers.vcores는 vcores 수를 사용자 지정 값으로 덮어쓸 수 있게 합니다. 이 파라미터가 동작하려면 YARN 클러스터에서 CPU 스케줄링을 활성화해야 합니다.
실패한 컨테이너(JobManager 포함)는 YARN이 교체합니다. JobManager 컨테이너 재시작의 최대 횟수는 yarn.application-attempts(기본 1)로 구성합니다. 모든 시도를 소진하면 YARN Application이 실패합니다.
YARN에서의 고가용성 (High-Availability on YARN)
YARN에서의 고가용성은 YARN과 high availability service의 조합으로 달성됩니다.
HA 서비스가 구성되면 JobManager 메타데이터를 영속화하고 리더 선출을 수행합니다.
YARN이 실패한 JobManager 재시작을 담당합니다. JobManager 재시작의 최대 횟수는 두 구성 파라미터로 정의됩니다. 첫째, Flink의 yarn.application-attempts 구성은 기본 2입니다. 이 값은 YARN의 yarn.resourcemanager.am.max-attempts로 제한되는데, 이것도 기본 2입니다.
YARN에 배포할 때 Flink가 high-availability.cluster-id 구성 파라미터를 관리한다는 점에 유의하세요. Flink는 기본적으로 이를 YARN 애플리케이션 id로 설정합니다. YARN에 HA 클러스터를 배포할 때 이 파라미터를 덮어쓰지 마세요. 클러스터 ID는 HA 백엔드(예: Zookeeper)에서 여러 HA 클러스터를 구분하는 데 사용됩니다. 이 구성 파라미터를 덮어쓰면 여러 YARN 클러스터가 서로에게 영향을 줄 수 있습니다.
컨테이너 종료 동작 (Container Shutdown Behaviour)
- YARN 2.3.0 < version < 2.4.0. application master가 실패하면 모든 컨테이너가 재시작됩니다.
- YARN 2.4.0 < version < 2.6.0. TaskManager 컨테이너는 application master 실패를 넘어 유지됩니다. 이는 시작 시간이 더 빠르고 사용자가 컨테이너 리소스를 다시 얻기 위해 기다릴 필요가 없다는 장점이 있습니다.
- YARN 2.6.0 <= version: 시도 실패 유효성 간격(attempt failure validity interval)을 Flink의 Pekko 타임아웃 값으로 설정합니다. attempt failure validity interval은 시스템이 한 간격 동안 최대 애플리케이션 시도 횟수를 본 후에만 애플리케이션이 종료된다는 것을 뜻합니다. 이는 장기 실행 작업이 애플리케이션 시도를 소진하는 것을 방지합니다.
위험: Hadoop YARN 2.4.0에는 재시작된 Application Master/Job Manager 컨테이너에서 컨테이너 재시작을 막는 주요 버그가 있습니다(2.5.0에서 수정됨). 자세한 내용은 FLINK-4142 참고. YARN의 고가용성 설정에는 최소 Hadoop 2.5.0 사용을 권장합니다.
지원되는 Hadoop 버전 (Supported Hadoop versions)
Flink on YARN은 Hadoop 2.10.2로 컴파일되어 있으며, Hadoop 3.x를 포함한 모든 Hadoop 버전 >= 2.10.2가 지원됩니다.
Flink에 필요한 Hadoop 의존성을 제공하려면 Getting Started / Preparation 섹션에서 이미 소개한 HADOOP_CLASSPATH 환경 변수를 설정하는 것을 권장합니다.
그것이 불가능하면 의존성을 Flink의 lib/ 폴더에 넣을 수도 있습니다.
Flink는 또한 웹사이트의 Downloads / Additional Components 섹션에서 lib/ 폴더에 넣기 위한 사전 번들된 Hadoop fat jar를 제공합니다. 이러한 사전 번들된 fat jar는 일반적인 라이브러리와의 의존성 충돌을 피하기 위해 셰이드(shaded)되어 있습니다. Flink 커뮤니티는 이러한 사전 번들된 jar에 대해 YARN 통합을 테스트하지 않습니다.
방화벽 뒤에서 Flink on YARN 실행 (Running Flink on YARN behind Firewalls)
일부 YARN 클러스터는 클러스터와 네트워크의 나머지 사이의 네트워크 트래픽을 제어하기 위해 방화벽을 사용합니다. 그러한 설정에서 Flink 작업은 클러스터 네트워크 내부(방화벽 뒤)에서만 YARN 세션에 제출할 수 있습니다. 이것이 프로덕션 사용에 적합하지 않다면, Flink는 클라이언트-클러스터 통신에 사용되는 REST 엔드포인트에 포트 범위를 구성할 수 있게 합니다. 이 범위를 구성하면 사용자가 방화벽을 넘어 Flink에 작업을 제출할 수도 있습니다.
REST 엔드포인트 포트를 지정하는 구성 파라미터는 rest.bind-port입니다. 이 구성 옵션은 단일 포트(예: "50010"), 범위("50000-50025"), 또는 둘의 조합을 받습니다.
사용자 jar와 클래스패스 (User jars & Classpath)
Session Mode
Yarn에서 Session Mode로 Flink를 배포할 때 시작 명령에 지정된 JAR 파일만 user-jar로 인식되어 user classpath에 포함됩니다.
Application Mode
Yarn에서 Application Mode로 Flink를 배포할 때 시작 명령에 지정된 JAR 파일과 Flink의 usrlib 폴더의 모든 JAR 파일이 user-jar로 인식됩니다. 기본적으로 Flink는 user-jar를 시스템 클래스패스에 포함합니다. 이 동작은 yarn.classpath.include-user-jar 파라미터로 제어할 수 있습니다.
이를 DISABLED로 설정하면 Flink는 jar를 사용자 클래스패스에 대신 포함합니다.
클래스패스에서 user-jar의 위치는 파라미터를 다음 중 하나로 설정하여 제어할 수 있습니다:
ORDER: (기본값) 사전순 순서에 따라 jar를 시스템 클래스패스에 추가.FIRST: jar를 시스템 클래스패스의 시작에 추가.LAST: jar를 시스템 클래스패스의 끝에 추가.
자세한 내용은 Debugging Classloading Docs를 참고하세요.