Flink on Kubernetes:云原生部署实践全解析

引言:为什么Flink要跑在Kubernetes上?

在实时计算领域,Apache Flink凭借低延迟、高吞吐、精确一次(Exactly-Once)语义的特性,成为流处理的事实标准。但随着业务规模的增长,传统部署方式(如Standalone集群、YARN/Mesos调度)逐渐暴露出以下痛点:

  1. 资源利用率低:固定规格的集群无法应对流量波动,闲时资源闲置、忙时资源不足;
  2. 运维复杂度高:需要手动管理集群节点、调整资源配额,故障恢复依赖人工干预;
  3. 多租户隔离弱:共享集群中的作业容易互相抢占资源,无法保证SLA;
  4. 云原生适配差:无法利用公有云/私有云的弹性资源(如自动扩缩容、按需付费)。

而Kubernetes(以下简称K8s)作为云原生时代的“操作系统”,恰好能解决这些问题:

  • 弹性伸缩:根据负载自动调整Pod数量,避免资源浪费;
  • 资源隔离:通过Namespaces、ResourceQuotas实现多租户资源管控;
  • 自动化运维:通过StatefulSet/Deployment管理应用生命周期,故障自动恢复;
  • 生态兼容:无缝对接云存储(S3/GCS)、监控(Prometheus)、日志(Loki)等工具。

简言之,Flink on Kubernetes是流处理走向云原生的必然选择——它将Flink的流处理能力与K8s的资源管理能力深度融合,实现“按需分配、自动运维、弹性伸缩”的生产级流处理系统。

基础概念铺垫:Flink与K8s的核心对应关系

在开始部署实践前,我们需要先明确Flink集群架构与K8s核心资源的对应关系,这是理解后续内容的关键。

1. Flink集群架构回顾

Flink集群由JobManagerTaskManager两类组件构成:

  • JobManager(JM):集群的“大脑”,负责作业调度、Checkpoint协调、故障恢复;
  • TaskManager(TM):集群的“工人”,负责执行具体的任务(Task),每个TM有多个Task Slot(任务槽,用于分配并行度);
  • Client:提交作业的客户端,将作业Jar包和配置发送给JM。

2. K8s核心资源速览

K8s通过以下资源实现应用的部署与管理:

  • Pod:K8s的最小调度单元,包含一个或多个容器(如Flink的JM/TM容器);
  • StatefulSet:用于管理有状态应用(如JM),提供稳定的网络标识(固定主机名)和持久存储;
  • Deployment:用于管理无状态应用(如TM),支持滚动更新、自动扩缩容;
  • Service:为Pod提供稳定的网络访问入口(如JM的RPC端口、Web UI端口);
  • ConfigMap:存储应用配置(如Flink的flink-conf.yaml);
  • PersistentVolumeClaim(PVC):申请持久化存储(如JM的HA元数据存储)。

3. 对应关系总结

Flink组件K8s资源类型原因说明
JobManagerStatefulSet + ServiceJM需要稳定的网络地址(供TM注册)和持久存储(HA元数据)
TaskManagerDeploymentTM是无状态的,可动态扩缩容
Flink配置ConfigMap集中管理配置,避免硬编码
持久化存储PVC存储Checkpoint/Savepoint、HA元数据

Flink on K8s的三种部署模式:原理与场景

Flink官方提供三种K8s部署模式,分别对应不同的业务场景。我们需要根据作业类型、资源需求、运维复杂度选择合适的模式。

1. Session Mode(会话模式)

原理

Session Mode是共享集群模式:先启动一个长期运行的Flink集群(包含JM和一组TM),然后通过Client向集群提交多个作业。所有作业共享集群的资源(Task Slot)。

适用场景
  • 小规模作业(如测试、Demo);
  • 作业运行时间短、资源需求小;
  • 希望复用集群资源,降低运维成本。
优缺点
  • 优点:集群复用,启动快;
  • 缺点:资源隔离差(作业间抢占Slot)、故障影响范围大(JM挂掉会导致所有作业失败)。

2. Job Cluster Mode(作业集群模式)

原理

Job Cluster Mode是专属集群模式:为每个作业启动一个独立的Flink集群(JM + 专属TM)。作业完成后,集群自动销毁。

适用场景
  • 大规模生产作业(如实时数据同步、用户画像);
  • 作业资源需求大、需要严格隔离;
  • 希望作业失败不影响其他任务。
