前置配置

NFS 起到缓存的作用

  • 这个可以考虑去掉,其实可以直接从minio拉取,只是要每次多拉取一个jar.
(base) [root@hadoop04 flink-cluster]# cat /etc/exports
/flink-cluster/minio_data *(rw,sync,no_subtree_check,no_root_squash)
/flink-cluster/log *(rw,sync,no_subtree_check,no_root_squash)
/flink-cluster/artifacts *(rw,sync,no_subtree_check,no_root_squash)
/flink-cluster/model_config *(rw,sync,no_subtree_check,no_root_squash)

# 上面的目录要先创建好,然后log目录要给权限不然日志写不进来 chmod 777 -R /flink-cluster/log
systemctl start nfs-server
exportfs -ra

PV和PVC要提前创建好

job_pv.yaml

# job_pv.yaml
apiVersion: v1
kind: PersistentVolume
metadata:
  name: flink-log
  labels: 
    app: flink
spec:
  capacity:
    storage: 5Gi
  accessModes:
  - ReadWriteMany
  persistentVolumeReclaimPolicy: Retain # 设置为 Retain
  nfs:
    path: /flink-cluster/log
    server: 100.100.30.233  # 替换为你的 IP
---
apiVersion: v1
kind: PersistentVolume
metadata:
  name: flink-artifacts
  labels: 
    app: flink
spec:
  capacity:
    storage: 5Gi
  accessModes:
  - ReadWriteMany
  persistentVolumeReclaimPolicy: Retain # 设置为 Retain
  nfs:
    path: /flink-cluster/artifacts
    server: 100.100.30.233  # 替换为你的 IP
---
apiVersion: v1
kind: PersistentVolume
metadata:
  name: flink-model-config
  labels: 
    app: flink
spec:
  capacity:
    storage: 5Gi
  accessModes:
  - ReadWriteMany
  persistentVolumeReclaimPolicy: Retain # 设置为 Retain
  nfs:
    path: /flink-cluster/model_config
    server: 100.100.30.233  # 替换为你的 IP

job_pvc.yaml

apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: flink-log  # PVC 名称,自定义
spec:
  volumeName: flink-log # PV名称 
  accessModes:
  - ReadWriteMany  # 与 PV 的 RWX 匹配
  resources:
    requests:
      storage: 5Gi  # 与 PV 的容量匹配
---
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: flink-artifacts  # PVC 名称,自定义
spec:
  volumeName: flink-artifacts # PV名称
  accessModes:
  - ReadWriteMany  # 与 PV 的 RWX 匹配
  resources:
    requests:
      storage: 5Gi  # 与 PV 的容量匹配
---
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: flink-model-config  # PVC 名称,自定义
spec:
  volumeName: flink-model-config # PV名称
  accessModes:
  - ReadWriteMany  # 与 PV 的 RWX 匹配
  resources:
    requests:
      storage: 5Gi  # 与 PV 的容量匹配
  • minIO,后期可能会引入redis

去看声明式配置里面的
账号是: minioadmin minioadmin
预先准备好目录
	model-config
		存放任务运行配置
	artifacts
		存放运行
  • 准备flink权限账号, 这个也是要预先执行的

# flink_before.yaml 
apiVersion: v1
kind: ServiceAccount
metadata:
  name: flink
  namespace: default
---

apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: flink-role
  namespace: default
rules:
  - apiGroups: [""]
    resources: ["pods", "pods/log"]
    verbs: ["get", "list", "watch", "create", "delete"]
  - apiGroups: [""]
    resources: ["services", "configmaps"]
    verbs: ["get", "list", "watch", "create", "delete"]
  - apiGroups: ["apps"]
    resources: ["deployments"]
    verbs: ["get", "list", "watch", "create", "delete"]  # 添加对 deployments 的权限
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: flink-role-binding
  namespace: default
subjects:
  - kind: ServiceAccount
    name: flink
    namespace: default
roleRef:
  kind: Role
  name: flink-role
  apiGroup: rbac.authorization.k8s.io
  • 准备cluster-pod模版,这个模版和Dockerfile在同级目录,要打包到镜像里面去的

# template.yaml 
apiVersion: v1
kind: Pod
metadata:
  labels:
    app: flink
