📌 本文是「Flink 教程」系列第 三 篇 · 系列目录
3.1 部署模式
三种部署模式对比:
| 模式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| Standalone | 部署简单、无需依赖 | 资源无法弹性伸缩 | 开发环境、简单生产 |
| YARN | 资源统一管理、弹性伸缩 | 依赖 Hadoop 集群 | 已有 Hadoop 环境 |
| Kubernetes | 云原生、自动伸缩、DevOps 友好 | 需 K8s 运维能力 | 云原生部署、容器化环境 |
3.1.1 独立集群(Standalone)
适用于开发和简单生产环境:
# 启动集群
bin/start-cluster.sh
# 启动 JobManager
bin/jobmanager.sh start
# 启动 TaskManager
bin/taskmanager.sh start
配置(flink-conf.yaml):
jobmanager.rpc.address: localhost
jobmanager.rpc.port: 6123
jobmanager.memory.process.size: 1600m
taskmanager.memory.process.size: 4096m
taskmanager.numberOfTaskSlots: 4
parallelism.default: 2
3.1.2 YARN 集群
适用于 Hadoop 集群环境:
# Session 模式
bin/flink-yarn.sh -jm 1024m -tm 4096m -s 4
# 分离模式提交
bin/flink run -m yarn-cluster -yqu root.my-queue -yjm 1024m -ytm 4096m myjob.jar
# 指定每个 TaskManager 的 slot 数
bin/flink run -m yarn-cluster -yys 4 myjob.jar
3.1.3 Kubernetes
适用于云原生环境:
Native Kubernetes:
- Flink 自动管理 Pod 创建
- 支持 Session 模式和 Application 模式
配置:
kubernetes.cluster-id: my-flink-cluster
kubernetes.namespace: flink
kubernetes.container.image: flink:1.14
3.2 高可用配置
两种 HA 方案对比:
| 对比项 | ZooKeeper HA | Kubernetes HA |
|---|---|---|
| 依赖 | 需部署 Zookeeper 集群 | 依赖 Kubernetes |
| 存储后端 | HDFS | S3/NFS/OSS |
| 适用部署 | Standalone / YARN | Kubernetes |
| 元数据恢复 | 从 ZK 获取 JobManager 地址 | 从 K8s API 获取 |
3.2.1 ZooKeeper HA
# flink-conf.yaml
high-availability: zookeeper
high-availability.storageDir: hdfs://namenode:40010/flink/ha
high-availability.zookeeper.quorum: zk1:2181,zk2:2181,zk3:2181
high-availability.zookeeper.path.root: /flink
3.2.2 Kubernetes HA
# flink-conf.yaml
high-availability: kubernetes
high-availability.cluster-id: my-flink-cluster
high-availability.storageDir: s3://my-bucket/flink/ha
kubernetes.namespace: flink
3.3 内存配置
3.3.1 TaskManager 内存
# 总内存(推荐)
taskmanager.memory.process.size: 4096m
# JVM 堆内存
taskmanager.memory.heap.size: 3072m
# 托管内存(用于 RocksDB 等)
taskmanager.memory.managed.size: 512m
taskmanager.memory.managed.fraction: 0.4
# 网络内存
taskmanager.memory.network.min: 64m
taskmanager.memory.network.max: 1g
taskmanager.memory.network.fraction: 0.1
3.3.2 生产环境内存配置参考
# JobManager 配置 (控制节点)
# 控制节点压力较小,通常 2GB-4GB 足够,除非 Checkpoint 元数据极大
jobmanager.memory.process.size: 2g
# TaskManager 配置 (工作节点) - 假设单容器 8GB 内存 / 4 CPU
taskmanager.memory.process.size: 8g
# [核心配置] 托管内存比例 (Managed Memory Fraction)
# 默认 0.4。如果你使用 RocksDB,建议保持 0.4 或调高到 0.5。
# Flink 会把这部分内存划给 RocksDB 用作 MemTable 和 Block Cache。
# ⚠️ 使用 RocksDB 时,状态存储在 Off-Heap 的 Managed Memory 中,而不是 JVM Heap
taskmanager.memory.managed.fraction: 0.4
# JVM Metaspace (防止加载大量类时 OOM)
taskmanager.memory.jvm-metaspace.size: 256m
# 网络缓冲区 (处理反压的关键)
# 如果网络 Shuffle 数据量大,可能需要适当调大
taskmanager.memory.network.fraction: 0.1
taskmanager.memory.network.min: 64mb
taskmanager.memory.network.max: 1gb
# Slot 配置:推荐 Slot 数 = CPU 核心数,避免线程上下文切换开销
taskmanager.numberOfTaskSlots: 4
# 全局默认并行度 (建议设置为总 Slot 数)
parallelism.default: 4
3.4 Checkpoint 配置
# 开启 Checkpoint (毫秒) - 例如每 60 秒一次
execution.checkpointing.interval: 60s
# Checkpoint 语义:EXACTLY_ONCE 保证精确一次,AT_LEAST_ONCE 性能更好
execution.checkpointing.mode: EXACTLY_ONCE
# Checkpoint 超时时间 (毫秒):如果状态较大,需要调大此值
execution.checkpointing.timeout: 600s
# 两个 Checkpoint 之间的最小间隔 (毫秒):防止 Checkpoint 过于频繁
execution.checkpointing.min-pause: 30s
# 允许最大并发 Checkpoint 数 (通常设为 1,防止 IO 风暴)
execution.checkpointing.max-concurrent-checkpoints: 1
# 存储配置
state.checkpoints.storage: filesystem
# Checkpoint 存储路径 (必须是分布式文件系统 HDFS/S3/MinIO)
state.checkpoints.dir: hdfs://namenode:40010/flink/checkpoints
# 非对齐 Checkpoint(Flink 1.11+):解决严重反压导致 Checkpoint 超时的问题
# Barrier 插队优先处理,但需将缓冲区数据一起存盘
execution.checkpointing.unaligned: true
Checkpoint 监控与优化细节
Checkpoint Delay Time 监控(生产环境重要):
checkpoint_start_delay = end_to_end_duration - async_duration - sync_duration
- 如果该时间过长,说明算子正在进行 Barrier 对齐,等待上游算子将数据写入到当前算子中
- 这表明系统正处于反压状态下
- 可以通过 Flink Web UI → Job 详情 → Checkpointing → Summary 监控这个指标
状态数据压缩:
ExecutionConfig executionConfig = new ExecutionConfig();
executionConfig.setUseSnapshotCompression(true); // 启用 Snappy 压缩
- 目前仅支持 Snappy 压缩算法
- 压缩在 Key-Group 层面进行
- 可以减少 Checkpoint 存储空间
并行 Checkpoints:
- 默认情况下只有一个 Checkpoint 可以运行
setMaxConcurrentCheckpoints(2)可以并行执行多个 Checkpoints- 会增加资源占用,但可以提升整体效率
3.5 状态后端配置
# MemoryStateBackend:状态存储在 JVM 堆内存,适合开发测试或小状态场景
state.backend: jobmanager
# FsStateBackend:状态存储在 JVM 堆内存,Checkpoint 时写入文件系统
# ⚠️ 注意:state.backend: filesystem 是 Flink 1.13 之前的旧写法
# 等同于 HashMapStateBackend,运行时数据在 JVM 堆内存
# state.backend: filesystem
# 启用增量 Checkpoint:只备份变化的数据块,大幅降低 Checkpoint 耗时
state.backend.incremental: true
# 保留的 Checkpoint 数量
state.checkpoints.num-retained: 3
# RocksDBStateBackend:状态存储在本地 RocksDB,适合大状态场景(>1GB)
# 突破内存限制:利用本地磁盘存储 TB 级状态数据
state.backend: rocksdb
# 启用增量 Checkpoint:只备份变化的数据块,极大降低 Checkpoint 耗时
state.backend.incremental: true
# RocksDB 本地状态存储路径(务必配置在 NVMe SSD 提升性能)
state.backend.rocksdb.localdir: /tmp/rocksdb
# 是否启用托管内存:Flink 自动管理 RocksDB 内存,建议保持 true
state.backend.rocksdb.memory.managed: true
3.6 RocksDB 配置
3.6.1 RocksDB 核心概念
RocksDB 是一个嵌入式、持久化、键值对(KV)存储引擎,可以想象成 “不仅存在于内存,还能高效溢写到硬盘上的超级 HashMap”。
| 特性 | 说明 |
|---|---|
| 嵌入式 | 不是独立数据库服务,而是 .jar 包,运行在 TaskManager 的 JVM 进程内部,无网络开销 |
| LSM-Tree 架构 | 写入极快(追加写),内存写满后 Flush 到磁盘,后台线程合并文件 |
| 突破内存限制 | 利用本地磁盘(推荐 SSD)存储 TB 级状态数据 |
| 增量 Checkpoint | 只备份变化的数据块,极大降低 Checkpoint 耗时 |
RocksDB 数据流转图:
flowchart TD
subgraph "JVM Heap"
MemTable["MemTable<br/>(内存写缓冲区)"]
end
subgraph "本地磁盘"
SST["SST Files<br/>(排序合并文件)"]
WAL["WAL<br/>(预写日志)"]
end
subgraph "远程存储"
Checkpoint["Checkpoint<br/>(增量备份到 HDFS/S3)"]
end
Data["Flink State"] -->|"写入"| WAL
WAL -->|"Flush"| MemTable
MemTable -->|"写满后 Flush"| SST
SST -->|"后台 Compaction"| Checkpoint
style MemTable fill:#e3f2fd
style SST fill:#fff3e0
style Checkpoint fill:#e8f5e9
3.6.2 RocksDB 配置参数速查
| 参数 | 默认值 | 说明与调优建议 |
|---|---|---|
state.backend.rocksdb.predefined-options | - | 官方配置模板,如 FLASH_SSD_OPTIMIZED 针对 SSD 优化 |
state.backend.rocksdb.localdir | - | 本地状态存储路径,务必配置在 NVMe SSD |
state.backend.rocksdb.checkpoint.transfer.thread.num | 1 | 异步上传线程数,建议调大到 4-8 以缩短 Checkpoint 耗时 |
state.backend.rocksdb.timer-service.factory | HEAP | 定时器存储位置,海量定时器场景设为 ROCKSDB 防止 OOM |
state.backend.rocksdb.memory.managed | true | 是否启用托管内存,建议保持 true 让 Flink 自动管理 |
state.backend.rocksdb.memory.fixed-per-slot | - | 每个 Slot 的固定 RocksDB 内存(不使用托管内存时配置) |
state.backend.rocksdb.memory.write-buffer-ratio | 0.5 | 写缓冲区占托管内存比例 |
state.backend.rocksdb.memory.high-pin-pinned-ratio | 0.1 | 锁定内存比例 |
3.6.3 生产环境推荐配置模板
# ===== RocksDB 优化配置 =====
state.backend: rocksdb
# 官方配置模板:FLASH_SSD_OPTIMIZED 针对 SSD 进行优化
state.backend.rocksdb.predefined-options: FLASH_SSD_OPTIMIZED
# 本地状态存储路径(务必使用 SSD,提升写入性能)
state.backend.rocksdb.localdir: /data/flink/rocksdb
# 异步上传线程数:建议调大到 4-8,缩短 Checkpoint 耗时
state.backend.rocksdb.checkpoint.transfer.thread.num: 4
# 定时器存入 RocksDB:海量定时器场景设为 ROCKSDB 防止 OOM
# 默认 HEAP 模式定时器存储在堆内存,数量过多会导致 JVM 内存压力
state.backend.rocksdb.timer-service.factory: ROCKSDB
# ===== 内存配置 =====
# 托管内存比例:使用 RocksDB 时建议保持 0.4 或调高到 0.5
taskmanager.memory.managed.fraction: 0.4
# JVM Metaspace:防止加载大量类时 OOM
taskmanager.memory.jvm-metaspace.size: 256mb
3.7 重启策略配置
# 策略:固定延迟重启 (还有 fixed-delay, fallback, none 等策略)
restart-strategy: fixed-delay
# 尝试重启次数:最大重试次数
restart-strategy.fixed-delay.attempts: 3
# 两次重启之间的等待时间:给外部系统(如 Kafka)恢复的时间
# 场景:针对网络抖动,不要让任务立即失败,尝试重试
restart-strategy.fixed-delay.delay: 10 s
3.8 常用运维命令
# 📤 提交作业
bin/flink run -p 4 path/to/myjob.jar
# 📋 查看运行中的作业
bin/flink list
# ⛔ 取消作业
bin/flink cancel <jobId>
# 💾 带 Savepoint 取消
bin/flink cancel -s hdfs://path/to/savepoint <jobId>
# 📸 触发 Savepoint
bin/flink savepoint <jobId> [targetDirectory]
# 🔄 从 Savepoint 恢复
bin/flink run -s <savepointPath> [runArgs]
# 🌐 查看 JobManager WebUI
# 默认地址: http://localhost:8081
3.9 常见问题排查
| 现象 | 可能原因 | 检查项 |
|---|---|---|
| RocksDB 频繁 Warn: “Stall” | 磁盘写入太慢 | 1. 检查 localdir 是否是 SSD2. 调大 setMaxBackgroundJobs3. 减小 Checkpoint 频率 |
| Checkpoint 超时/失败 | 状态过大 或 上传带宽不足 | 1. 检查是否开启增量 Checkpoint 2. 调大 checkpoint.transfer.thread.num |
| TaskManager OOM | 容器内存设置不合理 | 检查 taskmanager.memory.process.size 是否超过 Docker/K8s Limit |
| 反压 (Backpressure) 高 | 下游处理慢 或 网络堵塞 | 1. 检查 Sink 端(如 MySQL)连接池 2. 检查 taskmanager.memory.network 是否不够 |
参见: 01-应用场景 | 02-架构与原理 | 04-Java开发实践 | 附录
标签: #flink #运维 #部署 #配置 #大数据