优缺点
  • 优点:资源隔离好、故障影响范围小;
  • 缺点:集群启动时间长(每个作业都要创建JM/TM)、资源利用率略低。

3. Application Mode(应用模式)

原理

Application Mode是云原生最优模式:将应用代码Flink集群打包成一个Docker镜像,直接在K8s上运行。作业的生命周期与集群一致(启动集群→运行作业→销毁集群)。

适用场景
  • 生产级流处理作业(如实时推荐、日志分析);
  • 希望避免Client端依赖(如Jar包版本冲突);
  • 追求“一键部署”的云原生体验。
优缺点
  • 优点:无Client依赖、镜像化部署、运维简单;
  • 缺点:需要构建自定义镜像,初期学习成本略高。

模式对比总结

维度Session ModeJob Cluster ModeApplication Mode
集群复用
资源隔离
启动速度
Client依赖需要需要不需要
生产适用性

部署实践:从0到1搭建Flink on K8s集群

接下来,我们将通过**本地K8s集群(Kind)**演示三种模式的部署过程。所有操作均基于Flink 1.17.0(最新稳定版)和K8s 1.27.0。

前置环境准备

1. 安装依赖工具
  • Kind:本地K8s集群搭建工具(替代Minikube,启动更快);
  • Docker:镜像构建与运行工具;
  • kubectl:K8s命令行工具;
  • Flink CLI:Flink命令行客户端(用于提交作业)。
2. 搭建本地K8s集群

使用Kind创建一个单节点K8s集群:

kind create cluster --name flink-demo

验证集群状态:

kubectl cluster-info --context kind-flink-demo

模式1:Session Mode部署

步骤1:创建Flink配置ConfigMap

首先,我们需要将Flink的核心配置flink-conf.yaml存储为ConfigMap,方便集群共享。

创建flink-config.yaml文件:

apiVersion: v1
kind: ConfigMap
metadata:
  name: flink-config
  labels:
    app: flink
data:
  flink-conf.yaml: |+
    # JobManager配置
    jobmanager.rpc.address: flink-jobmanager  # JM的Service名称
    jobmanager.memory.process.size: 1024m     # JM进程总内存
    # TaskManager配置
    taskmanager.memory.process.size: 2048m    # TM进程总内存
    taskmanager.numberOfTaskSlots: 2          # 每个TM的Task Slot数量
    # 网络配置
    blob.server.port: 6124
    jobmanager.rpc.port: 6123
    taskmanager.rpc.port: 6122
    # HA配置(使用K8s的ConfigMap存储HA元数据,简化部署)
    high-availability: org.apache.flink.kubernetes.highavailability.KubernetesHaServicesFactory
    high-availability.storageDir: hdfs:///flink/ha  # 可替换为S3/GCS路径
    # Metrics配置(用于Prometheus监控)
    metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter
    metrics.reporter.prom.port: 9250-9260
  log4j-console.properties: |+
    log4j.rootLogger=INFO, console
    log4j.appender.console=org.apache.log4j.ConsoleAppender
    log4j.appender.console.layout=org.apache.log4j.PatternLayout
    log4j.appender.console.layout.ConversionPattern=%d{yyyy-MM-dd HH:mm:ss,SSS} %-5p %-60c %x - %m%n

应用ConfigMap:

kubectl apply -f flink-config.yaml
步骤2:部署JobManager(StatefulSet + Service)

JM需要稳定的网络地址和持久存储,因此使用StatefulSet部署。

创建flink-jobmanager.yaml文件:

apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: flink-jobmanager
spec:
  serviceName: flink-jobmanager  # StatefulSet的Headless Service
  replicas: 1
  template:
    metadata:
      labels:
        app: flink
        component: jobmanager
    spec:
      containers:
      - name: jobmanager
        image: flink:1.17.0-scala_2.12-java11  # Flink官方镜像
        args: ["jobmanager"]  # 启动JM进程
        ports:
        - containerPort: 6123  # RPC端口
          name: rpc
        - containerPort: 6124  # Blob Server端口
          name: blob-server
        - containerPort: 8081  # Web UI端口
          name: webui
        - containerPort: 9250  # Prometheus Metrics端口
          name: metrics
        env:
        - name: JOB_MANAGER_RPC_ADDRESS
          value: flink-jobmanager  # 与Service名称一致
        volumeMounts:
        - name: flink-config-volume  # 挂载ConfigMap中的配置
          mountPath: /opt/flink/conf/
      volumes:
      - name: flink-config-volume
        configMap:
          name: flink-config
          items:
          - key: flink-conf.yaml
            path: flink-conf.yaml
          - key: log4j-console.properties
            path: log4j-console.properties
  selector:
    matchLabels:
      app: flink
      component: jobmanager