spec:
  serviceAccountName: flink
  volumes:
    - name: flink-cache
      persistentVolumeClaim:
        claimName: flink-main-container
  initContainers:
    - name: init-downloader
      image: minio/mc
      imagePullPolicy: IfNotPresent
      command: ["/bin/sh", "-c"]
      args:
        - |
          set -e -x
          echo "==> 设置 MinIO alias"
          mc alias set myminio http://minio-service.default.svc.cluster.local:9000 minioadmin minioadmin
          mkdir -p /opt/flink/defined/model-config
          if [ ! -f /opt/flink/defined/flinkOffline-jar-with-dependencies.jar ]; then
            echo "==> 下载 JAR 包"
            mc cp myminio/artifacts/flinkOffline-jar-with-dependencies.jar /opt/flink/defined/ || { echo "Failed to download JAR"; exit 1; }
            ls -alF /opt/flink/defined/
          fi
          if [ ! -f /opt/flink/defined/model-config/test.json ]; then
            echo "==> 下载 JSON 配置"
            mc cp myminio/model-config/test.json /opt/flink/defined/model-config/ || { echo "Failed to download JSON"; exit 1; }
            ls -alF /opt/flink/defined/model-config/
          fi
      volumeMounts:
        - name: flink-cache
          mountPath: /opt/flink/defined
  containers:
    - name: flink-main-container
      volumeMounts:
        - name: flink-cache
          mountPath: /opt/flink/defined
      securityContext:
        runAsUser: 9999
        fsGroup: 9999
  • 准备镜像

[root@hadoop04 empty]# cat Dockerfile 

# 使用官方 Flink 镜像作为基础镜像
FROM flink:1.20.0
RUN mkdir -p /opt/flink/template
RUN mkdir -p /opt/flink/defined
# 将 JAR 文件复制到容器内的某个目录,例如 /opt/flink/usrlib/
COPY template.yaml /opt/flink/template/pod-template.yaml
# 可选:设置工作目录
WORKDIR /opt/flink

# --------------------------下面是命令-------------

docker build -t flink_offline:1.0-with-template .
docker save -o flink.tar flink_offline:1.0-with-template
ctr -n k8s.io images rm flink_offline:1.0-with-template
ctr -n k8s.io images rm 
# 镜像要分发到所有节点     scp ./flink.tar hadoop03:/root

模版配置要提前创建好

podTemplateConfigMap.yaml

  • kubectl apply -f podTemplateConfigMap.yaml
apiVersion: v1
kind: ConfigMap
metadata:
  name: flink-log-config-map
  namespace: default
  labels:
    app: flink
data:
  pod-template.yaml: |
    apiVersion: v1
    kind: Pod
    metadata:
      labels:
        app: flink
    spec:
      serviceAccountName: flink
      volumes:
        - name: flink-model-config
          persistentVolumeClaim:
            claimName: flink-model-config
        - name: flink-log
          persistentVolumeClaim:
            claimName: flink-log
        - name: flink-artifacts
          persistentVolumeClaim:
            claimName: flink-artifacts
        - name: log-config
          configMap:
            name: flink-log-config-map
        - name: log-upload-status               # ✅ 新增共享 volume
          emptyDir: {}

      initContainers:
        - name: init-downloader
          image: minio/mc
          imagePullPolicy: IfNotPresent
          command: ["/bin/sh", "-c"]
          args:
            - |
              set -e
              echo "==> 设置 MinIO alias"
              mc alias set myminio http://minio-service.default.svc.cluster.local:9000 minioadmin minioadmin
              echo "==> 下载 Jar 包"
              if [ ! -f /opt/flink/defined/artifacts/flinkOffline-jar-with-dependencies.jar ]; then
                mc cp myminio/artifacts/flinkOffline-jar-with-dependencies.jar /opt/flink/defined/artifacts/ || echo "下载失败"
              else
                echo "==> Jar包已存在,跳过下载"
              fi
              LOG_FILE=""
              for f in /opt/flink/log/*"${HOSTNAME}"*; do
                if [ -f "$f" ]; then
                  LOG_FILE="$f"
                  break
                fi
              done
              if [ -n "$LOG_FILE" ]; then
                echo "==> 删除历史日志 $LOG_FILE"
                rm -rf $LOG_FILE
              else
                echo "==> 没有历史日志"
              fi
          volumeMounts:
            - name: flink-artifacts
              mountPath: /opt/flink/defined/artifacts
            - name: flink-log
              mountPath: /opt/flink/log

      containers:
        - name: flink-main-container
          volumeMounts:
            - name: flink-model-config
              mountPath: /opt/flink/defined/model-config
            - name: flink-log
              mountPath: /opt/flink/log
            - name: flink-artifacts
              mountPath: /opt/flink/defined/artifacts
            - name: log-config
              mountPath: /opt/flink/template/pod-template.yaml
              subPath: pod-template.yaml
            - name: log-upload-status                 # ✅ 挂载共享标志目录
              mountPath: /status
          lifecycle:                                  # ✅ 新增 preStop hook
            preStop:
              exec:
                command:
                  - /bin/sh
                  - -c
                  - |
                    echo "==> preStop: 通知 sidecar 上传日志"
                    touch /status/upload.flag
          securityContext:
            runAsUser: 9999
            fsGroup: 9999

        - name: minio-upload-sidecar
          image: minio/mc
          imagePullPolicy: IfNotPresent
          env:
            - name: HOSTNAME
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
          command: [ "/bin/sh", "-c" ]
          args:
            - |
              set -e
              echo "==> 设置 MinIO alias"
              mc alias set myminio http://minio-service.default.svc.cluster.local:9000 minioadmin minioadmin

              upload_log() {
                LOG_FILE=""
                for f in /opt/flink/log/*"${HOSTNAME}"*; do
                  if [ -f "$f" ]; then
                    LOG_FILE="$f"
                    break
                  fi
                done
                if [ -n "$LOG_FILE" ]; then
                  echo "==> 上传日志文件: $LOG_FILE"
                  FILE_NAME="${LOG_FILE##*/}"
                  mc cp "$LOG_FILE" myminio/model-log/$FILE_NAME || {
                    echo "==> 上传失败"
                    return 1
                  }
                  echo "==> 上传成功"
                else
                  echo "==> 没有日志文件"
                fi
              }

              while true; do
                if [ -f /status/upload.flag ]; then
                  echo "==> 检测到主容器退出信号"
                  upload_log
                  echo "==> 上传完成,sidecar 退出"
                  exit 0
                fi
                sleep 1
              done
          volumeMounts:
            - name: flink-log
              mountPath: /opt/flink/log
            - name: log-upload-status                # ✅ 挂载共享目录以监听退出信号
              mountPath: /status

