📝 摘要:基于 Docker 单机部署 Flink 1.20.1(Java 17 镜像),一份 docker-compose 搞定 JobManager + TaskManager,含 checkpoint/savepoint 挂载与 OSHI 系统指标;再集成 Prometheus 采集 + 任务停止/频繁重启/Checkpoint 失败/反压/TM 内存 GC/Slot 不足等生产级告警规则;最后用 Grafana 官方看板(ID 14911)一键出图,部署、监控、告警、可视化一条龙落地。

本文手把手教你用 Docker 部署 Flink 1.20,从零到跑通任务只需 5 分钟。配套 Prometheus 监控告警规则,让你的 Flink 任务稳如老狗。

前言

想玩 Flink 却被复杂的集群部署劝退?想在本地快速验证一个 SQL 任务却懒得折腾环境?

Docker 单机部署了解一下——一条命令启动,开箱即用,测试环境首选。

本文基于 Flink 1.20.1 + Java 17 镜像,一条龙搞定:

  • 完整的 docker-compose 配置(JobManager + TaskManager)
  • Prometheus 监控指标采集
  • 生产级告警规则(拿来即用)
  • Grafana 官方看板一键导入

废话不多说,直接开干。

⚠️ 先说清适用边界

本文这份配置没有开高可用(HA)——JobManager 一挂,正在跑的作业连同它的 JobGraph 一起消失,得人工重新提交。所以它适合的是本地开发、功能验证、跑完即走的小任务。

但**「单机」不等于「不能上生产」:只要加上 ZooKeeper HA,同样是一个 JobManager,作业就能在 JM 重启后自动从最后一个 checkpoint 续上——这是一种真实可用的生产形态**,我自己的生产环境就是这么跑的。

差别到底在哪、要改哪几行,见文末 §八 从「能跑」到「能上生产」


一、环境准备

1.1 目录结构

先把目录建好,别等会儿挂载的时候报错:

mkdir -p /data/flink/{lib,conf,sql,checkpoints,savepoints,jars}
目录 用途
lib 额外的 jar 包(连接器、驱动等)
conf 自定义配置文件
sql Flink SQL 脚本
checkpoints Checkpoint 持久化目录
savepoints Savepoint 持久化目录
jars 监控相关的依赖包

1.2 监控依赖(可选)

如果需要采集系统资源指标(CPU、内存、磁盘等),需要额外下载 OSHI 库和 JNA 依赖:

# OSHI:Java 获取系统信息的神器
wget -P /data/flink/jars https://repo1.maven.org/maven2/com/github/oshi/oshi-core/6.4.11/oshi-core-6.4.11.jar

# JNA:OSHI 的底层依赖,用于调用操作系统 API
wget -P /data/flink/jars https://repo1.maven.org/maven2/net/java/dev/jna/jna/5.13.0/jna-5.13.0.jar
wget -P /data/flink/jars https://repo1.maven.org/maven2/net/java/dev/jna/jna-platform/5.13.0/jna-platform-5.13.0.jar

不需要监控?这步可以跳过——但记得把后面 compose 里对应的 volumes 挂载metrics.system-resource: true 一起去掉。只注释卷挂载、却留着 metrics.system-resource,Flink 会因为找不到 OSHI 库在启动日志里一直刷告警。


二、Docker Compose 配置

创建 docker-compose.yml,这是全文的核心:

version: '3.8'