---
# 为JM创建Service,提供稳定的网络入口
apiVersion: v1
kind: Service
metadata:
  name: flink-jobmanager
spec:
  type: ClusterIP
  ports:
  - name: rpc
    port: 6123
    targetPort: 6123
  - name: blob-server
    port: 6124
    targetPort: 6124
  - name: webui
    port: 8081
    targetPort: 8081
  - name: metrics
    port: 9250
    targetPort: 9250
  selector:
    app: flink
    component: jobmanager

应用JM部署:

kubectl apply -f flink-jobmanager.yaml
步骤3:部署TaskManager(Deployment)

TM是无状态的,使用Deployment部署,支持自动扩缩容。

创建flink-taskmanager.yaml文件:

apiVersion: apps/v1
kind: Deployment
metadata:
  name: flink-taskmanager
spec:
  replicas: 2  # 初始2个TM,每个TM有2个Slot,总Slot数4
  template:
    metadata:
      labels:
        app: flink
        component: taskmanager
    spec:
      containers:
      - name: taskmanager
        image: flink:1.17.0-scala_2.12-java11
        args: ["taskmanager"]  # 启动TM进程
        ports:
        - containerPort: 6122  # RPC端口
          name: rpc
        - containerPort: 9250  # Metrics端口
          name: metrics
        env:
        - name: JOB_MANAGER_RPC_ADDRESS
          value: flink-jobmanager  # 连接JM的Service
        volumeMounts:
        - name: flink-config-volume
          mountPath: /opt/flink/conf/
      volumes:
      - name: flink-config-volume
        configMap:
          name: flink-config
          items:
          - key: flink-conf.yaml
            path: flink-conf.yaml
          - key: log4j-console.properties
            path: log4j-console.properties
  selector:
    matchLabels:
      app: flink
      component: taskmanager

应用TM部署:

kubectl apply -f flink-taskmanager.yaml
步骤4:验证集群状态

查看Pod状态(确保STATUSRunning):

kubectl get pods -l app=flink

输出示例:

NAME                          READY   STATUS    RESTARTS   AGE
flink-jobmanager-0            1/1     Running   0          5m
flink-taskmanager-7f89d6b7c5-2xqkf   1/1     Running   0          3m
flink-taskmanager-7f89d6b7c5-5z7k8   1/1     Running   0          3m

访问Flink Web UI(通过kubectl port-forward暴露端口):

kubectl port-forward service/flink-jobmanager 8081:8081

打开浏览器访问http://localhost:8081,即可看到Flink集群的Dashboard(显示总Slot数4,可用Slot数4)。

步骤5:提交作业

使用Flink CLI向Session集群提交作业(以官方WordCount示例为例):

flink run -d -t kubernetes-session \
  --kubernetesNamespace default \
  --kubernetesJobManagerServiceName flink-jobmanager \
  ./examples/streaming/WordCount.jar

参数说明:

  • -d:后台运行;
  • -t kubernetes-session:指定Session模式;
  • --kubernetesNamespace:K8s命名空间;
  • --kubernetesJobManagerServiceName:JM的Service名称。

提交成功后,在Web UI的“Running Jobs”页面可看到作业状态。

模式2:Job Cluster Mode部署

Job Cluster Mode为每个作业创建独立的JM和TM。我们直接使用Flink CLI的run-application命令部署。

步骤1:提交作业
flink run-application -t kubernetes-application \
  --kubernetesNamespace default \
  --kubernetesClusterId wordcount-job-cluster \  # 集群ID(唯一)
  --image flink:1.17.0-scala_2.12-java11 \
  --image-pull-policy IfNotPresent \
  ./examples/streaming/WordCount.jar
步骤2:验证集群状态

查看作业对应的Pod:

kubectl get pods -l app=flink,cluster-id=wordcount-job-cluster

输出示例:

