📌 本文是「Flink 教程」系列第 二 篇 · 系列目录
2.1 Flink 架构
Flink 是典型的 Master-Slave (主从) 架构,类似于 ” 包工头与工人 ” 的关系。
2.1.1 集群组件
Flink 运行时由两种类型的进程组成:JobManager 和 TaskManager。
JobManager 组件:
| 组件 | 职责 |
|---|---|
| ResourceManager | 资源提供、回收、分配,管理 task slots |
| Dispatcher | 提供 REST 接口提交作业,运行 WebUI |
| JobMaster | 负责管理单个 JobGraph 的执行 |
TaskManager:
- 执行作业流的 task
- 缓存和交换数据流
- 资源调度的最小单位是 task slot
- 每个 slot 代表 TaskManager 资源的固定子集
2.1.2 Task Slot 和资源
Slot 数量:
- 每个 TaskManager 至少有一个 slot
- Slot 数量表示并发处理 task 的数量
- 一个 slot 中可以执行多个算子(通过算子链)
Slot 共享:
| 特性 | 说明 |
|---|---|
| 共享规则 | 同一作业的不同 Task 的 SubTask 可以共享 Slot |
| 最大共享 | 整个作业管道可以放在一个 Slot 中执行 |
| 核心优势 | 资源利用率更高,无需计算总 Task 数 |
Slot Sharing 示例 Source(1) → Map(1) → Sink(1),这 3 个 SubTask 可共享 1 个 Slot
2.1.3 核心架构详解
组件角色
| 组件 | 角色 | 职责 |
|---|---|---|
| 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)
Flink 集群架构图
graph TD
subgraph "客户端"
Client
end
subgraph "JobManager"
Dispatcher[Dispatcher<br/>WebUI]
JobMaster[JobMaster<br/>作业管理]
RM[ResourceManager<br/>资源管理]
end
subgraph "TaskManager 集群"
TM1["TaskManager 1<br/>Slot: ★ ★ ★"]
TM2["TaskManager 2<br/>Slot: ★ ★ ★"]
TMn["TaskManager N<br/>Slot: ★ ★ ★"]
end
Client --> Dispatcher
Dispatcher --> JobMaster
JobMaster --> RM
RM <-->|"Akka RPC<br/>控制流"| TM1
RM <-->|"Akka RPC<br/>控制流"| TM2
RM <-->|"Akka RPC<br/>控制流"| TMn
TM1 <-->|"Netty<br/>数据交换"| TM2
TM2 <-->|"Netty<br/>数据交换"| TMn
TM1 <-->|"Netty<br/>数据交换"| TMn
subgraph "TM1 内部"
S1[Slot 1: Source→Map→Sink]
S2[Slot 2: Filter]
S3[Slot 3: Agg]
end
2.1.4 集群模式
| 模式 | 生命周期 | 资源隔离 | 适用场景 |
|---|---|---|---|
| Session 集群 | 长期运行,作业共享 | 资源竞争 | 交互式分析、短作业 |
| Job 集群 | 作业启动时创建 | 作业隔离 | 大型作业、长时运行 |
| Application 集群 | 随应用启动 | 最佳隔离 | 云原生部署 |
2.2 有状态流处理
有状态流处理是 Flink 的核心概念,指在处理事件流时维护和更新状态信息。
2.2.1 什么是状态
在数据流中,许多操作(如窗口操作)需要记住跨多个事件的信息,这些操作称为有状态操作。
状态类型:
// Keyed State - 按 Key 分区的状态
ValueState<T> // 单值状态,存储单个值
ListState<T> // 列表状态,存储元素列表
MapState<K, V> // 映射状态,存储键值对
ReducingState<T> // 聚合状态,自动聚合
AggregatingState<IN, OUT> // 聚合状态,自定义聚合
// Operator State - Operator 级别的状态(无 Key 分组)
Keyed State 的特点:
- 仅在 keyed streams 上可用(即 keyBy 操作之后)
- 状态与特定的 key 绑定
- 访问仅限于当前事件的 key
- 状态更新是本地操作,保证一致性
2.2.2 Checkpoint 机制
Flink 使用检查点和流重播的组合来实现容错。
Checkpoint 原理:
- 在每个输入流的特定点标记检查点
- 同时保存每个操作符的对应状态
- 从检查点恢复时,重置状态并从该点重放记录
Barrier(屏障):
- 屏障注入数据流,作为数据的一部分流动
- 屏障不会超过记录,它们严格按顺序流动
- 一个屏障分隔当前快照和下一个快照的记录
Checkpoint 流程:
- 屏障从源头注入(记录在 Kafka 中的偏移量)
- 屏障向下游流动
- 中间操作符从所有输入流收到屏障后,发送屏障到输出流
- Sink 操作符收到所有输入的屏障后,确认检查点完成
Checkpoint 屏障流动图
sequenceDiagram
participant Source
participant Map
participant KeyBy
participant Sink
Note over Source: Barrier n 注入<br/>记录 Kafka 偏移量
Source->>Map: 数据 + Barrier n
Map->>KeyBy: 数据 + Barrier n
KeyBy->>Sink: 数据 + Barrier n
Note over Sink: 收到所有输入 Barrier n<br/>确认 Checkpoint 完成 ✓
Checkpoint vs Savepoint 对比
| 特性 | Checkpoint | Savepoint |
|---|---|---|
| 触发 | 自动(周期性) | 手动 |
| 生命周期 | 任务停止即删(默认) | 永久保留 |
| 用途 | 活下来(防挂) | 活得更好(升级) |
2.2.3 非对齐 Checkpoint
Flink 1.11 引入了非对齐 Checkpoint,基本思想是检查点可以超过所有飞行中的数据,只要飞行中的数据成为操作符状态的一部分。
特点:
- 确保屏障以最快速度到达 sink
- 特别适合慢速数据路径的场景
- 避免了长时间的对齐等待
非对齐 Checkpoint:解决严重反压导致 Checkpoint 超时的问题,Barrier 插队优先处理,但需将缓冲区数据一起存盘
2.2.4 内存管理与反压
TaskManager 内存估算
| 方式 | 推荐 | 风险 |
|---|---|---|
| 增量聚合 | 使用 Reduce/Aggregate Function | 内存只存累加结果,占用极小 |
| 全量缓存 | 使用 ProcessWindowFunction | 内存存原始数据 List,极大可能 OOM |
估算公式:Key 数量 × 状态对象大小(需考虑 Java 对象头开销)
反压 vs Watermark
| 对比项 | 反压 | Watermark 不一致 |
|---|---|---|
| 本质 | 物理现象,下游处理慢 | 逻辑现象,水位线滞后 |
| 表现 | 缓冲区填满,堵塞上游 | 窗口无法关闭,状态堆积 |
| 传播方向 | 下游 → 上游(反向) | 上游 → 下游(正向) |
关系图示:
flowchart TD
subgraph "正常情况"
A1[Source] -->|数据| B1[Map]
B1 -->|数据| C1[KeyBy]
C1 -->|数据| D1[Sink]
D1 -.->|"✓ 处理正常"| B1
end
subgraph "反压场景"
A2[Source] -->|数据| B2[Map]
B2 -->|数据| C2[KeyBy]
C2 -->|"⚠️ 缓冲区满"| D2[Sink]
D2 -.->|"🚫 反压 Backpressure"| C2
C2 -.->|"🚫 反压"| B2
B2 -.->|"🚫 反压"| A2
end
style D2 fill:#ffcccc
style C2 fill:#ffdddd
style B2 fill:#ffeeee
核心关系:
- 反压 → Watermark:反压会导致缓冲区积压,Watermark 传递变慢
- Watermark → 反压:Watermark 滞后导致状态无限堆积,最终引发反压
2.3 时间语义与水位线
2.3.1 时间概念
Flink 支持三种时间概念:
| 时间类型 | 说明 | 特点 |
|---|---|---|
| 事件时间(Event Time) | 事件实际发生的时间 | 结果可重现,与处理时间无关 |
| 处理时间(Processing Time) | 系统处理事件的时间 | 简单,但结果不可重现 |
| 摄入时间(Ingestion Time) | Flink 读取事件的时间 | 介于两者之间 |
处理时间 vs 事件时间:
处理时间是最简单的时间概念,不需要协调流和机器之间的协调,提供最佳性能和最低延迟。但在分布式和异步环境中,处理时间不提供确定性,因为它受记录到达系统的速度影响。
事件时间处理会产生完全一致和确定性的结果,无论事件何时到达或如何排序。然而,除非已知事件按顺序到达,否则事件时间处理会在等待乱序事件时产生一些延迟。
2.3.2 水位线(Watermark)
水位线是 Flink 中测量事件时间进度的机制。
水位线的作用:
- 声明事件时间到达某个时间点
- 标记该时间之前的事件应该都已到达
- 触发窗口关闭和计算
水位线策略:
WatermarkStrategy<Event> strategy = WatermarkStrategy
.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(20)) // 最大无序边界
.withTimestampAssigner((event, timestamp) -> event.timestamp);
DataStream<Event> withWatermarks = stream.assignTimestampsAndWatermarks(strategy);
Watermark(t) 声明:
- 声明事件时间已到达 t
- 不应有更多时间戳
t' <= t的元素 - 即时间戳小于等于水印的事件都应该已到达
乱序流中的水位线:
- 一般来说,水位线是对某个时间点之前所有事件都已到达的声明
- 一旦水位线到达操作符,操作符可以将其内部事件时间时钟推进到水位线的值
Watermark 多并行度细节
- Watermark 在 Source Operator 中生成,并且在每个 Source Operator 的子 Task 中都会独立生成 Watermark
- Watermark 会逐步更新下游算子中的 Watermark 水位线
- 如果多个 Watermark 同时更新一个算子 Task 的当前事件时间,Flink 会选择最小的水位线来更新
- 当一个 Window 算子 Task 中水位线大于了 Window 结束时间,就会立即触发窗口计算
- 慢分区会拖慢整体进度:如果某个分区 Watermark 落后太多,可能是该分区数据积压或处理慢,建议监控各并行度的 Watermark 进度
水位线核心公式
- 公式:
Watermark = MaxEventTime - MaxOutOfOrderness(最大允许乱序时间) - 传播:多对一场景下遵循木桶效应,下游取所有上游输入中最小的 Watermark
木桶效应图示:
flowchart LR
subgraph "上游 Source (并行度=3)"
P1["Partition 1<br/>WM: 12:05"]
P2["Partition 2<br/>WM: 11:58"]
P3["Partition 3<br/>WM: 12:02"]
end
subgraph "下游 Window"
W["Window<br/>取最小值: 11:58"]
end
P1 -->|"WM: 12:05"| W
P2 -->|"WM: 11:58<br/>⚠️ 最慢分区"| W
P3 -->|"WM: 12:02"| W
style P2 fill:#ffcccc,stroke:#ff0000
style W fill:#ffdddd,stroke:#ff0000
公式详解:
Watermark = 当前分区最大事件时间 - 最大允许乱序时间
例如:最大事件时间 12:05,最大无序边界 5 分钟
Watermark = 12:05 - 5min = 12:00
→ 声明:12:00 之前的数据都已到达
2.3.3 延迟处理
允许的延迟(Allowed Lateness):
- 延迟事件是那些在水位线之后到达的事件
- 可以配置允许延迟的时间窗口
- 在延迟窗口内,延迟事件仍会被处理
侧输出(Side Output):
- 将被丢弃的延迟事件发送到侧输出流
- 用于后续单独处理延迟数据
迟到数据处理(三道防线)
- Watermark 自带缓冲:设置
MaxOutOfOrderness(如 5s) - AllowedLateness(窗口宽容期):窗口触发后保留状态一段时间(如 1min),迟到数据会导致窗口再次触发
- Side Output(侧输出流):彻底迟到的数据兜底存入侧输出流,人工处理
2.4 状态后端(State Backends)
常见误区
state.backend: filesystem是 Flink 1.13 之前的旧写法,等同于 HashMapStateBackend,运行时数据在 JVM 堆内存,Checkpoint 时才写入文件系统。
Flink 提供三种状态后端:
| 状态后端 | 存储位置 | 适用场景 |
|---|---|---|
| MemoryStateBackend | JVM 堆内存 | 开发测试、小状态 |
| FsStateBackend | 文件系统 + 堆内存 | 中等状态、生产环境 |
| RocksDBStateBackend | RocksDB + 文件系统 | 大状态、生产环境 |
2.4.1 状态后端对比详解
| 类型 | 场景 | 存储位置 | 配置 |
|---|---|---|---|
| HashMapStateBackend | 低延迟、状态量可控(< 几 GB) | TM 堆内存 | state.backend: hashmap |
| EmbeddedRocksDBStateBackend | 超大状态(TB 级)、长窗口 | 本地磁盘(RocksDB)+ 内存缓存 | state.backend: rocksdb |
💡 Tips: RocksDB JNI 传输限制(生产环境重要)
- RocksDB 通过 JNI 的方式进行数据交互
- 每次能够传输的最大数据量为 2^31 字节(约 2GB)
- 在 RocksDBStateBackend 合并的状态数据量大小不能超过此限制
- 否则将会导致状态数据无法同步,这是 RocksDB 采用 JNI 方式的限制
- 用户在使用过程中应当注意,避免单条状态记录超过 2GB
状态后端选择建议:
| 场景 | 推荐后端 | 原因 |
|---|---|---|
| 开发调试 | MemoryStateBackend | 快速、简单 |
| 小状态(<1GB) | FsStateBackend | 性能好 |
| 大状态(>1GB) | RocksDBStateBackend | 容量大 |
| 状态 > 内存 | RocksDBStateBackend | 磁盘存储 |
2.5 Savepoint
Savepoint 是手动触发的检查点,用于:
- 作业停止和恢复
- 应用程序升级
- 集群迁移
- A/B 测试
特点:
- 手动触发,不会自动过期
- 依赖常规检查点机制
- 从 savepoint 恢复时进行精确的状态恢复
Savepoint 要点:
- 必须为有状态算子设置 UID(
.uid("name")),否则代码修改后状态无法恢复
2.6 窗口机制
2.6.1 窗口触发时间
窗口触发条件:Watermark ≥ WindowEndTime
窗口触发时间线:
flowchart LR
subgraph "数据到达阶段"
A1((12:00<br/>窗口开启))
A2(12:02<br/>事件A)
A3(12:05<br/>事件B)
A4(12:08<br/>事件C)
end
subgraph "等待水位线阶段"
B1((12:10<br/>窗口关闭时间点))
B2{等待<br/>Watermark}
end
subgraph "窗口触发阶段"
C1((12:15<br/>Watermark>=12:10<br/>窗口计算触发))
end
A1 --> A2 --> A3 --> A4 --> B1 --> B2 --> C1
style A1 fill:#e1f5fe
style C1 fill:#c8e6c9
style B2 fill:#fff3e0
关键公式:
看到结果时刻 = 窗口结束时间 + Watermark(最大无序边界)
示例:窗口 [12:00, 12:10),MaxOutOfOrderness=5min
必须等到 Watermark >= 12:10
即:12:05 之后的数据最大时间戳 >= 12:05
最早 12:10 看到结果,最迟 12:15 看到结果
2.6.2 窗口类型
| 类型 | 特点 | 风险 |
|---|---|---|
| 滚动窗口(Tumbling) | 数据只属于 1 个窗口 | - |
| 滑动窗口(Sliding) | 数据会属于 N 个 窗口 | 资源消耗大(数据逻辑膨胀) |
2.6.3 算子与窗口的关系
- 并非 1 个窗口占用 1 个算子
- 1 个 WindowOperator (SubTask) 管理着该 Slot 内所有 Key 的成千上万个窗口对象
- Slot 之间隔离,但 Slot 内部的窗口共享内存(大户 Key 会导致 OOM)
2.7 动态表
Flink 的 Table API 和 SQL 基于动态表概念工作。
核心原理:
- 流可以转换为表(Changelog 流上的查询)
- 表可以转换回流(输出)
- 支持连续查询(Continuous Query)
流与表的转换:
流 --(asTable)--> 动态表 --(toChangelogStream)--> 流
Append-Only Stream <--> 动态表 <--> Retract Stream
Upsert Stream <--> 动态表 <--> Upsert Stream
参见: 01-应用场景 | 03-部署与运维 | 04-Java开发实践 | 附录
标签: #flink #架构 #原理 #大数据 #流处理