services:
  jobmanager:
    image: flink:1.20.1-scala_2.12-java17
    hostname: jobmanager
    ports:
      - "8081:8081"   # Web UI,日常看任务状态
      - "6123:6123"   # RPC 端口,TaskManager 连接用
      - "9020:9020"   # Prometheus 指标端口
    command: jobmanager
    environment:
      - TZ=Asia/Shanghai
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
        state.backend: filesystem
        state.checkpoints.dir: file:///opt/flink/checkpoints/
        state.savepoints.dir: file:///opt/flink/savepoints/
        execution.checkpointing.interval: 30000
        metrics.reporters: prom                   # 基于 Prometheus 进行监控,注意注意注意,此行不能有注释
        metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory
        metrics.reporter.prom.port: 9020
        metrics.system-resource: true
    volumes:
      # Flink 的系统资源指标依赖 OSHI 库
      - /data/flink/jars/oshi-core-6.4.11.jar:/opt/flink/lib/oshi-core-6.4.11.jar
      # oshi-core(OSHI 库)依赖另一个库:JNA(Java Native Access) 来调用操作系统底层 API
      - /data/flink/jars/jna-5.13.0.jar:/opt/flink/lib/jna-5.13.0.jar
      - /data/flink/jars/jna-platform-5.13.0.jar:/opt/flink/lib/jna-platform-5.13.0.jar
      # - /data/flink/conf:/opt/flink/conf          # 如果把 conf目录 copy 出来,这行可以取消注释
      # - /data/flink/sql:/opt/flink/sql            # 如果需要sql 等,这行可以取消注释
      - /tmp:/tmp                                 # 如果不需要安装其他中间件(如 flink-cdc),可以注释掉此行
      - /data/flink/checkpoints:/opt/flink/checkpoints
      - /data/flink/savepoints:/opt/flink/savepoints

  taskmanager:
    image: flink:1.20.1-scala_2.12-java17
    depends_on:
      - jobmanager
    command: taskmanager
    ports:
      - "9021:9021"   # Prometheus 指标端口
    environment:
      - TZ=Asia/Shanghai
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
        state.backend: filesystem
        state.checkpoints.dir: file:///opt/flink/checkpoints/
        state.savepoints.dir: file:///opt/flink/savepoints/
        taskmanager.numberOfTaskSlots: 4
        metrics.reporters: prom
        metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory
        metrics.reporter.prom.port: 9021-9040
        metrics.system-resource: true
    volumes:
      # Flink 的系统资源指标依赖 OSHI 库
      - /data/flink/jars/oshi-core-6.4.11.jar:/opt/flink/lib/oshi-core-6.4.11.jar
      # oshi-core(OSHI 库)依赖另一个库:JNA(Java Native Access) 来调用操作系统底层 API
      - /data/flink/jars/jna-5.13.0.jar:/opt/flink/lib/jna-5.13.0.jar
      - /data/flink/jars/jna-platform-5.13.0.jar:/opt/flink/lib/jna-platform-5.13.0.jar
      - /data/flink/checkpoints:/opt/flink/checkpoints
      - /data/flink/savepoints:/opt/flink/savepoints

配置说明

配置项 说明
state.backend: filesystem 使用文件系统存储状态,生产环境建议用 RocksDB
execution.checkpointing.interval: 30000 每 30 秒做一次 Checkpoint
taskmanager.numberOfTaskSlots: 4 每个 TaskManager 提供 4 个 Slot
metrics.reporter.prom.port: 9021-9040 TaskManager 的指标端口范围,支持扩容

踩坑提醒metrics.reporters: prom 这行配置后面不能加注释,否则会解析失败。别问我怎么知道的。


三、启动与验证

3.1 启动服务

docker compose up -d

查看容器状态:

docker compose ps

3.2 访问 Web UI

浏览器打开:http://<你的IP>:8081

看到这个界面就说明部署成功了:

Flink 1.20 基于 Docker 单机部署成功后的 Web UI 页面


四、提交任务

两种方式,看你喜欢哪种:

下面的容器名 flink-jobmanager-1 是 Compose 按「目录名-服务名-序号」自动生成的——如果你的 docker-compose.yml 不在 flink 目录下,名字会不一样,先用 docker compose ps 查一眼实际容器名再替换。

方式一:进入容器执行

# 进入 JobManager 容器
docker exec -it flink-jobmanager-1 /bin/bash