NAME                                      READY   STATUS    RESTARTS   AGE
wordcount-job-cluster-jobmanager-0        1/1     Running   0          2m
wordcount-job-cluster-taskmanager-5f7d8   1/1     Running   0          1m
步骤3:清理集群

作业完成后,手动删除集群:

kubectl delete clusterrolebinding wordcount-job-cluster-role-binding
kubectl delete deployment wordcount-job-cluster-taskmanager
kubectl delete statefulset wordcount-job-cluster-jobmanager
kubectl delete service wordcount-job-cluster-jobmanager

模式3:Application Mode部署

Application Mode的核心是将应用代码打包进Docker镜像,避免Client端依赖。我们需要先构建自定义镜像,再部署。

步骤1:构建自定义镜像

创建Dockerfile(将WordCount示例Jar包打包进镜像):

# 基于Flink官方镜像
FROM flink:1.17.0-scala_2.12-java11

# 将应用Jar包复制到Flink的用户库目录
COPY ./examples/streaming/WordCount.jar /opt/flink/usrlib/WordCount.jar

# 指定入口命令(Application Mode无需手动启动JM/TM)
CMD ["application"]

构建镜像:

docker build -t my-flink-app:1.0 .

将镜像推送到Kind集群(本地测试无需推送到远程仓库):

kind load docker-image my-flink-app:1.0 --name flink-demo
步骤2:部署Application集群

使用Flink CLI提交Application:

flink run-application -t kubernetes-application \
  --kubernetesNamespace default \
  --kubernetesClusterId my-flink-application \
  --image my-flink-app:1.0 \
  --image-pull-policy Never \  # 本地镜像,无需拉取
  --appArgs "localhost:9999"  # 应用参数(WordCount的输入端口)
步骤3:验证应用状态

查看Pod状态:

kubectl get pods -l app=flink,cluster-id=my-flink-application

访问Web UI(需先暴露JM的Service端口):

kubectl port-forward service/my-flink-application-jobmanager 8081:8081

核心原理深入:Flink如何在K8s上运行?

1. 资源调度流程

Flink on K8s的资源调度由Flink ResourceManagerK8s Scheduler协同完成,流程如下(Mermaid流程图):

ClientJobManager (JM)ResourceManager (RM)Kubernetes Scheduler (K8s)TaskManager Pod (TM)JMRMK8sTM提交作业(包含JobGraph、配置)解析JobGraph,计算所需Task Slot数量请求N个Task Slot调用K8s API创建TM Pod(指定CPU、内存、镜像)将Pod调度到可用Node启动后向JM注册,报告可用Slot将Task分配到Slot,开始执行ClientJobManager (JM)ResourceManager (RM)Kubernetes Scheduler (K8s)TaskManager Pod (TM)JMRMK8sTM

2. 状态管理:如何保证Exactly-Once?

Flink的Checkpoint机制是实现精确一次语义的核心,在K8s上需要解决状态持久化问题。常见方案:

方案1:使用云存储(S3/GCS)

将Checkpoint存储到云对象存储(如AWS S3),配置flink-conf.yaml

# 状态后端:文件系统(支持S3/GCS)
state.backend: filesystem
# Checkpoint存储路径(S3)
state.checkpoints.dir: s3://my-flink-checkpoints/
# Savepoint存储路径(可选)
state.savepoints.dir: s3://my-flink-savepoints/
# S3配置(AWS)
s3.access-key: AKIAXXXXXXXXXXXXXXX
s3.secret-key: XXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX
s3.endpoint: s3.us-east-1.amazonaws.com
方案2:使用K8s PVC

将Checkpoint存储到K8s的PersistentVolume(PV),需先创建PVC:

apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: flink-checkpoint-pvc
spec:
  accessModes:
    - ReadWriteOnce
  resources:
    requests:
      storage: 10Gi

然后配置flink-conf.yaml

state.backend: filesystem
state.checkpoints.dir: file:///opt/flink/checkpoints/

在TM的Deployment中挂载PVC:

volumeMounts:
- name: checkpoint-volume
  mountPath: /opt/flink/checkpoints/
volumes:
- name: checkpoint-volume
  persistentVolumeClaim:
    claimName: flink-checkpoint-pvc

3. 故障恢复:K8s如何保证高可用?

Flink on K8s的高可用依赖K8s的自愈能力Flink的HA机制

