쿠버네티스의 스탠드얼론 배포

쿠버네티스의 스탠드얼론 배포 (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.yamlargs 속성은 사용자 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] 가이드에서 설명합니다.

Configuration

모든 구성 옵션은 [configuration page]에 나열되어 있습니다. 구성 옵션은 flink-configuration-configmap.yaml config map의 [Flink configuration file] 섹션에 추가할 수 있습니다.

Accessing Flink in Kubernetes

다양한 방법으로 Flink UI에 접근하고 job을 제출할 수 있습니다:

$ ./bin/flink run -m localhost:8081 ./examples/streaming/TopSpeedWindowing.jar
  • jobmanager의 rest service에 NodePort service를 생성합니다:

  • kubectl create -f jobmanager-rest-service.yaml을 실행하여 jobmanager에 NodePort service를 생성합니다. 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.yamlreplicas 또는 jobmanager-application-ha.yamlparallelism1보다 큰 값으로 구성하여 대기(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

더 알아보기 (Learn more)