flink账号,要提前创建好

flinkAccount.yaml

  • kubectl apply -f flinkAccount.yaml
apiVersion: v1
kind: ServiceAccount
metadata:
  name: flink
  namespace: default
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: flink-role
  namespace: default
rules:
  - apiGroups: [""]
    resources: ["pods", "pods/log"]
    verbs: ["get", "list", "watch", "create", "delete"]
  - apiGroups: [""]
    resources: ["services", "configmaps"]
    verbs: ["get", "list", "watch", "create", "delete"]
  - apiGroups: ["apps"]
    resources: ["deployments"]
    verbs: ["get", "list", "watch", "create", "delete"]  # 添加对 deployments 的权限
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: flink-role-binding
  namespace: default
subjects:
  - kind: ServiceAccount
    name: flink
    namespace: default
roleRef:
  kind: Role
  name: flink-role
  apiGroup: rbac.authorization.k8s.io

镜像问题

Dockerfile

  • docker build -t flink_offline:1.0-with-template .
  • docker save -o flink.tar flink_offline:1.0-with-template
  • ctr -n k8s.io images import flink.tar
  • # 如果要重新导入镜像就先删掉旧的 ctr -n k8s.io images rm flink_offline:1.0-with-template
# 使用官方 Flink 镜像作为基础镜像
FROM flink:1.20.0


RUN mkdir -p /opt/flink/template /opt/flink/defined/artifacts /opt/flink/defined/model-config 
# 模版配置可以预先打包到镜像里面去的,如果已经固定的话
# 可选:设置工作目录
WORKDIR /opt/flink

直接执行的

job.yaml

  • 确定好minio准备好且上面有jar包
  • 确定configMap和PVC好了,log目录好了
  • 确定好flink账号创建好了
  • 下面这个${JOB_ID}根据生产环境修改. JOB_ID长度不能大于31位
apiVersion: batch/v1
kind: Job
metadata:
  name: flink-job-${JOB_ID}