(1)JobManager故障恢复
  • JM使用StatefulSet部署,K8s会自动重启故障的JM Pod;
  • JM的HA元数据存储在云存储/PVC中,重启后可恢复作业状态;
  • TM会重新向新的JM注册,恢复Task执行。
(2)TaskManager故障恢复
  • TM使用Deployment部署,K8s会自动创建新的TM Pod替换故障节点;
  • JM会将故障Task的状态从最近的Checkpoint恢复,重新分配到新的TM Slot。
(3)集群整体故障恢复
  • 若整个K8s集群故障(如节点宕机),重启后K8s会重新调度所有Flink Pod;
  • 只要Checkpoint/Savepoint存储在可靠介质(如S3),作业可从Checkpoint恢复。

高级运维:监控、日志与弹性伸缩

1. 监控:Prometheus + Grafana

Flink内置Prometheus Metrics报告器,可将Metrics暴露给Prometheus,再用Grafana可视化。

步骤1:配置Flink Metrics

flink-conf.yaml中添加:

metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter
metrics.reporter.prom.port: 9250-9260  # 每个TM/JM的Metrics端口范围
步骤2:部署Prometheus

创建prometheus.yaml(配置抓取规则):

apiVersion: v1
kind: ConfigMap
metadata:
  name: prometheus-config
data:
  prometheus.yml: |
    global:
      scrape_interval: 15s
    scrape_configs:
      - job_name: 'flink'
        static_configs:
          - targets: ['flink-jobmanager:9250']  # JM的Metrics地址
          - targets: ['flink-taskmanager:9250'] # TM的Metrics地址
---
apiVersion: apps/v1
kind: Deployment
metadata:
  name: prometheus
spec:
  replicas: 1
  template:
    metadata:
      labels:
        app: prometheus
    spec:
      containers:
      - name: prometheus
        image: prom/prometheus:v2.45.0
        ports:
        - containerPort: 9090
        volumeMounts:
        - name: config-volume
          mountPath: /etc/prometheus/
      volumes:
      - name: config-volume
        configMap:
          name: prometheus-config
---
apiVersion: v1
kind: Service
metadata:
  name: prometheus
spec:
  type: NodePort
  ports:
  - port: 9090
    targetPort: 9090
    nodePort: 30090
  selector:
    app: prometheus

应用配置:

kubectl apply -f prometheus.yaml
步骤3:部署Grafana

创建grafana.yaml

apiVersion: apps/v1
kind: Deployment
metadata:
  name: grafana
spec:
  replicas: 1
  template:
    metadata:
      labels:
        app: grafana
    spec:
      containers:
      - name: grafana
        image: grafana/grafana:9.5.2
        ports:
        - containerPort: 3000
        env:
        - name: GF_SECURITY_ADMIN_PASSWORD
          value: admin  # 初始密码
---
apiVersion: v1
kind: Service
metadata:
  name: grafana
spec:
  type: NodePort
  ports:
  - port: 3000
    targetPort: 3000
    nodePort: 30030
  selector:
    app: grafana

应用配置:

