📌 本文是「Flink 教程」系列第 五 篇 · 系列目录
文档说明 Flink 学习笔记核心总结,包含架构解析、内存配置、状态管理和故障恢复等关键知识点。
1. 核心架构 (Architecture)
Flink 是典型的 Master-Slave (主从) 架构,类似于”包工头与工人”的关系。
组件角色
| 组件 | 角色 | 职责 |
|---|---|---|
| JobManager | Master/包工头 | 作业调度、Checkpoint 协调、故障恢复、管理 Resource Manager 和 Dispatcher |
| TaskManager | Worker/工人 | 实际执行计算的 JVM 进程,每个 TM 包含一定数量的 Slot |
| Slot | 工位 | TM 资源的固定子集(主要是内存隔离,CPU 共享),决定并行执行数 |
| SubTask | 子任务 | 算子根据并行度复制出来的物理执行单元(线程) |
核心机制
- Slot Sharing (槽位共享): 允许同一个 Job 的不同 Task 的 SubTask 共享同一个 Slot(如 Source + Map + Sink 可在一个 Slot 跑),提高资源利用率。
- 通信:
- 控制流: Akka (RPC) —— JM 与 TM 之间指令交互
- 数据流: Netty —— 跨 TM 的数据传输(Shuffle)
2. 状态后端 (State Backends)
常见误区
state.backend: filesystem是 Flink 1.13 之前的旧写法,等同于 HashMapStateBackend,运行时数据在 JVM 堆内存,Checkpoint 时才写入文件系统。
对比
| 类型 | 场景 | 存储位置 | 配置 |
|---|---|---|---|
| HashMapStateBackend | 低延迟、状态量可控(< 几 GB) | TM 堆内存 | state.backend: hashmap |
| EmbeddedRocksDBStateBackend | 超大状态(TB 级)、长窗口 | 本地磁盘 (RocksDB) + 内存缓存 | state.backend: rocksdb |
3. 容错机制 (Checkpoint & Savepoint)
Checkpoint vs Savepoint
| 特性 | Checkpoint | Savepoint |
|---|---|---|
| 触发 | 自动 (周期性) | 手动 |
| 生命周期 | 任务停止即删 (默认) | 永久保留 |
| 用途 | 活下来 (防挂) | 活得更好 (升级) |
Checkpoint 原理
- 算法: Chandy-Lamport 算法 (Barrier 对齐)
- 流程: JM 注入 Barrier → Source 随数据流广播 → 算子收到 Barrier 后对状态做快照并异步上传 HDFS
- 非对齐 Checkpoint: 解决严重反压导致 Checkpoint 超时的问题,Barrier 插队优先处理,但需将缓冲区数据一起存盘
Savepoint 要点
- 必须为有状态算子设置 UID (
.uid("name")),否则代码修改后状态无法恢复
4. 水位线 (Watermark)
核心定义
Watermark 是解决乱序数据与结果确定性矛盾的机制,是一种特殊标记,意味着:“小于该时间的数据都已到齐”。
- 公式:
Watermark = MaxEventTime - MaxOutOfOrderness(最大允许乱序时间) - 传播: 多对一场景下遵循木桶效应,下游取所有上游输入中最小的 Watermark
迟到数据处理(三道防线)
- Watermark 自带缓冲: 设置
MaxOutOfOrderness(如 5s) - AllowedLateness (窗口宽容期): 窗口触发后保留状态一段时间(如 1min),迟到数据会导致窗口再次触发
- Side Output (侧输出流): 彻底迟到的数据兜底存入侧输出流,人工处理
5. 窗口机制 (Windows)
触发时间
窗口触发条件:Watermark ≥ WindowEndTime
看到结果的时刻 =
窗口结束时间 + Watermark 延迟配置例:
[12:00, 12:10)窗口,延迟 5min,必须等到 12:15 的数据到达时窗口才关闭
窗口类型
| 类型 | 特点 | 风险 |
|---|---|---|
| 滚动窗口 (Tumbling) | 数据只属于 1 个窗口 | - |
| 滑动窗口 (Sliding) | 数据会属于 N 个 窗口 | 资源消耗大(数据逻辑膨胀) |
算子与窗口的关系
- 并非 1 个窗口占用 1 个算子
- 1 个 WindowOperator (SubTask) 管理着该 Slot 内所有 Key 的成千上万个窗口对象
- Slot 之间隔离,但 Slot 内部的窗口共享内存(大户 Key 会导致 OOM)
6. 内存管理与反压
TaskManager 内存估算
| 方式 | 推荐 | 风险 |
|---|---|---|
| 增量聚合 | 使用 Reduce/Aggregate Function | 内存只存累加结果,占用极小 |
| 全量缓存 | 使用 ProcessWindowFunction | 内存存原始数据 List,极大可能 OOM |
估算公式: Key 数量 × 状态对象大小(需考虑 Java 对象头开销)
反压 vs Watermark
- 反压: 物理现象,下游处理慢导致缓冲区填满,堵塞上游
- Watermark 不一致: 逻辑现象,某个 Slot 水位线滞后导致下游窗口无法关闭,状态无限堆积
- 关系: 反压会导致 Watermark 传递变慢;Watermark 导致的状态爆炸也会引发反压
7. 生产部署建议
- 部署模式: 生产环境首选 Application Mode(JM 在集群端生成图)
- UID 设置: 对所有 Stateful Operator 显式调用
.uid("...") - RocksDB: 如果 Key 量级无法预估或很大,务必开启 RocksDB StateBackend
- Side Output: 在涉及金额/审计的场景,务必开启侧输出流处理迟到数据
8. 常见问题排查
| 现象 | 可能原因 | 检查项 |
|---|---|---|
| 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 是否不够 |
9. 关键配置速查
Cluster 资源配置 (flink-conf.yaml)
# JobManager 配置
jobmanager.memory.process.size: 2g
# TaskManager 配置 (单容器 8GB 内存 / 4 CPU)
taskmanager.memory.process.size: 8g
taskmanager.memory.managed.fraction: 0.4 # RocksDB 建议 0.4-0.5
taskmanager.memory.jvm-metaspace.size: 256m
taskmanager.memory.network.fraction: 0.1
taskmanager.numberOfTaskSlots: 4
parallelism.default: 4
容错与高可用配置
# Checkpoint 设置
execution.checkpointing.interval: 180000 # 每 3 分钟
state.checkpoints.dir: hdfs://namenode:8020/flink/checkpoints
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION
execution.checkpointing.max-concurrent-checkpoints: 1
# 重启策略
restart-strategy: fixed-delay
restart-strategy.fixed-delay.attempts: 3
restart-strategy.fixed-delay.delay: 10 s
RocksDB 优化配置
# 官方模板 (SSD 优化)
state.backend.rocksdb.predefined-options: FLASH_SSD_OPTIMIZED
# 本地状态存储路径 (务必配置在 NVMe SSD)
state.backend.rocksdb.localdir: /path/to/ssd
# 异步上传线程数 (建议 4-8)
state.backend.rocksdb.checkpoint.transfer.thread.num: 4
# 定时器存储位置 (海量定时器设为 RocksDB)
state.backend.rocksdb.timer-service.factory: ROCKSDB
💡 提示: 查看 Flink配置相关 获取完整配置模板。
标签: #flink #架构 #疑问解答 #FAQ #核心概念