# 执行 SQL 任务
cd /opt/flink && ./bin/sql-client.sh -f /opt/flink/sql/your-task.sql

方式二:一行命令搞定

docker exec -it flink-jobmanager-1 /opt/flink/bin/sql-client.sh -f /opt/flink/sql/your-task.sql

推荐方式二,简洁优雅,适合写到脚本里。


五、Prometheus 监控集成

没有监控的服务就像没有仪表盘的汽车——出了问题两眼一抹黑。

5.1 配置 Prometheus 采集

prometheus.yml 中添加:

scrape_configs:
  # JobManager 指标
  - job_name: 'flink-jm'
    static_configs:
      - targets: ['<JobManager-IP>:9020']
        labels:
          group: jm

  # TaskManager 指标
  - job_name: 'flink-tm'
    static_configs:
      - targets: ['<TaskManager-IP>:9021']
        labels:
          group: tm

5.2 告警规则

创建 flink-rules.yml,这套规则覆盖了 Flink 运维的核心场景:

groups:
  # ==================== 任务级告警 ====================
  - name: flink-job-alerts
    rules:
      # 任务停止运行 - 最高优先级
      - alert: Flink任务停止运行
        expr: flink_jobmanager_job_uptime == 0
        for: 1m
        labels:
          severity: critical
          group: flink
        annotations:
          summary: "Flink 任务停止运行"
          description: "任务 {{ $labels.job_name }} 已停止运行超过 1 分钟"

      # 任务重启 - 需要关注
      - alert: Flink任务重启
        expr: increase(flink_jobmanager_job_numRestarts[5m]) > 0
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink 任务发生重启"
          description: "任务 {{ $labels.job_name }} 在过去 5 分钟内重启了 {{ $value }} 次"

      # 频繁重启 - 说明有大问题
      - alert: Flink任务频繁重启
        expr: increase(flink_jobmanager_job_numRestarts[30m]) > 3
        labels:
          severity: critical
          group: flink
        annotations:
          summary: "Flink 任务频繁重启"
          description: "任务 {{ $labels.job_name }} 在过去 30 分钟内重启了 {{ $value }} 次,请立即检查"

      # Checkpoint 失败
      - alert: Flink Checkpoint失败
        expr: increase(flink_jobmanager_job_numberOfFailedCheckpoints[10m]) > 0
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink Checkpoint 失败"
          description: "任务 {{ $labels.job_name }} 在过去 10 分钟内有 {{ $value }} 次 Checkpoint 失败"

      # Checkpoint 耗时过长
      - alert: Flink Checkpoint耗时过长
        expr: flink_jobmanager_job_lastCheckpointDuration > 120000
        for: 5m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink Checkpoint 耗时过长"
          description: "任务 {{ $labels.job_name }} Checkpoint 耗时超过 2 分钟(当前:{{ $value }}ms)"

      # 长时间没有成功的 Checkpoint
      - alert: Flink长时间无Checkpoint
        expr: time() - flink_jobmanager_job_lastCheckpointRestoreTimestamp > 1800
        for: 5m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink 长时间无 Checkpoint"
          description: "任务 {{ $labels.job_name }} 超过 30 分钟没有成功的 Checkpoint"

  # ==================== TaskManager 告警 ====================
  - name: flink-taskmanager-alerts
    rules:
      # 堆内存使用过高
      - alert: Flink TM堆内存过高
        expr: flink_taskmanager_Status_JVM_Memory_Heap_Used / flink_taskmanager_Status_JVM_Memory_Heap_Max > 0.85
        for: 5m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "TaskManager 堆内存使用过高"
          description: "TaskManager {{ $labels.tm_id }} 堆内存使用率达 {{ printf \"%.1f\" $value }}%"

      # GC 时间过长
      - alert: Flink TM GC时间过长
        expr: increase(flink_taskmanager_Status_JVM_GarbageCollector_G1_Old_Generation_Time[5m]) > 30000
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "TaskManager GC 时间过长"
          description: "TaskManager {{ $labels.tm_id }} 过去 5 分钟 Old GC 耗时 {{ $value }}ms"

      # 输入数据停滞
      - alert: Flink任务无输入数据
        expr: rate(flink_taskmanager_job_task_numRecordsIn[5m]) == 0 and flink_taskmanager_job_task_numRecordsIn > 0
        for: 10m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink 任务无输入数据"
          description: "任务 {{ $labels.job_name }} 的 {{ $labels.task_name }} 超过 10 分钟无输入,检查数据源"

      # 输出数据停滞
      - alert: Flink任务无输出数据
        expr: rate(flink_taskmanager_job_task_numRecordsOut[5m]) == 0 and flink_taskmanager_job_task_numRecordsOut > 0
        for: 10m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink 任务无输出数据"
          description: "任务 {{ $labels.job_name }} 的 {{ $labels.task_name }} 超过 10 分钟无输出,检查 Sink"

      # 反压告警
      - alert: Flink任务反压
        expr: flink_taskmanager_job_task_backPressuredTimeMsPerSecond > 500
        for: 5m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "Flink 任务出现反压"
          description: "任务 {{ $labels.job_name }} 的 {{ $labels.task_name }} 反压严重(每秒 {{ $value }}ms)"

  # ==================== 集群级告警 ====================
  - name: flink-cluster-alerts
    rules:
      # JobManager 挂了
      - alert: Flink JobManager不可用
        expr: up{job="flink-jm"} == 0
        for: 1m
        labels:
          severity: critical
          group: flink
        annotations:
          summary: "JobManager 不可用"
          description: "Flink JobManager {{ $labels.instance }} 已停止响应"

      # TaskManager 挂了
      - alert: Flink TaskManager不可用
        expr: up{job="flink-tm"} == 0
        for: 1m
        labels:
          severity: critical
          group: flink
        annotations:
          summary: "TaskManager 不可用"
          description: "Flink TaskManager {{ $labels.instance }} 已停止响应"

      # TaskManager 数量不足
      - alert: Flink TaskManager数量减少
        expr: flink_jobmanager_numRegisteredTaskManagers < 1
        for: 2m
        labels:
          severity: critical
          group: flink
        annotations:
          summary: "TaskManager 数量不足"
          description: "当前注册的 TaskManager 数量为 {{ $value }},请检查集群状态"

      # Slot 不足
      - alert: Flink可用Slot不足
        expr: flink_jobmanager_taskSlotsAvailable == 0
        for: 5m
        labels:
          severity: warning
          group: flink
        annotations:
          summary: "可用 Slot 不足"
          description: "当前没有可用的 Task Slot,可能影响新任务提交"