spec:
  backoffLimit: 0
  template:
    metadata:
      labels:
        job-id: ${JOB_ID}
    spec:
      serviceAccountName: flink
      restartPolicy: Never
      terminationGracePeriodSeconds: 30
      volumes:
        - name: flink-model-config
          persistentVolumeClaim:
            claimName: flink-model-config
        - name: flink-log
          persistentVolumeClaim:
            claimName: flink-log
        - name: flink-artifacts
          persistentVolumeClaim:
            claimName: flink-artifacts
        - name: log-config
          configMap:
            name: flink-log-config-map
      initContainers:
        - name: init-downloader
          image: minio/mc
          imagePullPolicy: IfNotPresent
          command: ["/bin/sh", "-c"]
          args:
            - |
              set -e
              if [ ! -f /opt/flink/defined/model-config/${JOB_ID}.json ]; then
                echo "==> 删除历史配置 ${JOB_ID}.json"
                rm -rf /opt/flink/defined/model-config/${JOB_ID}.json
              fi
              echo "==> 设置 MinIO alias"
              mc alias set myminio http://minio-service.default.svc.cluster.local:9000 minioadmin minioadmin
              echo "==> 下载配置"
              mc cp myminio/model-config/${JOB_ID}.json /opt/flink/defined/model-config/${JOB_ID}.json || { echo "下载配置失败"; exit 1; }
              ls -alF /opt/flink/defined/model-config/
          volumeMounts:
            - name: flink-model-config
              mountPath: /opt/flink/defined/model-config
      containers:
        - name: flink-job
          volumeMounts:
            - name: flink-model-config
              mountPath: /opt/flink/defined/model-config
            - name: flink-log
              mountPath: /opt/flink/log
            - name: flink-artifacts
              mountPath: /opt/flink/defined/artifacts
            - name: log-config
              mountPath: /opt/flink/template/pod-template.yaml
              subPath: pod-template.yaml
          image: docker.io/library/flink_offline:1.0-with-template
          imagePullPolicy: IfNotPresent
          command:
            - /opt/flink/bin/flink
            - run-application
            - -t
            - kubernetes-application
            - -Drestart-strategy.type=none
            - -Dkubernetes.cluster-id=fcluster-${JOB_ID}
            - -Dkubernetes.jobmanager.cpu=${POD_CPU_LIMIT} # 配置 jobmanager和taskManager cpu,这个不能太少不然是真的慢 示例: 0.5 相当于半个CPU
            - -Dtaskmanager.memory.process.size=${POD_MEMORY_LIMIT}mb       # 配置内存 进程大小 示例 500
            - -Dtaskmanager.numberOfTaskSlots=${SLOT_SIZE}
            - -Dkubernetes.container.image=docker.io/library/flink_offline:1.0-with-template
            - -Dkubernetes.pod-template-file=/opt/flink/template/pod-template.yaml
            - -c
            - org.kunmtech.app.ModelApplication
            - /opt/flink/defined/artifacts/flinkOffline-jar-with-dependencies.jar
            - -localConfig
            - /opt/flink/defined/model-config/${JOB_ID}.json
          resources:
            requests:
              cpu: "${JOB_CPU_LIMIT}m" #  这个限制的是job 示例 500 相当于半个CPU
              memory: "${JOB_MEMORY_LIMIT}Mi"
            limits:
              cpu: "${JOB_CPU_LIMIT}m"
              memory: "${JOB_MEMORY_LIMIT}Mi"

补充下Minio的配置

  • 有一点要注意系统时间必须校准,不然会出现客户端访问MinIO没有权限之类的报错
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: minio-pvc
spec:
  accessModes:
    - ReadWriteMany
  resources:
    requests:
      storage: 5Gi
---
apiVersion: v1
kind: Pod
metadata:
  name: minio
  labels:
    app: minio
spec:
  containers:
    - name: minio
      image: quay.io/minio/minio:latest
      command: ["/bin/sh", "-c"]
      args:
        - |
          export MINIO_SERVER_URL="http://${HOST_IP}:30900" && \
          exec minio server /data --console-address ":9001"
      env:
        - name: MINIO_ROOT_USER
          value: "minioadmin"
        - name: MINIO_ROOT_PASSWORD
          value: "minioadmin"
        - name: HOST_IP
          valueFrom:
            fieldRef:
              fieldPath: status.hostIP
      ports:
        - containerPort: 9000  # API 接口
        - containerPort: 9001  # 新增 WebUI 容器端口
	  resources:
		requests:
		  memory: "2048Mi"
		  cpu: "1000m"
		limits:
		  memory: "4Gi"
		  cpu: "4"
      volumeMounts:
        - name: minio-storage
          mountPath: /data
  volumes:
    - name: minio-storage
      persistentVolumeClaim:
        claimName: minio-pvc


---
apiVersion: v1
kind: PersistentVolume
metadata:
  name: minio-pv
spec:
  capacity:
    storage: 5Gi
  accessModes:
  - ReadWriteMany
  persistentVolumeReclaimPolicy: Retain # 设置为 Retain
  nfs:
    path: /flink-cluster/minio_data
    server: 100.100.30.233  # 替换为你的 IP
---
apiVersion: v1
kind: Service
metadata:
  name: minio-service
spec:
  type: NodePort
  ports:
    - name: api
      port: 9000
      targetPort: 9000
      nodePort: 30900  # 可从宿主机通过 30900 端口访问
    - name: webui
      port: 9001
      targetPort: 9001
      nodePort: 30901  # 新增:Web UI 的 NodePort
  selector:
    app: minio

官方文档

https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-main/docs/operations/metrics-logging/
https://nightlies.apache.org/flink/flink-docs-master/docs/deployment/resource-providers/native_kubernetes/#logging
要自定义配置的看这个文档
配置 | Apache Flink — Configuration | Apache Flink

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