kubectl apply -f grafana.yaml
步骤4:导入Flink Dashboard
  1. 访问Grafana(http://localhost:30030,用户名admin,密码admin);
  2. 点击“+”→“Import”,输入Flink官方Dashboard ID:11046
  3. 选择Prometheus数据源,点击“Import”。

即可看到Flink的关键Metrics:

  • 作业状态(Running/Failed/Cancelled);
  • Task Slot利用率;
  • Checkpoint成功率/延迟;
  • 输入/输出吞吐量;
  • 任务延迟(Latency)。

2. 日志收集:Loki + Promtail

Flink的日志默认输出到stdout,可通过**Loki(日志存储)Promtail(日志收集)**实现集中化日志管理。

步骤1:部署Loki

创建loki.yaml

apiVersion: apps/v1
kind: Deployment
metadata:
  name: loki
spec:
  replicas: 1
  template:
    metadata:
      labels:
        app: loki
    spec:
      containers:
      - name: loki
        image: grafana/loki:2.8.0
        ports:
        - containerPort: 3100
---
apiVersion: v1
kind: Service
metadata:
  name: loki
spec:
  type: ClusterIP
  ports:
  - port: 3100
    targetPort: 3100
  selector:
    app: loki

应用配置:

kubectl apply -f loki.yaml
步骤2:部署Promtail

创建promtail.yaml(配置日志收集规则):

apiVersion: v1
kind: ConfigMap
metadata:
  name: promtail-config
data:
  promtail.yml: |
    server:
      http_listen_port: 9080
      grpc_listen_port: 0
    positions:
      filename: /tmp/positions.yaml
    clients:
      - url: http://loki:3100/loki/api/v1/push
    scrape_configs:
      - job_name: kubernetes-pods
        kubernetes_sd_configs:
          - role: pod
        relabel_configs:
          # 仅收集Flink Pod的日志
          - source_labels: [__meta_kubernetes_pod_label_app]
            regex: flink
            action: keep
          # 仅收集JM/TM容器的日志
          - source_labels: [__meta_kubernetes_pod_container_name]
            regex: (jobmanager|taskmanager)
            action: keep
          # 添加Pod名称标签
          - source_labels: [__meta_kubernetes_pod_name]
            target_label: pod
          # 添加命名空间标签
          - source_labels: [__meta_kubernetes_namespace]
            target_label: namespace
---
apiVersion: apps/v1
kind: DaemonSet
metadata:
  name: promtail
spec:
  selector:
    matchLabels:
      app: promtail
  template:
    metadata:
      labels:
        app: promtail
    spec:
      containers:
      - name: promtail
        image: grafana/promtail:2.8.0
        volumeMounts:
        - name: config-volume
          mountPath: /etc/promtail/
        - name: var-log
          mountPath: /var/log
        - name: pods
          mountPath: /var/log/pods
          readOnly: true
      volumes:
      - name: config-volume
        configMap:
          name: promtail-config
      - name: var-log
        hostPath:
          path: /var/log
      - name: pods
        hostPath:
          path: /var/log/pods

应用配置:

kubectl apply -f promtail.yaml
步骤3:在Grafana中查询日志
  1. 进入Grafana,点击“Explore”;
  2. 选择“Loki”数据源;
  3. 使用LogQL查询日志,例如:
    • 查询所有Flink Pod的日志:{app="flink"}
    • 查询JM的错误日志:{app="flink", container="jobmanager"} |= "ERROR"

3. 弹性伸缩:HPA + Flink Operator

Flink on K8s的弹性伸缩有两种方式:基于K8s HPA(简单)和基于Flink Operator(智能)。

方式1:K8s HPA(基于CPU/内存)

HPA(Horizontal Pod Autoscaler)可根据Pod的CPU/内存利用率自动调整副本数。

创建flink-taskmanager-hpa.yaml

apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: flink-taskmanager-hpa
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: flink-taskmanager  # 目标Deployment
  minReplicas: 2  # 最小副本数
  maxReplicas: 10  # 最大副本数
  metrics:
  - type: Resource
    resource:
      name: cpu
      target:
        type: Utilization
        averageUtilization: 70  # CPU利用率超过70%时扩展

应用配置:

kubectl apply -f flink-taskmanager-hpa.yaml
方式2:Flink Operator(基于作业负载)

Flink Operator是官方推出的K8s Operator,可根据作业的负载指标(如待处理记录数、延迟)自动调整并行度。

部署Flink Operator:

kubectl create -f https://github.com/apache/flink-kubernetes-operator/releases/download/v1.5.0/flink-operator.yaml

创建Flink作业CRD(Custom Resource Definition):

apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: my-flink-job
spec:
  image: my-flink-app:1.0
  flinkVersion: v1_17
  flinkConfiguration:
    taskmanager.numberOfTaskSlots: "2"
  serviceAccount: flink
  jobManager:
    replicas: 1
    resource:
      memory: "1024m"
      cpu: 1
  taskManager:
    resource:
      memory: "2048m"
      cpu: 1
  job:
    jarURI: local:///opt/flink/usrlib/WordCount.jar
    parallelism: 4  # 初始并行度
    autoScaling:
      enabled: true  # 开启自动扩缩容
      minParallelism: 2
      maxParallelism: 10
      scalingPolicy:
        type: "Threshold"
        threshold:
          # 当待处理记录数超过1000时扩展,低于200时缩容
          pendingRecordsThreshold: 1000
          scaledownDelay: 300s  # 缩容延迟(避免频繁调整)

应用CRD:

kubectl apply -f my-flink-job.yaml

Flink Operator会持续监控作业的负载指标,自动调整TM的数量和并行度,比HPA更智能。

实际应用场景:Flink on K8s的生产实践

场景1:电商实时推荐系统

业务需求
  • 实时收集用户浏览、点击、购买行为;
  • 计算用户实时兴趣特征(如最近浏览的商品类别);
  • 将特征实时推送给推荐系统,实现“千人千面”推荐。
技术方案
  • 部署模式:Application Mode(每个推荐模型对应一个独立集群,资源隔离);
  • 状态管理:使用S3存储Checkpoint(保证Exactly-Once);
  • 弹性伸缩:Flink Operator根据用户流量自动调整并行度;
  • 监控:Prometheus + Grafana监控作业延迟、吞吐量;
  • 日志:Loki + Promtail收集日志,快速排查问题。
收益
  • 低延迟:处理延迟从秒级降至毫秒级;
  • 高可用:故障恢复时间从分钟级降至30秒内;
  • 资源利用率:闲时资源使用率从20%提升至50%,忙时自动扩展。

场景2:日志实时分析系统

业务需求
  • 实时收集服务器、应用的日志;
  • 解析日志中的错误、警告信息;
  • 实时告警(如CPU利用率超过90%、请求失败率超过5%)。
技术方案
  • 部署模式:Job Cluster Mode(每个日志源对应一个作业,隔离故障);
  • 状态管理:使用PVC存储Checkpoint(本地存储,低延迟);
  • 弹性伸缩:K8s HPA根据日志吞吐量调整TM数量;
  • 集成:与ELK(Elasticsearch + Logstash + Kibana)结合,实现日志查询与可视化。
收益
  • 实时性:日志从产生到告警的时间从5分钟降至1分钟;
  • 可扩展性:支持日均10TB日志的处理;
  • 运维效率:自动扩缩容减少了90%的手动干预。

工具与资源推荐

1. 部署工具

  • Flink Kubernetes Operator:官方Operator,管理Flink集群生命周期;
  • Kind:本地K8s集群搭建工具;
  • Kubectl:K8s命令行工具;
  • Helm:K8s包管理工具(可快速部署Flink Operator、Prometheus等)。

2. 监控与日志

  • Prometheus: metrics 收集与存储;
  • Grafana:可视化Dashboard;
  • Loki:日志存储;
  • Promtail:日志收集。

3. 学习资源

  • Flink官方文档:https://nightlies.apache.org/flink/flink-docs-stable/
  • K8s官方文档:https://kubernetes.io/docs/
  • Flink on K8s实战:《Apache Flink 实战》第10章;
  • 社区资源:Flink中文社区(https://flink.sojb.cn/)、K8s中文社区(https://kubernetes.io/zh-cn/)。

未来趋势与挑战

1. 未来趋势

  • Serverless Flink:公有云厂商(如阿里云、AWS)推出Serverless Flink服务,用户无需管理集群,按使用量付费;
  • GitOps:使用Git管理Flink集群配置,通过Argo CD自动同步到K8s,实现“基础设施即代码”;
  • AI辅助调优:通过机器学习模型预测作业负载,自动调整并行度、Checkpoint间隔等参数;
  • 流批一体:Flink的Table API支持流批一体,K8s的弹性资源可支持流作业与批作业共享集群。

2. 面临的挑战

  • 资源调度延迟:K8s创建Pod需要时间(秒级),对于低延迟作业(如金融交易)可能影响SLA;
  • 状态持久化性能:云存储(如S3)的延迟可能导致Checkpoint失败;
  • 多租户隔离:需要精细配置ResourceQuotas、LimitRanges,避免作业间资源抢占;
  • 学习成本:需要掌握Flink、K8s、监控、日志等多门技术,门槛较高。

总结

Flink on Kubernetes是云原生时代流处理的最佳实践——它将Flink的流处理能力与K8s的资源管理能力深度融合,解决了传统部署方式的痛点,实现了“弹性、高可用、易运维”的生产级流处理系统。

虽然部署和运维需要一定的学习成本,但通过选择合适的部署模式(如Application Mode)、使用官方工具(如Flink Operator)、结合云原生生态(如Prometheus、Loki),我们可以快速搭建稳定、高效的Flink on K8s集群。

随着Serverless、GitOps、AI调优等技术的发展,Flink on K8s的未来将更加光明——它将成为企业实时计算的“基础设施”,支撑更多复杂的业务场景(如实时推荐、实时风控、实时数仓)。

如果你正在考虑将Flink迁移到K8s,现在就是最佳时机!

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