六、Grafana 可视化看板

Prometheus 有了指标和告警,还差最后一块——一个能一眼看清任务健康度的看板。好消息是:Flink 的指标已经被 Prometheus 采走了(§五),Grafana 这边不用像 ClickHouse 那样单独装数据源插件,直接复用 Prometheus 数据源导一个现成看板即可。

6.1 确认 Prometheus 数据源

打开 Grafana → Connections → Data sources,确认已经有一个指向你 Flink 采集端的 Prometheus 数据源。没有就 Add data source → Prometheus → URL 填 http://<Prometheus-IP>:9090 → Save & Test。

Grafana / Prometheus 本身怎么搭、告警怎么接飞书钉钉,见配套的 《Prometheus + Grafana + AlertManager 监控体系搭建:Docker 一把梭》

6.2 导入 Flink 官方看板

  1. Grafana → Dashboards → New → Import
  2. Import via grafana.com 填看板 ID,点 Load
  3. 数据源选上一步的 Prometheus,点 Import

两个现成看板,按需选:

ID 名称 特点
14911 Apache Flink (2021) Dashboard for Job / Task Manager 按 JobManager / TaskManager 维度分组,Heap / Direct / Metaspace / GC / Checkpoint / 反压一屏看全,信息最全,推荐
10369 Flink Dashboard 经典基础版:CPU、内存、Task Slot、运行 Job 数总览,简洁够用

