커맨드라인 인터페이스
커맨드라인 인터페이스 (Command-Line Interface)
Flink는 JAR 파일로 패키징된 프로그램을 실행하고 그 실행을 제어하기 위해 커맨드라인 인터페이스(CLI) bin/flink를 제공합니다. CLI는 모든 Flink 설정의 일부이며, 로컬 단일 노드 설정과 분산 설정 모두에서 사용할 수 있습니다. Flink 구성 파일에 지정된 실행 중인 JobManager에 연결합니다.
출처: 문서
본문
작업 수명주기 관리
이 섹션에 나열된 명령이 작동하기 위한 전제 조건은 Kubernetes, YARN 또는 다른 사용 가능한 옵션 같은 실행 중인 Flink 배포가 있다는 것입니다. 자신의 머신에서 명령을 시도하려면 Flink 클러스터를 로컬로 시작하세요.
작업 제출
작업 제출은 작업의 JAR와 관련 의존성을 Flink 클러스터에 업로드하고 작업 실행을 시작하는 것을 의미합니다. 이 예시를 위해 examples/streaming/StateMachineExample.jar 같은 장기 실행 작업을 선택합니다. examples/ 폴더의 다른 JAR 아카이브를 선택하거나 자체 작업을 배포할 수 있습니다.
$ ./bin/flink run \
--detached \
./examples/streaming/StateMachineExample.jar
--detached로 작업을 제출하면 제출이 완료된 후 명령이 반환됩니다. 출력에는 (다른 것들 외에도) 새로 제출된 작업의 ID가 포함됩니다.
Usage with built-in data generator: StateMachineExample [--error-rate <probability-of-invalid-transition>] [--sleep <sleep-per-record-in-ms>]
Usage with Kafka: StateMachineExample --kafka-topic <topic> [--brokers <brokers>]
Options for both the above setups:
[--backend <file|rocks>]
[--checkpoint-dir <filepath>]
[--async-checkpoints <true|false>]
[--incremental-checkpoints <true|false>]
[--output <filepath> OR null for stdout]
Using standalone source with error rate 0.000000 and sleep delay 1 millis
Job has been submitted with JobID cca7bc1061d61cf15238e92312c2fc20
출력된 사용 정보는 필요하면 작업 제출 명령 끝에 추가할 수 있는 작업 관련 매개변수를 나열합니다. 가독성을 위해 아래 명령에서 반환된 JobID가 변수 JOB_ID에 저장되어 있다고 가정합니다.
$ export JOB_ID="cca7bc1061d61cf15238e92312c2fc20"
run 명령은 -D 인수를 통해 추가 구성 매개변수 전달을 지원합니다. 예를 들어 작업의 최대 병렬도를 -Dpipeline.max-parallelism=120으로 설정할 수 있습니다. 이 인수는 구성 파일을 변경하지 않고 구성 파일을 변경하지 않고도 클러스터에 어떤 구성 매개변수든 전달할 수 있으므로 애플리케이션 모드 클러스터를 구성하는 데 매우 유용합니다.
기존 세션 클러스터에 작업을 제출할 때는 실행 구성 매개변수만 지원됩니다.
작업 모니터링
list 액션으로 실행 중인 작업을 모니터링할 수 있습니다.
$ ./bin/flink list
Waiting for response...
------------------ Running/Restarting Jobs -------------------
30.11.2020 16:02:29 : cca7bc1061d61cf15238e92312c2fc20 : State machine job (RUNNING)
--------------------------------------------------------------
No scheduled jobs.
제출되었지만 아직 시작되지 않은 작업은 "Scheduled Jobs" 아래에 나열됩니다.
Savepoint 생성
Savepoint은 작업이 있는 현재 상태를 저장하기 위해 생성할 수 있습니다. 필요한 것은 JobID뿐입니다.
$ ./bin/flink savepoint \
$JOB_ID \
/tmp/flink-savepoints
Triggering savepoint for job cca7bc1061d61cf15238e92312c2fc20.
Waiting for response...
Savepoint completed. Path: file:/tmp/flink-savepoints/savepoint-cca7bc-bb1e257f0dab
You can resume your program from this savepoint with the run command.
savepoint 폴더는 선택 사항이며 execution.checkpointing.savepoint-dir이 설정되지 않은 경우 지정해야 합니다.
마지막으로 savepoint의 이진 포맷이 무엇이어야 하는지 선택적으로 제공할 수 있습니다.
savepoint 경로는 나중에 Flink 작업 재시작에 사용할 수 있습니다.
작업의 상태가 꽤 크면 클라이언트는 savepoint이 완료될 때까지 기다려야 하므로 타임아웃 예외를 얻게 됩니다.
Triggering savepoint for job bec5244e09634ad71a80785937a9732d.
Waiting for response...
--------------------------------------------------------------
The program finished with the following exception:
org.apache.flink.util.FlinkException: Triggering a savepoint for the job bec5244e09634ad71a80785937a9732d failed.
at org.apache.flink.client.cli.CliFrontend.triggerSavepoint(CliFrontend. java:828)
at org.apache.flink.client.cli.CliFrontend.lambda$savepopint$8(CliFrontend.java:794)
at org.apache.flink.client.cli.CliFrontend.runClusterAction(CliFrontend.java:1078)
at org.apache.flink.client.cli.CliFrontend.savepoint(CliFrontend.java:779)
at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1150)
at org.apache.flink.client.cli.CliFrontend.lambda$mainInternal$9(CliFrontend.java:1226)
at org.apache.flink.runtime.security.contexts.NoOpSecurityContext.runSecured(NoOpSecurityContext.java:28)
at org.apache.flink.client.cli.CliFrontend.mainInternal(CliFrontend.java:1226)
at org.apache.flink.client.cli.CliFrontend.main(CliFronhtend.java:1194)
Caused by: java.util.concurrent.TimeoutException
at java.util.concurrent.CompletableFuture.timedGet(CompletableFuture.java:1784)
at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1928)
at org.apache.flink.client.cli.CliFrontend.triggerSavepoint(CliFrontend.java:822)
... 8 more
이 경우 "-detached" 옵션을 사용해 분리(detached) savepoint을 트리거할 수 있으며, 클라이언트는 트리거 id가 반환되는 즉시 반환합니다.
$ ./bin/flink savepoint \
$JOB_ID \
/tmp/flink-savepoints
-detached
Triggering savepoint in detached mode for job bec5244e09634ad71a80785937a9732d.
Successfully trigger manual savepoint, triggerId: 2505bbd12c5b58fd997d0f193db44b97
분리 savepoint의 상태는 rest api로 알 수 있습니다.
Savepoint 폐기
savepoint 액션은 savepoint을 제거하는 데도 사용할 수 있습니다. 해당 savepoint 경로와 함께 --dispose를 추가해야 합니다.
$ ./bin/flink savepoint \
--dispose \
/tmp/flink-savepoints/savepoint-cca7bc-bb1e257f0dab \
$JOB_ID
Disposing savepoint '/tmp/flink-savepoints/savepoint-cca7bc-bb1e257f0dab'.
Waiting for response...
Savepoint '/tmp/flink-savepoints/savepoint-cca7bc-bb1e257f0dab' disposed.
사용자 지정 상태 인스턴스(예: 사용자 지정 reducing state 또는 RocksDB 상태)를 사용한다면 savepoint이 트리거된 프로그램 JAR 경로를 지정해야 합니다. 그렇지 않으면 ClassNotFoundException이 발생합니다.
$ ./bin/flink savepoint \
--dispose <savepointPath> \
--jarfile <jarFile>
savepoint 액션으로 savepoint 폐기를 트리거하는 것은 저장소에서 데이터를 제거할 뿐만 아니라 Flink가 savepoint 관련 메타데이터도 정리하게 합니다.
Checkpoint 생성
Checkpoint도 수동으로 생성하여 현재 상태를 저장할 수 있습니다. checkpoint와 savepoint의 차이는 Checkpoints vs. Savepoints를 참조하세요. checkpoint를 수동으로 트리거하는 데 필요한 것은 JobID뿐입니다.
$ ./bin/flink checkpoint \
$JOB_ID
Triggering checkpoint for job 99c59fead08c613763944f533bf90c0f.
Waiting for response...
Checkpoint(CONFIGURED) 26 for job 99c59fead08c613763944f533bf90c0f completed.
You can resume your program from this checkpoint with the run command.
작업이 주기적으로 증분 checkpoint를 트리거하는 동안 전체 checkpoint를 트리거하려면 --full 옵션을 사용하세요.
$ ./bin/flink checkpoint \
$JOB_ID \
--full
Triggering checkpoint for job 99c59fead08c613763944f533bf90c0f.
Waiting for response...
Checkpoint(FULL) 78 for job 99c59fead08c613763944f533bf90c0f completed.
You can resume your program from this checkpoint with the run command.
작업 종료
최종 Savepoint을 생성하며 작업을 정상적으로 중지
작업을 중지하는 또 다른 액션은 stop입니다. stop은 소스에서 싱크로 흐르므로 실행 중인 스트리밍 작업을 중지하는 더 우아한 방법입니다. 사용자가 작업 중지를 요청하면 모든 소스가 savepoint을 트리거할 마지막 checkpoint barrier를 보내도록 요청되며, 그 savepoint이 성공적으로 완료된 후 자신의 cancel() 메서드를 호출하여 종료합니다.
$ ./bin/flink stop \
--savepointPath /tmp/flink-savepoints \
$JOB_ID
Suspending job "cca7bc1061d61cf15238e92312c2fc20" with a savepoint.
Savepoint completed. Path: file:/tmp/flink-savepoints/savepoint-cca7bc-bb1e257f0dab
execution.checkpointing.savepoint-dir이 설정되지 않은 경우 --savepointPath로 savepoint 폴더를 지정해야 합니다.
--drain 플래그가 지정되면 마지막 checkpoint barrier 전에 MAX_WATERMARK가 방출됩니다. 이렇게 하면 등록된 모든 이벤트 시간 타이머가 실행되어 특정 워터마크를 기다리는 모든 상태(예: 윈도우)가 플러시됩니다. 작업은 모든 소스가 제대로 종료될 때까지 계속 실행됩니다. 이는 작업이 진행 중(in-flight) 데이터를 모두 처리할 수 있게 하며, 중지 중에 savepoint을 찍은 후 처리할 일부 레코드를 생성할 수 있습니다.
작업을 영구적으로 종료하려면
--drain플래그를 사용하세요. 나중에 작업을 재개하려면 파이프라인을 drain하지 마세요. 작업이 재개될 때 잘못된 결과를 초래할 수 있습니다.
savepoint을 분리 모드로 트리거하려면 명령에 -detached 옵션을 추가하세요.
마지막으로 savepoint의 이진 포맷이 무엇이어야 하는지 선택적으로 제공할 수 있습니다.
작업을 비정상적으로 취소
cancel 액션으로 작업을 취소할 수 있습니다.
$ ./bin/flink cancel $JOB_ID
Cancelling job cca7bc1061d61cf15238e92312c2fc20.
Cancelled job cca7bc1061d61cf15238e92312c2fc20.
해당 작업의 상태는 Running에서 Cancelled로 전환됩니다. 모든 계산이 중지됩니다.
--withSavepoint플래그는 작업 취소의 일부로 savepoint을 만들 수 있게 합니다. 이 기능은 deprecated입니다. 대신 stop 액션을 사용하세요.
Savepoint에서 작업 시작
run 액션으로 savepoint에서 작업을 시작할 수 있습니다.
$ ./bin/flink run \
--detached \
--fromSavepoint /tmp/flink-savepoints/savepoint-cca7bc-bb1e257f0dab \
./examples/streaming/StateMachineExample.jar
Usage with built-in data generator: StateMachineExample [--error-rate <probability-of-invalid-transition>] [--sleep <sleep-per-record-in-ms>]
Usage with Kafka: StateMachineExample --kafka-topic <topic> [--brokers <brokers>]
Options for both the above setups:
[--backend <file|rocks>]
[--checkpoint-dir <filepath>]
[--async-checkpoints <true|false>]
[--incremental-checkpoints <true|false>]
[--output <filepath> OR null for stdout]
Using standalone source with error rate 0.000000 and sleep delay 1 millis
Job has been submitted with JobID 97b20a0a8ffd5c1d656328b0cd6436a6
명령이 초기 run 명령과 동일하지만, 이전에 중지된 작업의 상태를 참조하는 데 사용되는 --fromSavepoint 매개변수만 다릅니다. 작업을 유지하는 데 사용할 수 있는 새 JobID가 생성됩니다.
기본적으로 우리는 전체 savepoint 상태를 제출되는 작업에 매칭하려고 시도합니다. 새 작업으로 복원할 수 없는 savepoint 상태를 건너뛰는 것을 허용하려면 --allowNonRestoredState 플래그를 설정할 수 있습니다. savepoint이 트리거될 때 프로그램의 일부였던 연산자를 프로그램에서 제거했고 여전히 savepoint을 사용하려면 이를 허용해야 합니다.
$ ./bin/flink run \
--fromSavepoint <savepointPath> \
--allowNonRestoredState ...
이는 프로그램이 savepoint의 일부였던 연산자를 버린 경우 유용합니다.
또한 savepoint에 사용해야 하는 claim mode를 선택할 수 있습니다. 이 모드는 지정된 savepoint의 파일 소유권을 누가 취하는지 제어합니다.
CLI 액션
다음은 Flink의 CLI 도구가 지원하는 액션의 개요입니다.
| Action | Purpose |
|---|---|
run |
이 액션은 작업을 실행합니다. 작업을 포함하는 jar가 최소한 필요합니다. 필요하면 Flink- 또는 작업 관련 인수를 전달할 수 있습니다. |
info |
이 액션은 전달된 작업의 최적화된 실행 그래프를 출력하는 데 사용할 수 있습니다. 다시 작업을 포함하는 jar를 전달해야 합니다. |
list |
이 액션은 실행 중이거나 예약된 모든 작업을 나열합니다. |
savepoint |
이 액션은 주어진 작업에 대한 savepoint을 만들거나 폐기하는 데 사용할 수 있습니다. Flink configuration file에 execution.checkpointing.savepoint-dir 매개변수가 지정되지 않은 경우 JobID 외에도 savepoint 디렉터리를 지정해야 할 수 있습니다. |
checkpoint |
이 액션은 주어진 작업에 대한 checkpoint를 만드는 데 사용할 수 있습니다. checkpoint 타입(전체 또는 증분)을 지정할 수 있습니다. |
cancel |
이 액션은 JobID를 기반으로 실행 중인 작업을 취소하는 데 사용할 수 있습니다. |
stop |
이 액션은 cancel과 savepoint 액션을 결합하여 실행 중인 작업을 중지하지만 다시 시작할 savepoint도 만듭니다. |
모든 액션과 그 매개변수에 대한 더 세밀한 설명은 bin/flink --help 또는 각 개별 액션의 사용 정보 bin/flink <action> --help로 접근할 수 있습니다.
고급 CLI
REST API
Flink 클러스터는 REST API로도 관리할 수 있습니다. 이전 섹션에서 설명한 명령은 Flink의 REST 엔드포인트가 제공하는 것의 일부입니다. 따라서 curl 같은 도구를 사용해 Flink에서 더 많은 것을 얻을 수 있습니다.
배포 대상 선택
Flink는 Kubernetes나 YARN 같은 여러 클러스터 관리 프레임워크와 호환되며, 이는 Resource Provider 섹션에서 자세히 설명합니다. 작업은 서로 다른 배포 모드(Deployment Modes)로 제출될 수 있습니다. 작업 제출의 매개변수화는 하부 프레임워크와 배포 모드에 따라 다릅니다.
bin/flink는 서로 다른 옵션을 처리하는 --target 매개변수를 제공합니다. 또한 작업은 run을 사용해 제출해야 합니다(Session 및 Application Mode용). 다음 매개변수 조합 요약을 참조하세요.
- YARN
./bin/flink run --target yarn-session: 이미 실행 중인 Flink on YARN 클러스터에 제출./bin/flink run --target yarn-application: Application Mode로 YARN 클러스터에서 Flink를 구동하며 제출
- Kubernetes
./bin/flink run --target kubernetes-session: 이미 실행 중인 Flink on Kubernetes 클러스터에 제출./bin/flink run --target kubernetes-application: Application Mode로 Kubernetes 클러스터에서 Flink를 구동하며 제출
- Standalone:
./bin/flink run --target local: Session Mode에서 MiniCluster를 사용한 로컬 제출./bin/flink run --target remote: 이미 실행 중인 Flink 클러스터에 제출
--target은 Flink 구성 파일에 지정된 execution.target을 덮어씁니다.
명령과 사용 가능한 옵션에 대한 자세한 내용은 문서의 Resource Provider별 페이지를 참조하세요.
PyFlink 작업 제출
현재 사용자는 CLI를 통해 PyFlink 작업을 제출할 수 있습니다. Java 작업 제출과 달리 JAR 파일 경로나 엔트리 main 클래스를 지정할 필요가 없습니다.
flink run으로 Python 작업을 제출할 때 Flink는 "python" 명령을 실행합니다. 현재 환경의 python 실행 파일이 지원되는 Python 버전 3.9+를 가리키는지 확인하려면 다음 명령을 실행하세요.
$ python --version
# the version printed here must be 3.9+
다음 명령은 다양한 PyFlink 작업 제출 사용 사례를 보여줍니다.
- PyFlink 작업 실행:
$ ./bin/flink run --python examples/python/table/word_count.py
- 추가 소스 및 리소스 파일과 함께 PyFlink 작업 실행.
--pyFiles에 지정된 파일은PYTHONPATH에 추가되므로 Python 코드에서 사용할 수 있습니다.
$ ./bin/flink run \
--python examples/python/table/word_count.py \
--pyFiles file:///user.txt,hdfs:///$namenode_address/username.txt
- Java UDF 또는 외부 커넥터를 참조할 PyFlink 작업 실행.
--jarfile에 지정된 JAR 파일은 클러스터에 업로드됩니다.
$ ./bin/flink run \
--python examples/python/table/word_count.py \
--jarfile <jarFile>
--pyModule에 지정된 메인 엔트리 모듈과 pyFiles로 PyFlink 작업 실행:
$ ./bin/flink run \
--pyModule word_count \
--pyFiles examples/python/table
- 호스트
<jobmanagerHost>에서 실행 중인 특정 JobManager에 PyFlink 작업 제출(명령을 그에 맞게 조정):
$ ./bin/flink run \
--jobmanager <jobmanagerHost>:8081 \
--python examples/python/table/word_count.py
- Application Mode의 YARN 클러스터를 사용해 PyFlink 작업 실행:
$ ./bin/flink run -t yarn-application \
-Djobmanager.memory.process.size=1024m \
-Dtaskmanager.memory.process.size=1024m \
-Dyarn.application.name=<ApplicationName> \
-Dyarn.ship-files=/path/to/shipfiles \
-pyarch shipfiles/venv.zip \
-pyclientexec venv.zip/venv/bin/python3 \
-pyexec venv.zip/venv/bin/python3 \
-pyfs shipfiles \
-pym word_count
참고: 작업 실행에 필요한 Python 의존성이 이미 /path/to/shipfiles 디렉터리에 있다고 가정합니다. 예를 들어 위 예시에서는 venv.zip과 word_count.py를 포함해야 합니다.
참고: YARN application mode에서 JobManager에서 작업을 실행하므로 -pyarch와 -pyfs에 지정된 경로는 배송된 파일의 디렉터리 이름인 shipfiles에 대한 상대 경로입니다. 엔트리포인트의 절대 경로나 상대 경로를 알 수 없으므로 -py 대신 -pym을 사용해 프로그램 엔트리포인트를 지정하는 것이 좋습니다.
참고: -pyarch로 지정된 아카이브 파일은 blob server를 통해 TaskManager에 배포되며 파일 크기 제한은 2 GB입니다. 아카이브 파일 크기가 2 GB보다 크면 분산 파일 시스템에 업로드한 후 -pyarch 커맨드라인 옵션에서 그 경로를 사용할 수 있습니다.
- 클러스터 ID
<ClusterId>가 있는 네이티브 Kubernetes 클러스터에서 PyFlink 애플리케이션 실행. PyFlink가 설치된 docker 이미지가 필요합니다. Enabling PyFlink in docker 참조:
$ ./bin/flink run \
--target kubernetes-application \
--parallelism 8 \
-Dkubernetes.cluster-id=<ClusterId> \
-Dtaskmanager.memory.process.size=4096m \
-Dkubernetes.taskmanager.cpu=2 \
-Dtaskmanager.numberOfTaskSlots=4 \
-Dkubernetes.container.image.ref=<PyFlinkImageName> \
--pyModule word_count \
--pyFiles /opt/flink/examples/python/table/word_count.py
사용 가능한 옵션에 대한 자세한 내용은 Resource Provider 섹션에서 더 자세히 설명하는 Kubernetes 또는 YARN을 참조하세요.
위에서 언급한 --pyFiles, --pyModule, --python 외에도 다른 Python 관련 옵션이 있습니다. Flink의 CLI 도구가 지원하는 run 액션의 모든 Python 관련 옵션 개요는 다음과 같습니다.
| Option | Description |
|---|---|
-py,--python |
프로그램 엔트리 포인트가 있는 Python 스크립트입니다. 의존 리소스는 --pyFiles 옵션으로 구성할 수 있습니다. |
-pym,--pyModule |
프로그램 엔트리 포인트가 있는 Python 모듈입니다. 이 옵션은 --pyFiles와 함께 사용해야 합니다. |
-pyfs,--pyFiles |
작업에 사용자 지정 파일을 첨부합니다. .py/.egg/.zip/.whl 같은 표준 리소스 파일 접미사 또는 디렉터리가 모두 지원됩니다. 이 파일들은 로컬 클라이언트와 원격 python UDF worker 모두의 PYTHONPATH에 추가됩니다. .zip으로 끝나는 파일은 추출되어 PYTHONPATH에 추가됩니다. 여러 파일을 지정하려면 컴마(',')를 구분자로 사용할 수 있습니다(예: --pyFiles file:///tmp/myresource.zip,hdfs:///$namenode_address/myresource2.zip). |
-pyarch,--pyArchives |
작업에 python 아카이브 파일을 추가합니다. 아카이브 파일은 python UDF worker의 작업 디렉터리에 추출됩니다. 각 아카이브 파일에 대해 대상 디렉터리를 지정할 수 있습니다. 대상 디렉터리 이름이 지정되면 아카이브 파일은 지정된 이름의 디렉터리에 추출됩니다. 그렇지 않으면 아카이브 파일은 아카이브 파일과 같은 이름의 디렉터리에 추출됩니다. 이 옵션으로 업로드된 파일은 상대 경로로 접근할 수 있습니다. '#'을 아카이브 파일 경로와 대상 디렉터리 이름의 구분자로 사용할 수 있습니다. 여러 아카이브 파일을 지정하려면 컴마(',')를 구분자로 사용할 수 있습니다. 이 옵션은 가상 환경, Python UDF에서 사용되는 데이터 파일을 업로드하는 데 사용할 수 있습니다(예: --pyArchives file:///tmp/py37.zip,file:///tmp/data.zip#data --pyExecutable py37.zip/py37/bin/python). 데이터 파일은 Python UDF에서 접근할 수 있습니다. 예: f = open('data/data.txt', 'r'). |
-pyclientexec,--pyClientExecutable |
"flink run"으로 Python 작업을 제출하거나 Python UDF를 포함하는 Java 작업을 컴파일할 때 Python 프로세스를 시작하는 데 사용되는 Python 인터프리터의 경로입니다. (예: --pyArchives file:///tmp/py37.zip --pyClientExecutable py37.zip/py37/python) |
-pyexec,--pyExecutable |
python UDF worker를 실행하는 데 사용되는 python 인터프리터의 경로를 지정합니다(예: --pyExecutable /usr/local/bin/python3). python UDF worker는 Python 3.9+, Apache Beam(버전 >= 2.54.0, <= 2.61.0), Pip(버전 >= 20.3), SetupTools(버전 >= 37.0.0)에 의존합니다. 지정된 환경이 위 요구사항을 충족하는지 확인하세요. |
-pyreq,--pyRequirements |
타사 의존성을 정의하는 requirements.txt 파일을 지정합니다. 이 의존성들은 설치되어 python UDF worker의 PYTHONPATH에 추가됩니다. 이러한 의존성의 설치 패키지를 포함하는 디렉터리를 선택적으로 지정할 수 있습니다. 선택 매개변수가 있으면 '#'을 구분자로 사용합니다(예: --pyRequirements file:///tmp/requirements.txt#file:///tmp/cached_dir). |
작업 제출 중 커맨드라인 옵션 외에도 코드 내부의 구성 또는 Python API로 의존성을 지정하는 것도 지원합니다. 자세한 내용은 dependency management를 참조하세요.