Flink 1.20 从裸奔到全副武装:Docker 部署 + Prometheus 监控告警 + Grafana 看板一条龙
📝 摘要:基于 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-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 官方看板
- Grafana → Dashboards → New → Import
- 在 Import via grafana.com 填看板 ID,点 Load
- 数据源选上一步的 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 数一眼可见:

14911 详细——JVM 各区内存、GC、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 放本地盘就锁死在单机,要多机就得换共享存储。
如果这篇文章对你有帮助,欢迎点赞收藏。有问题欢迎评论区交流。
延伸阅读
- Flink 1.20 实战:MySQL 实时同步到 ClickHouse —— 部署好 Flink 后的第一个实战:把 MySQL binlog 实时同步到 CK
- Flink 1.20 + CDC 3.5 实战:MongoDB 实时同步到 ClickHouse,从踩坑到上线 —— 换个数据源,用同一套 Flink 环境同步 MongoDB
- ClickHouse 25.4 基于 Docker 单机部署实战指南 —— 同步链路的下游存储,同样一条 docker 命令搞定
- Flink Jenkinsfile 怎么写不出 bug:10 条设计要点 + 完整 Demo —— 把 Flink 作业的部署做成 savepoint 优雅停启的 CI/CD 流水线
- Prometheus + Grafana + AlertManager 监控体系搭建:Docker 一把梭 —— 本文 Grafana 看板的底座,附 15 种组件的看板 ID 与告警接飞书/钉钉
- ZooKeeper 从裸奔到全副武装:搭建、可视化、监控、告警、看板一条龙 —— §8.2 走 standalone / YARN 路线时,ZK 集群怎么搭、怎么监控
外部权威:
- Flink 1.20 官方:High Availability 概览 —— 两种 HA 实现的官方说明(ZooKeeper HA 需要 quorum,Kubernetes HA 仅限 K8s)
- Flink 1.20 官方:Kubernetes HA —— §8.1 配置的出处,含 ConfigMap 只存指针、元数据落
storageDir的机制 - Flink 1.20 官方:ZooKeeper HA —— §8.2 配置的出处
📅 最后更新:2026-08-05,新增 §八「从『能跑』到『能上生产』」。
⚠️ 修正一处结论:原文末尾写的「生产环境还是老老实实上 K8s 或 YARN 吧」并不准确——Docker Compose 这套单机形态加上 ZooKeeper HA,是能稳定跑生产的,我自己的生产环境就是这么跑的。
关键是把「高可用」拆成两件事:作业不丢(开 HA 即可,一个 JobManager 就够)vs 秒级接管(才需要多 JobManager)。混淆这两点就会得出「单 JM 不能上生产」的错误结论。详见新增的 §八。
更多推荐



所有评论(0)