10369 总览——Task Slot 用量、运行中 Job 数一眼可见:

Flink Grafana 看板 10369:JobManager 与 TaskManager 的 CPU 负载、内存、可用 Task Slot、运行中 Job 数量总览

14911 详细——JVM 各区内存、GC、Checkpoint 分组下钻:

Flink Grafana 官方看板 14911:Job/Task Manager 的 JVM Heap、Direct、Mapped、Metaspace 内存与 GC、Slots、Checkpoint 分组面板

导入后看板空白 / 没数据?按顺序查三件事:① Prometheus → Status → Targets 里 flink-jm / flink-tm 是不是 UP;② compose 里 metrics.reporters: prom 那行有没有被行内注释带崩(见 §二踩坑);③ 看板右上角 / 变量里选的数据源是不是你的 Prometheus。


七、常见问题

Q1: 容器启动后 TaskManager 连不上 JobManager?

检查 jobmanager.rpc.address 配置是否正确,确保两个容器在同一个 Docker 网络中。

Q2: Checkpoint 目录权限问题?

确保宿主机的 /data/flink/checkpoints 目录有读写权限:

chmod -R 777 /data/flink/checkpoints /data/flink/savepoints

Q3: 需要扩展更多 TaskManager?

修改 docker-compose.yml,复制 taskmanager 服务配置,改个名字即可。别忘了调整端口映射避免冲突。


八、从「能跑」到「能上生产」:要补的其实不多

网上一提「Flink 上生产」,标配答案就是「上 K8s 或 YARN」。但实际上,Docker Compose 这套单机形态,加上 ZooKeeper HA 之后就能稳定跑生产——我自己的生产环境就是这么跑的,一个 JobManager + ZK quorum + restart: always

要理解为什么够用,得先分清一件常被混淆的事。

8.1 「高可用」其实是两件事,很多人把它们当成一件

能力 需要什么 没有它会怎样
① 作业不丢(故障恢复) 开 HA 即可,一个 JobManager 就够 JM 一挂,JobGraph 连同运行中的作业一起消失,要人工重新提交
② 秒级接管(故障转移) 必须有 ≥2 个 JobManager JM 挂到被拉起来这段时间,作业停止调度(有中断,但不丢)

绝大多数场景真正怕的是 ——凌晨作业挂了没人重提,第二天数据断一天。而① 只要开 HA 就有了,不需要多 JobManager

官方文档把 high-availability.storageDir 的职责写得很清楚:

“Persisting state which is required for the successor to resume the job execution (JobGraphs, user code jars, completed checkpoints)”

也就是说 JobGraph 和已完成的 checkpoint 是被持久化下来的,ZooKeeper 只存「谁是 leader」和指针。所以配上 restart: always 之后:

JobManager 挂 → Docker 自动拉起 → 重新获得 leadership
              → 从 HA store 读回 JobGraph → 从最后一个 checkpoint 续跑

全程无人值守。 这就是单 JM + ZK HA 能扛生产的原因——它换来的不是「零中断」,而是**「中断有界、作业不丢」**。

8.2 最小改动:在本文配置上加 5 行

high-availability.type: zookeeper
high-availability.zookeeper.quorum: zk1:2181,zk2:2181,zk3:2181   # 3 节点
high-availability.storageDir: file:///opt/flink/ha/               # 见 8.4 的取舍
high-availability.zookeeper.path.root: /flink
high-availability.cluster-id: flink-prod                          # 多集群共用 ZK 时靠它隔离

配合 restart: always(本文的 compose 已经有了),JobManager 挂掉就能自愈。

