Cloudera 에서 Flink 는 Cloudera Streaming Analytics(CSA) 로 배포되며 파셀 형태로 설치된다. Cloudera Manager 에는 Flink 서비스로 등록되고 실행은 YARN 위에서 이뤄진다. 이 문서는 설치된 환경에서 동작을 확인하는 최소 절차를 다룬다.
파셀 경로와 버전을 확인한다. 파셀 이름에 Flink 버전과 CSA 버전이 함께 들어간다.
ls -d /opt/cloudera/parcels/FLINK-*
flink --version
예를 들어 FLINK-1.19.1-csa1.13.2.0-... 이면 Flink 1.19.1 을 담은 CSA 1.13.2.0 이다. 예제 JAR 은 파셀 안의 Flink 배포 디렉터리에 있다.
FLINK_HOME=$(ls -d /opt/cloudera/parcels/FLINK-*/lib/flink)
ls ${FLINK_HOME}/examples/streaming/
Kerberos 가 켜진 클러스터라면 제출 전에 인증을 받는다. 이 단계를 건너뛰면 ZooKeeper JAAS 설정이 비어 생성되어 인증 오류가 뒤따른다.
kinit <principal>
klist
Flink 1.15 부터 job cluster 방식은 폐기됐다. -m yarn-cluster 대신 application 모드를 쓴다.
flink run-application -t yarn-application \
${FLINK_HOME}/examples/streaming/WordCount.jar
입력을 직접 넣어 보려면 소켓 예제를 쓴다. 한쪽 터미널에서 포트를 열어 두고 문장을 입력하면 다른 쪽에서 단어 수가 집계된다.
nc -lk 9000
flink run-application -t yarn-application \
${FLINK_HOME}/examples/streaming/SocketWindowWordCount.jar \
--hostname <host> --port 9000
제출하면 YARN 애플리케이션 ID 와 JobManager 웹 인터페이스 주소가 로그에 찍힌다.
Submitted application application_<ts>_<n>
Found Web Interface <host>:<port> of application 'application_<ts>_<n>'
Job has been submitted with JobID <job-id>
상태는 다음으로 본다.
flink list
yarn application -list
웹 UI 는 YARN ResourceManager 화면의 해당 애플리케이션에서 ApplicationMaster 링크를 따라가는 것이 확실하다. 포트가 매번 달라지기 때문이다.
SLF4J: Class path contains multiple SLF4J bindings. 는 Flink 파셀과 CDH 파셀의 로깅 바인딩이 함께 잡혀 나온다. 실제 바인딩이 무엇인지 뒤에 찍히며, 동작에는 영향이 없다.
Job Clusters are deprecated since Flink 1.15. 는 -m yarn-cluster 로 제출했을 때 나온다. application 모드로 바꾼다.
Cannot use kerberos delegation token manager, no valid kerberos credentials provided. 는 인증 없이 제출했다는 뜻이다. Flink 작업 제출 시 ZooKeeper SASL 설정이 없다는 오류 를 참고한다.
Trying to access closed classloader 는 작업 종료 후 정리 과정에서 나오는 경고로, 예제 실행에서는 무시해도 된다.
Flink 를 스트리밍으로 써 보려면 Kafka 토픽을 소스로 쓰는 것이 자연스럽다. 토픽을 만들고 콘솔 도구로 메시지 흐름을 먼저 확인한 뒤 Flink 작업을 붙인다.
kafka-topics --create --topic test-topic \
--bootstrap-server <broker>:9092 --partitions 1 --replication-factor 1
kafka-console-producer --topic test-topic --bootstrap-server <broker>:9092
kafka-console-consumer --topic test-topic --bootstrap-server <broker>:9092 --from-beginning
Kerberos 가 켜진 Kafka 라면 --command-config 또는 --consumer.config 로 보안 설정 파일을 함께 넘겨야 한다.