flink以Application模式部署在k8s集群
·
前置配置
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 的容量匹配
去看声明式配置里面的
账号是: minioadmin minioadmin
预先准备好目录
model-config
存放任务运行配置
artifacts
存放运行
# 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
# 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-templatectr -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
更多推荐




所有评论(0)