ZK 要奇数节点(3 或 5),选主靠多数派;单节点 ZK 没有意义——那只是把单点从 JobManager 挪到了 ZK 身上。ZK 集群怎么搭、怎么接监控,见 ZooKeeper 从裸奔到全副武装:搭建、可视化、监控、告警、看板一条龙

⚠️ ZK 一旦成为 Flink 的依赖,它自己就必须被监控。ZK 挂了,leader 选举跟着失效——别让兜底组件变成新的单点zk_up 和节点数告警属于必配。

8.3 不想引入 ZooKeeper?跑在 K8s 上就不用

Flink 从 1.12 起提供了原生 Kubernetes HA,用 K8s 的 ConfigMap 做 leader 选举:

high-availability.type: kubernetes
high-availability.storageDir: s3://your-bucket/flink/ha
kubernetes.cluster-id: your-flink-cluster
HA 方式 要 ZooKeeper 吗 适用
high-availability.type: zookeeper ✅ 要 quorum 任意部署形态(Docker / standalone / YARN)
high-availability.type: kubernetes 不用 仅限跑在 K8s 上

所以「Flink 上生产必须搭 ZK」是个过时说法——只有不在 K8s 上跑时才必须。已经有 K8s 的团队,这条路能少维护一个 ZK 集群。

8.4 真正该注意的取舍:storageDir 放本地盘意味着什么

high-availability.storageDir 和 checkpoint 目录可以放本地盘file:///… + 挂载宿主机目录),单机形态下完全能跑。但要清楚它换来的约束:

放本地盘 后果
✅ 简单,不依赖对象存储 单机形态够用,JM 重启照样恢复
⚠️ 横向扩不了 想再加一个 JobManager 做 standby、或把 TaskManager 挪到别的机器,对方读不到 HA 元数据和 checkpoint,HA 直接失效
⚠️ 磁盘即单点 这块盘坏了,HA 元数据和 checkpoint 一起没,恢复无从谈起

判断标准很简单:作业只跑在这一台机器上、且这块盘有备份或能接受重跑,本地盘就够;一旦要多机,storageDir 和 checkpoint 必须换成所有节点都能访问的存储(S3 / HDFS / NFS)。

8.5 另外两处按需调

本文默认 什么时候要改
state.backend: filesystem 状态存内存、checkpoint 时落盘 状态大到内存扛不住就换 RocksDB(支持增量 checkpoint)。状态小的同步类作业不用改
taskmanager.numberOfTaskSlots 默认较小 按作业并行度调,slot 不够任务提交不上去

总结

Docker 部署 Flink 的优势:

  • 快速:一条命令启动,环境隔离干净
  • 灵活:挂载目录,配置随改随生效
  • 可观测:Prometheus 采集 + 告警规则 + Grafana 官方看板,部署、监控、告警、可视化一步到位

本文这份配置没开 HA,适合本地开发、功能验证、跑完即走的小任务。

但从这里到生产,距离比想象中短——加上 §8.2 那 5 行 ZooKeeper HA 配置,一个 JobManager 也能做到「挂了自愈、作业不丢」。真正要想清楚的不是「够不够格上生产」,而是 §8.4 那个取舍:storageDir 放本地盘就锁死在单机,要多机就得换共享存储。


如果这篇文章对你有帮助,欢迎点赞收藏。有问题欢迎评论区交流。


延伸阅读

外部权威


📅 最后更新:2026-08-05,新增 §八「从『能跑』到『能上生产』」。

⚠️ 修正一处结论:原文末尾写的「生产环境还是老老实实上 K8s 或 YARN 吧」并不准确——Docker Compose 这套单机形态加上 ZooKeeper HA,是能稳定跑生产的,我自己的生产环境就是这么跑的。

关键是把「高可用」拆成两件事:作业不丢(开 HA 即可,一个 JobManager 就够)vs 秒级接管(才需要多 JobManager)。混淆这两点就会得出「单 JM 不能上生产」的错误结论。详见新增的 §八。

Logo

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

更多推荐