쿠버네티스의 스탠드얼론 배포
쿠버네티스의 스탠드얼론 배포 (Standalone Kubernetes)
이 페이지에서는 Flink의 standalone 배포 방식을 사용하여 [Kubernetes] 위에 Session cluster를 배포하는 방법을 설명합니다. 신규 사용자에게는 [native Kubernetes deployments]를 사용해 Flink를 배포하는 것을 일반적으로 권장합니다. Apache Flink는 쿠버네티스에서 Flink 클러스터를 관리하기 위한 Kubernetes operator도 제공합니다.
출처: 문서
본문
Getting Started (시작하기)
이 Getting Started 가이드는 [Kubernetes] 위에 Session cluster를 배포하는 방법을 설명합니다.
Introduction (소개)
이 페이지는 Flink의 standalone 배포를 사용하여 Kubernetes 위에 [standalone] Flink 클러스터를 배포하는 방법을 설명합니다. 신규 사용자에게는 [native Kubernetes deployments]를 사용해 쿠버네티스에 Flink를 배포할 것을 일반적으로 권장합니다.
Apache Flink는 쿠버네티스에서 Flink 클러스터를 관리하기 위한 Kubernetes operator도 제공합니다. 이 operator는 standalone 및 native 배포 모드를 모두 지원하며, 쿠버네티스에서 Flink 리소스의 배포, 구성, 생명주기 관리를 크게 단순화합니다.
자세한 내용은 [Flink Kubernetes Operator documentation]을 참조하세요.
Preparation (준비 사항)
이 가이드는 쿠버네티스 환경이 존재한다고 가정합니다. kubectl get nodes와 같은 명령을 실행하여 쿠버네티스 설정이 제대로 작동하는지 확인할 수 있으며, 이 명령은 연결된 모든 Kubelet을 나열합니다.
로컬에서 쿠버네티스를 실행하려면 [MiniKube]를 사용할 것을 권장합니다.
MiniKube를 사용하는 경우 Flink 클러스터를 배포하기 전에 minikube ssh 'sudo ip link set docker0 promisc on'을 실행해야 합니다. 그렇지 않으면 Flink 컴포넌트가 쿠버네티스 service를 통해 스스로를 참조할 수 없습니다.
Starting a Kubernetes Cluster (Session Mode)
Flink Session cluster는 장기 실행(life-long running) 쿠버네티스 Deployment로 실행됩니다. Session cluster에서는 여러 Flink job을 실행할 수 있습니다. 각 job은 클러스터가 배포된 후 클러스터에 제출되어야 합니다.
쿠버네티스에서 Flink Session cluster 배포는 최소한 세 가지 컴포넌트로 구성됩니다:
- [JobManager]를 실행하는 Deployment
- [TaskManagers] 풀을 위한 Deployment
- JobManager's REST 및 UI 포트를 노출하는 Service
[the common resource definitions]에서 제공하는 파일 내용을 사용하여 다음 파일들을 만들고, kubectl 명령으로 각각의 컴포넌트를 생성합니다:
# Configuration and service definition
$ kubectl create -f flink-configuration-configmap.yaml
$ kubectl create -f jobmanager-service.yaml
# Create the deployments for the cluster
$ kubectl create -f jobmanager-session-deployment-non-ha.yaml
$ kubectl create -f taskmanager-session-deployment.yaml
다음으로, Flink UI에 접근하고 job을 제출하기 위해 포트 포워딩을 설정합니다:
kubectl port-forward ${flink-jobmanager-pod} 8081:8081을 실행하여 jobmanager의 web ui 포트를 로컬 8081로 포워딩합니다.- 브라우저에서 [http://localhost:8081]로 이동합니다.
- 또한 아래 명령을 사용하여 클러스터에 job을 제출할 수도 있습니다:
$ ./bin/flink run -m localhost:8081 ./examples/streaming/TopSpeedWindowing.jar
다음 명령으로 클러스터를 내려(tear down) 수 있습니다:
$ kubectl delete -f jobmanager-service.yaml
$ kubectl delete -f flink-configuration-configmap.yaml
$ kubectl delete -f taskmanager-session-deployment.yaml
$ kubectl delete -f jobmanager-session-deployment-non-ha.yaml
Deployment Modes (배포 모드)
Application Mode (애플리케이션 모드)
애플리케이션 모드의 전반적인 개념은 [deployment mode overview]를 참조하세요.
Flink Application cluster는 배포 시점에 제공되어야 하는 단일 애플리케이션을 실행하는 전용 클러스터입니다.
쿠버네티스에서 기본 Flink Application cluster 배포는 세 가지 컴포넌트로 구성됩니다:
- JobManager를 실행하는 Application
- TaskManagers 풀을 위한 Deployment
- JobManager's REST 및 UI 포트를 노출하는 Service
[the Application cluster specific resource definitions]를 확인하고 그에 맞게 조정합니다:
jobmanager-application-non-ha.yaml의 args 속성은 사용자 job의 메인 클래스를 지정해야 합니다. jobmanager-application-non-ha.yaml의 Flink 이미지에 다른 args를 전달하는 방법은 [how to specify the JobManager arguments]를 참조하세요.
job artifacts는 다음과 같은 방식으로 제공할 수 있습니다:
- job artifacts는 [the resource definition examples]의
job-artifacts-volume에서 사용할 수 있습니다. 정의 예제는 minikube 클러스터에서 컴포넌트를 생성한다고 가정하고 볼륨을 호스트의 로컬 디렉터리로 마운트합니다. minikube 클러스터를 사용하지 않는다면 쿠버네티스 클러스터에서 사용 가능한 다른 유형의 볼륨을 사용하여 job artifacts를 제공할 수 있습니다. - 대신 아티팩트를 이미 포함하는 [a custom image]를 빌드할 수 있습니다.
- 로컬, [remote DFS], 또는 HTTP(S) 엔드포인트를 통해 접근할 수 있는 [–jars] 옵션으로 아티팩트를 전달할 수 있습니다.
[the common cluster components]를 만든 후 [the Application cluster specific resource definitions]를 사용하여 kubectl 명령으로 클러스터를 실행합니다:
$ kubectl create -f jobmanager-application-non-ha.yaml
$ kubectl create -f taskmanager-job-deployment.yaml
단일 애플리케이션 클러스터를 종료하려면 [the common ones]와 함께 이 컴포넌트들을 kubectl 명령으로 삭제할 수 있습니다:
$ kubectl delete -f taskmanager-job-deployment.yaml
$ kubectl delete -f jobmanager-application-non-ha.yaml
Session Mode (세션 모드)
세션 모드의 전반적인 개념은 [deployment mode overview]를 참조하세요.
Session 클러스터의 배포는 이 페이지 상단의 [Getting Started] 가이드에서 설명합니다.
Flink on Standalone Kubernetes Reference
Configuration
모든 구성 옵션은 [configuration page]에 나열되어 있습니다. 구성 옵션은 flink-configuration-configmap.yaml config map의 [Flink configuration file] 섹션에 추가할 수 있습니다.
Accessing Flink in Kubernetes
다양한 방법으로 Flink UI에 접근하고 job을 제출할 수 있습니다:
-
kubectl proxy: -
터미널에서
kubectl proxy를 실행합니다. -
브라우저에서 [http://localhost:8001/api/v1/namespaces/default/services/flink-jobmanager:webui/proxy]로 이동합니다.
-
kubectl port-forward: -
kubectl port-forward ${flink-jobmanager-pod} 8081:8081을 실행하여 jobmanager의 web ui 포트를 로컬 8081로 포워딩합니다. -
브라우저에서 [http://localhost:8081]로 이동합니다.
-
또한 아래 명령을 사용하여 클러스터에 job을 제출할 수 있습니다:
$ ./bin/flink run -m localhost:8081 ./examples/streaming/TopSpeedWindowing.jar
-
jobmanager의 rest service에
NodePortservice를 생성합니다: -
kubectl create -f jobmanager-rest-service.yaml을 실행하여 jobmanager에NodePortservice를 생성합니다.jobmanager-rest-service.yaml의 예제는 [appendix]에서 찾을 수 있습니다. -
kubectl get svc flink-jobmanager-rest를 실행하여 이 service의node-port를 확인하고 브라우저에서 [http://: ]로 이동합니다. -
minikube를 사용한다면
minikube ip를 실행하여 public ip를 얻을 수 있습니다. -
port-forward솔루션과 유사하게 아래 명령으로 클러스터에 job을 제출할 수도 있습니다:
$ ./bin/flink run -m <public-node-ip>:<node-port> ./examples/streaming/TopSpeedWindowing.jar
Debugging and Log Access
많은 일반적인 오류는 Flink의 로그 파일을 확인하여 쉽게 감지할 수 있습니다. Flink의 웹 사용자 인터페이스에 접근할 수 있다면 거기서 JobManager 및 TaskManager 로그를 볼 수 있습니다.
Flink를 시작할 때 문제가 있으면 쿠버네티스 유틸리티로 로그에 접근할 수도 있습니다. kubectl get pods를 사용하여 실행 중인 모든 pod를 확인합니다. 위의 quickstart 예제에서는 세 개의 pod가 보여야 합니다:
$ kubectl get pods
NAME READY STATUS RESTARTS AGE
flink-jobmanager-589967dcfc-m49xv 1/1 Running 3 3m32s
flink-taskmanager-64847444ff-7rdl4 1/1 Running 3 3m28s
flink-taskmanager-64847444ff-nnd6m 1/1 Running 3 3m28s
이제 kubectl logs flink-jobmanager-589967dcfc-m49xv를 실행하여 로그에 접근할 수 있습니다.
High-Availability with Standalone Kubernetes
쿠버네티스에서 고가용성을 위해 [existing high availability services]를 사용할 수 있습니다.
Kubernetes High-Availability Services
Session Mode 및 Application Mode 클러스터는 [Kubernetes high availability service] 사용을 지원합니다. [flink-configuration-configmap.yaml]에 다음 Flink 구성 옵션을 추가해야 합니다.
참고: 구성된 HA 스토리지 디렉터리의 스킴에 해당하는 파일시스템이 런타임에서 사용 가능해야 합니다. 자세한 내용은 [custom Flink image]와 [enable plugins]를 참조하세요.
apiVersion: v1
kind: ConfigMap
metadata:
name: flink-config
labels:
app: flink
data:
config.yaml: |+
...
kubernetes.cluster-id: <cluster-id>
high-availability.type: kubernetes
high-availability.storageDir: hdfs:///flink/recovery
restart-strategy.type: fixed-delay
restart-strategy.fixed-delay.attempts: 10
...
또한 ConfigMap을 create, edit, delete할 수 있는 권한이 있는 service account로 JobManager 및 TaskManager pod를 시작해야 합니다. 자세한 내용은 [how to configure service accounts for pods]를 참조하세요.
고가용성이 활성화되면 Flink는 서비스 디스커버리에 자체 HA 서비스를 사용합니다. 따라서 JobManager pod는 Kubernetes service 대신 IP 주소를 jobmanager.rpc.address로 사용하여 시작해야 합니다. 전체 구성은 [appendix]를 참조하세요.
Standby JobManagers
보통 단일 JobManager pod만 시작하면 충분합니다. pod가 충돌하면 쿠버네티스가 다시 시작하기 때문입니다. 더 빠른 복구를 원한다면 jobmanager-session-deployment-ha.yaml의 replicas 또는 jobmanager-application-ha.yaml의 parallelism을 1보다 큰 값으로 구성하여 대기(standby) JobManager를 시작하세요.
Using Standalone Kubernetes with Reactive Mode
[Reactive Mode]는 Application Cluster가 항상 job 병렬 처리도를 사용 가능한 리소스에 맞게 조정하는 모드로 Flink를 실행할 수 있게 해줍니다. 쿠버네티스와 결합하면 TaskManager deployment의 replica 수가 사용 가능한 리소스를 결정합니다. replica 수를 늘리면 job이 확장(scale up)되고, 줄이면 축소(scale down)가 촉발됩니다. 이는 [Horizontal Pod Autoscaler]를 사용해 자동으로 수행할 수도 있습니다.
쿠버네티스에서 Reactive Mode를 사용하려면 [deploying a job using an Application Cluster]와 동일한 단계를 따르세요. 단, flink-configuration-configmap.yaml 대신 이 config map을 사용하세요: flink-reactive-mode-configuration-configmap.yaml. 여기에는 Flink용 scheduler-mode: reactive 설정이 포함되어 있습니다.
Application Cluster를 배포한 후 flink-taskmanager deployment의 replica 수를 변경하여 job을 확장하거나 축소할 수 있습니다.
Enabling Local Recovery Across Pod Restarts
pod 실패 시 복구를 가속화하려면 Flink의 [working directory] 기능을 로컬 복구와 함께 활용할 수 있습니다. working directory가 재시작된 TaskManager pod에 다시 마운트되는 영구 볼륨에 있도록 구성되면, Flink는 상태를 로컬로 복구할 수 있습니다. [StatefulSet]을 사용하면 pod를 영구 볼륨에 매핑하는 데 필요한 정확한 도구를 쿠버네티스가 제공합니다.
TaskManager를 StatefulSet으로 배포하면 TaskManager에 영구 볼륨을 마운트하는 데 사용되는 볼륨 클레임 템플릿을 구성할 수 있습니다. 또한 결정적인 taskmanager.resource-id를 구성해야 합니다. 적합한 값은 환경 변수로 노출하는 [pod name]입니다. 예제 StatefulSet 구성은 [appendix]를 참조하세요.
Appendix (부록)
Common cluster resource definitions
flink-configuration-configmap.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: flink-config
labels:
app: flink
data:
config.yaml: |+
jobmanager.rpc.address: flink-jobmanager
taskmanager.numberOfTaskSlots: 2
blob.server.port: 6124
jobmanager.rpc.port: 6123
taskmanager.rpc.port: 6122
jobmanager.memory.process.size: 1600m
taskmanager.memory.process.size: 1728m
parallelism.default: 2
log4j-console.properties: |+
# This affects logging for both user code and Flink
rootLogger.level = INFO
rootLogger.appenderRef.console.ref = ConsoleAppender
rootLogger.appenderRef.rolling.ref = RollingFileAppender
# Uncomment this if you want to _only_ change Flink's logging
#logger.flink.name = org.apache.flink
#logger.flink.level = INFO
# The following lines keep the log level of common libraries/connectors on
# log level INFO. The root logger does not override this. You have to manually
# change the log levels here.
logger.pekko.name = org.apache.pekko
logger.pekko.level = INFO
logger.kafka.name= org.apache.kafka
logger.kafka.level = INFO
logger.hadoop.name = org.apache.hadoop
logger.hadoop.level = INFO
logger.zookeeper.name = org.apache.zookeeper
logger.zookeeper.level = INFO
# Log all infos to the console
appender.console.name = ConsoleAppender
appender.console.type = CONSOLE
appender.console.layout.type = PatternLayout
appender.console.layout.pattern = %d{yyyy-MM-dd HH:mm:ss,SSS} %-5p %-60c %x - %m%n
# Log all infos in the given rolling file
appender.rolling.name = RollingFileAppender
appender.rolling.type = RollingFile
appender.rolling.append = false
appender.rolling.fileName = ${sys:log.file}
appender.rolling.filePattern = ${sys:log.file}.%i
appender.rolling.layout.type = PatternLayout
appender.rolling.layout.pattern = %d{yyyy-MM-dd HH:mm:ss,SSS} %-5p %-60c %x - %m%n
appender.rolling.policies.type = Policies
appender.rolling.policies.size.type = SizeBasedTriggeringPolicy
appender.rolling.policies.size.size=100MB
appender.rolling.strategy.type = DefaultRolloverStrategy
appender.rolling.strategy.max = 10
# Suppress the irrelevant (wrong) warnings from the Netty channel handler
logger.netty.name = org.jboss.netty.channel.DefaultChannelPipeline
logger.netty.level = OFF
flink-reactive-mode-configuration-configmap.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: flink-config
labels:
app: flink
data:
config.yaml: |+
jobmanager.rpc.address: flink-jobmanager
taskmanager.numberOfTaskSlots: 2
blob.server.port: 6124
jobmanager.rpc.port: 6123
taskmanager.rpc.port: 6122
jobmanager.memory.process.size: 1600m
taskmanager.memory.process.size: 1728m
parallelism.default: 2
scheduler-mode: reactive
execution.checkpointing.interval: 10s
log4j-console.properties: |+
# This affects logging for both user code and Flink
rootLogger.level = INFO
rootLogger.appenderRef.console.ref = ConsoleAppender
rootLogger.appenderRef.rolling.ref = RollingFileAppender
# Uncomment this if you want to _only_ change Flink's logging
#logger.flink.name = org.apache.flink
#logger.flink.level = INFO
# The following lines keep the log level of common libraries/connectors on
# log level INFO. The root logger does not override this. You have to manually
# change the log levels here.
logger.pekko.name = org.apache.pekko
logger.pekko.level = INFO
logger.kafka.name= org.apache.kafka
logger.kafka.level = INFO
logger.hadoop.name = org.apache.hadoop
logger.hadoop.level = INFO
logger.zookeeper.name = org.apache.zookeeper
logger.zookeeper.level = INFO
# Log all infos to the console
appender.console.name = ConsoleAppender
appender.console.type = CONSOLE
appender.console.layout.type = PatternLayout
appender.console.layout.pattern = %d{yyyy-MM-dd HH:mm:ss,SSS} %-5p %-60c %x - %m%n
# Log all infos in the given rolling file
appender.rolling.name = RollingFileAppender
appender.rolling.type = RollingFile
appender.rolling.append = false
appender.rolling.fileName = ${sys:log.file}
appender.rolling.filePattern = ${sys:log.file}.%i
appender.rolling.layout.type = PatternLayout
appender.rolling.layout.pattern = %d{yyyy-MM-dd HH:mm:ss,SSS} %-5p %-60c %x - %m%n
appender.rolling.policies.type = Policies
appender.rolling.policies.size.type = SizeBasedTriggeringPolicy
appender.rolling.policies.size.size=100MB
appender.rolling.strategy.type = DefaultRolloverStrategy
appender.rolling.strategy.max = 10
# Suppress the irrelevant (wrong) warnings from the Netty channel handler
logger.netty.name = org.jboss.netty.channel.DefaultChannelPipeline
logger.netty.level = OFF
jobmanager-service.yaml - 비-HA 모드에만 필요한 선택적 service입니다.
apiVersion: v1
kind: Service
metadata:
name: flink-jobmanager
spec:
type: ClusterIP
ports:
- name: rpc
port: 6123
- name: blob-server
port: 6124
- name: webui
port: 8081
selector:
app: flink
component: jobmanager
jobmanager-rest-service.yaml. jobmanager의 rest 포트를 공개 쿠버네티스 노드 포트로 노출하는 선택적 service입니다.
apiVersion: v1
kind: Service
metadata:
name: flink-jobmanager-rest
spec:
type: NodePort
ports:
- name: rest
port: 8081
targetPort: 8081
nodePort: 30081
selector:
app: flink
component: jobmanager
Session cluster resource definitions
jobmanager-session-deployment-non-ha.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: flink-jobmanager
spec:
replicas: 1
selector:
matchLabels:
app: flink
component: jobmanager
template:
metadata:
labels:
app: flink
component: jobmanager
spec:
containers:
- name: jobmanager
image: apache/flink:2.3.0-scala_2.12
args: ["jobmanager"]
ports:
- containerPort: 6123
name: rpc
- containerPort: 6124
name: blob-server
- containerPort: 8081
name: webui
livenessProbe:
tcpSocket:
port: 6123
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf
securityContext:
runAsUser: 9999 # refers to user _flink_ from official flink image, change if necessary
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
jobmanager-session-deployment-ha.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: flink-jobmanager
spec:
replicas: 1 # Set the value to greater than 1 to start standby JobManagers
selector:
matchLabels:
app: flink
component: jobmanager
template:
metadata:
labels:
app: flink
component: jobmanager
spec:
containers:
- name: jobmanager
image: apache/flink:2.3.0-scala_2.12
env:
- name: POD_IP
valueFrom:
fieldRef:
apiVersion: v1
fieldPath: status.podIP
# The following args overwrite the value of jobmanager.rpc.address configured in the configuration config map to POD_IP.
args: ["jobmanager", "$(POD_IP)"]
ports:
- containerPort: 6123
name: rpc
- containerPort: 6124
name: blob-server
- containerPort: 8081
name: webui
livenessProbe:
tcpSocket:
port: 6123
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf
securityContext:
runAsUser: 9999 # refers to user _flink_ from official flink image, change if necessary
serviceAccountName: flink-service-account # Service account which has the permissions to create, edit, delete ConfigMaps
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
taskmanager-session-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: flink-taskmanager
spec:
replicas: 2
selector:
matchLabels:
app: flink
component: taskmanager
template:
metadata:
labels:
app: flink
component: taskmanager
spec:
containers:
- name: taskmanager
image: apache/flink:2.3.0-scala_2.12
args: ["taskmanager"]
ports:
- containerPort: 6122
name: rpc
livenessProbe:
tcpSocket:
port: 6122
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf/
securityContext:
runAsUser: 9999 # refers to user _flink_ from official flink image, change if necessary
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
Application cluster resource definitions
jobmanager-application-non-ha.yaml
apiVersion: batch/v1
kind: Job
metadata:
name: flink-jobmanager
spec:
template:
metadata:
labels:
app: flink
component: jobmanager
spec:
restartPolicy: OnFailure
containers:
- name: jobmanager
image: apache/flink:2.3.0-scala_2.12
env:
args: ["standalone-job", "--job-classname", "com.job.ClassName", <optional arguments>, <job arguments>] # optional arguments: ["--job-id", "<job id>", "--jars", "/path/to/artifact1,/path/to/artifact2", "--fromSavepoint", "/path/to/savepoint", "--allowNonRestoredState"]
ports:
- containerPort: 6123
name: rpc
- containerPort: 6124
name: blob-server
- containerPort: 8081
name: webui
livenessProbe:
tcpSocket:
port: 6123
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf
- name: job-artifacts-volume
mountPath: /opt/flink/usrlib
securityContext:
runAsUser: 9999 # refers to user _flink_ from official flink image, change if necessary
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
- name: job-artifacts-volume
hostPath:
path: /host/path/to/job/artifacts
jobmanager-application-ha.yaml
apiVersion: batch/v1
kind: Job
metadata:
name: flink-jobmanager
spec:
parallelism: 1 # Set the value to greater than 1 to start standby JobManagers
template:
metadata:
labels:
app: flink
component: jobmanager
spec:
restartPolicy: OnFailure
containers:
- name: jobmanager
image: apache/flink:2.3.0-scala_2.12
env:
- name: POD_IP
valueFrom:
fieldRef:
apiVersion: v1
fieldPath: status.podIP
# The following args overwrite the value of jobmanager.rpc.address configured in the configuration config map to POD_IP.
args: ["standalone-job", "--host", "$(POD_IP)", "--job-classname", "com.job.ClassName", <optional arguments>, <job arguments>] # optional arguments: ["--job-id", "<job id>", "--jars", "/path/to/artifact1,/path/to/artifact2", "--fromSavepoint", "/path/to/savepoint", "--allowNonRestoredState"]
ports:
- containerPort: 6123
name: rpc
- containerPort: 6124
name: blob-server
- containerPort: 8081
name: webui
livenessProbe:
tcpSocket:
port: 6123
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf
- name: job-artifacts-volume
mountPath: /opt/flink/usrlib
securityContext:
runAsUser: 9999 # refers to user _flink_ from official flink image, change if necessary
serviceAccountName: flink-service-account # Service account which has the permissions to create, edit, delete ConfigMaps
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
- name: job-artifacts-volume
hostPath:
path: /host/path/to/job/artifacts
taskmanager-job-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: flink-taskmanager
spec:
replicas: 2
selector:
matchLabels:
app: flink
component: taskmanager
template:
metadata:
labels:
app: flink
component: taskmanager
spec:
containers:
- name: taskmanager
image: apache/flink:2.3.0-scala_2.12
env:
args: ["taskmanager"]
ports:
- containerPort: 6122
name: rpc
livenessProbe:
tcpSocket:
port: 6122
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf/
- name: job-artifacts-volume
mountPath: /opt/flink/usrlib
securityContext:
runAsUser: 9999 # refers to user _flink_ from official flink image, change if necessary
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
- name: job-artifacts-volume
hostPath:
path: /host/path/to/job/artifacts
Local Recovery Enabled TaskManager StatefulSet
apiVersion: v1
kind: ConfigMap
metadata:
name: flink-config
labels:
app: flink
data:
config.yaml: |+
jobmanager.rpc.address: flink-jobmanager
taskmanager.numberOfTaskSlots: 2
blob.server.port: 6124
jobmanager.rpc.port: 6123
taskmanager.rpc.port: 6122
state.backend.local-recovery: true
process.taskmanager.working-dir: /pv
---
apiVersion: v1
kind: Service
metadata:
name: taskmanager-hl
spec:
clusterIP: None
selector:
app: flink
component: taskmanager
---
apiVersion: apps/v1
kind: StatefulSet
metadata:
name: flink-taskmanager
spec:
serviceName: taskmanager-hl
replicas: 2
selector:
matchLabels:
app: flink
component: taskmanager
template:
metadata:
labels:
app: flink
component: taskmanager
spec:
securityContext:
runAsUser: 9999
fsGroup: 9999
containers:
- name: taskmanager
image: apache/flink:2.3.0-scala_2.12
env:
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
args: ["taskmanager", "-Dtaskmanager.resource-id=$(POD_NAME)"]
ports:
- containerPort: 6122
name: rpc
- containerPort: 6121
name: metrics
livenessProbe:
tcpSocket:
port: 6122
initialDelaySeconds: 30
periodSeconds: 60
volumeMounts:
- name: flink-config-volume
mountPath: /opt/flink/conf/
- name: pv
mountPath: /pv
volumes:
- name: flink-config-volume
configMap:
name: flink-config
items:
- key: config.yaml
path: config.yaml
- key: log4j-console.properties
path: log4j-console.properties
volumeClaimTemplates:
- metadata:
name: pv
spec:
accessModes: [ "ReadWriteOnce" ]
resources:
requests:
storage: 50Gi