跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
🌊 Flink 教程 · 3 / 7 查看系列简介 →

Flink 教程(三):部署与运维

文章目录
  1. 3.1 部署模式
  2. 3.1.1 独立集群(Standalone)
  3. 3.1.2 YARN 集群
  4. 3.1.3 Kubernetes
  5. 3.2 高可用配置
  6. 3.2.1 ZooKeeper HA
  7. 3.2.2 Kubernetes HA
  8. 3.3 内存配置
  9. 3.3.1 TaskManager 内存
  10. 3.3.2 生产环境内存配置参考
  11. 3.4 Checkpoint 配置
  12. 3.5 状态后端配置
  13. 3.6 RocksDB 配置
  14. 3.6.1 RocksDB 核心概念
  15. 3.6.2 RocksDB 配置参数速查
  16. 3.6.3 生产环境推荐配置模板
  17. 3.7 重启策略配置
  18. 3.8 常用运维命令
  19. 3.9 常见问题排查

📌 本文是「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 HAKubernetes HA
依赖需部署 Zookeeper 集群依赖 Kubernetes
存储后端HDFSS3/NFS/OSS
适用部署Standalone / YARNKubernetes
元数据恢复从 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.num1异步上传线程数,建议调大到 4-8 以缩短 Checkpoint 耗时
state.backend.rocksdb.timer-service.factoryHEAP定时器存储位置,海量定时器场景设为 ROCKSDB 防止 OOM
state.backend.rocksdb.memory.managedtrue是否启用托管内存,建议保持 true 让 Flink 自动管理
state.backend.rocksdb.memory.fixed-per-slot-每个 Slot 的固定 RocksDB 内存(不使用托管内存时配置)
state.backend.rocksdb.memory.write-buffer-ratio0.5写缓冲区占托管内存比例
state.backend.rocksdb.memory.high-pin-pinned-ratio0.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 是否是 SSD
2. 调大 setMaxBackgroundJobs
3. 减小 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 #运维 #部署 #配置 #大数据


RAG 智能问答

针对本文继续提问:《Flink 教程(三):部署与运维》