flink的架构
2026/9/4 22:45:16 网站建设 项目流程

Flink 是一个面向有状态计算的分布式流处理框架,也支持批处理。它的架构可以从“运行时组件、数据流模型、容错机制”三个层面理解。

一、整体架构

一个典型的 Flink 集群包含以下组件:

客户端 Client | | 提交 Job v 作业管理器 JobManager | | 调度任务、协调检查点、管理元数据 v 任务管理器 TaskManager 集群 | | 执行算子、处理数据、保存本地状态 v 外部系统:Kafka / MySQL / Elasticsearch / HDFS / OSS

1. Client

Client 负责:

  • 解析用户编写的 Flink 程序
  • 构建 JobGraph
  • 将作业提交给 JobManager
  • 提交后可以退出,作业仍由集群继续运行

Client 通常不是长期运行时组件。

2. JobManager

JobManager 是作业的控制中心,主要负责:

  • 接收和管理作业
  • 将作业转换为可执行任务
  • 调度任务到 TaskManager
  • 管理 Checkpoint 和 Savepoint
  • 处理故障恢复
  • 管理作业状态和资源需求

在现代 Flink 架构中,JobManager 内部通常可以理解为几个逻辑角色:

  • Dispatcher:接收 REST、命令行等提交请求,并启动 JobMaster
  • JobMaster:负责单个作业的调度和协调
  • ResourceManager:管理集群资源,并向外部资源管理系统申请资源

3. TaskManager

TaskManager 是真正执行计算任务的工作节点,主要负责:

  • 执行算子逻辑
  • 处理输入和输出数据
  • 管理网络数据交换
  • 保存算子状态
  • 向 JobManager 汇报任务状态

一个 TaskManager 中可以有多个Task Slot。Slot 是 Flink 对 TaskManager 资源进行逻辑隔离和分配的单位。

二、Flink 的作业执行模型

用户程序通常经过以下转换:

用户代码 -> StreamGraph -> JobGraph -> ExecutionGraph -> TaskManager 上的执行任务

1. StreamGraph

StreamGraph 是对用户逻辑的直接描述。例如:

Source -> Map -> Filter -> Sink

它包含 Source、Transformation、Sink 等节点。

2. JobGraph

JobGraph 是提交给 JobManager 的作业描述。Flink 会在这一阶段进行算子链合并等优化,把可以连续执行的算子合并为一个 JobVertex。

例如:

Source -> Map -> Filter

如果它们之间没有发生数据重分区,可能会被合并到同一个 Operator Chain 中,从而减少线程切换和网络通信。

3. ExecutionGraph

ExecutionGraph 是 JobManager 根据并行度和资源情况生成的实际执行计划。它会把 JobGraph 中的逻辑节点展开为多个并行实例。

例如并行度为 3:

Map 算子 Map-0 Map-1 Map-2

4. Operator Chain

多个连续算子可以在同一个线程中执行,形成 Operator Chain:

Source -> Map -> Filter

这样可以避免不必要的序列化、网络传输和线程切换,提升性能。

如果两个算子之间存在 Shuffle、KeyBy、广播或并行度变化,通常不能直接链在一起。

三、Task、SubTask 和 Slot

假设一个作业如下:

Source -> Map -> KeyBy -> Reduce -> Sink

并行度为 2 时,每个算子会产生两个 SubTask:

Source-0 -> Map-0 -> Reduce-0 -> Sink-0 Source-1 -> Map-1 -> Reduce-1 -> Sink-1

其中:

  • Operator:算子的逻辑定义
  • SubTask:算子按照并行度展开后的一个实例
  • Task:一个或多个 Operator Chain 的执行单元
  • Task Slot:TaskManager 提供的逻辑资源槽位

需要注意,Slot 主要用于资源调度和隔离,并不等于一个独立线程,也不一定对应一个 SubTask。

四、数据流模型

Flink 的核心是连续不断的数据流。即使处理有限数据集,也可以看作一个有界流。

Source -> Transformation -> Sink

常见组件如下:

  • Source:从 Kafka、文件、数据库等读取数据
  • Transformation:执行 Map、Filter、Join、Window 等计算
  • Sink:将结果写入数据库、消息队列或文件系统

KeyBy会按照 Key 对数据进行逻辑分区,使同一个 Key 的数据进入同一个下游并行实例:

相同 userId 的数据 | v 同一个 Reduce SubTask

这为按 Key 维护状态提供了基础。

五、状态管理

Flink 与普通无状态计算框架的一个重要区别是:它原生支持大规模、有一致性保障的状态管理。

状态通常分为:

  • Keyed State:绑定到某个 Key,例如每个用户的累计金额
  • Operator State:绑定到算子实例,例如 Kafka Source 的分区消费进度

Flink 的状态可以存储在:

  • 堆内存
  • RocksDB 等嵌入式状态后端
  • 其他状态后端实现

状态后端负责状态的实际存储,而 Checkpoint 负责定期持久化和恢复状态。

六、Checkpoint 容错机制

Flink 使用基于 Chandy-Lamport 思想的分布式快照机制实现一致性 Checkpoint。

简化流程如下:

JobManager 发出 Checkpoint Barrier | v Source 注入 Barrier | v Barrier 随数据流向下游传播 | v 各算子保存自己的状态 | v Checkpoint 完成并写入持久化存储

Checkpoint 保存的内容通常包括:

  • 算子状态
  • Source 消费位点
  • 其他恢复所需的运行时信息

任务失败后,Flink 可以:

  1. 从最近一次成功的 Checkpoint 恢复状态
  2. 恢复 Source 的消费位置
  3. 重新执行之后的数据
  4. 继续处理作业

常见语义包括:

  • At-most-once:最多一次,可能丢数据
  • At-least-once:至少一次,可能重复数据
  • Exactly-once:端到端恰好一次,但需要 Source、Flink 和 Sink 协同支持

七、网络与数据交换

TaskManager 之间通过网络传输数据。典型的数据交换方式包括:

  • Forward:上游一个 SubTask 对应下游一个 SubTask
  • Rebalance:轮询分发数据
  • Broadcast:发送给所有下游实例
  • KeyBy / Hash Partition:按 Key 哈希分区
  • Rescale:在本地范围内重新分配数据

keyBy()往往会产生网络 Shuffle,是数据流发生重新分区的关键位置。

TaskManager 内部通常通过网络缓冲区、Result Partition 和 Input Gate 组织数据传输。

八、部署模式

Flink 常见的部署方式有:

1. Session Cluster

预先启动一个 Flink 集群,多个作业共享这个集群。

优点:启动快,资源可以复用。

缺点:作业之间可能互相影响,资源隔离较弱。

2. Per-Job Cluster

每个作业单独创建一个集群,作业结束后集群释放。

优点:作业隔离较好。

缺点:启动和资源申请成本较高。

3. Application Mode

将应用程序和 Flink 集群一起部署,由集群直接执行应用入口逻辑。

常见运行环境包括:

  • Standalone
  • YARN
  • Kubernetes

在 Kubernetes 环境中,JobManager 通常运行在一个或多个 Pod 中,TaskManager 以 Pod 形式按需扩缩容。

九、Flink 与传统 Lambda 架构的区别

传统 Lambda 架构一般将系统拆成:

Batch Layer + Speed Layer + Serving Layer

同一份数据可能需要分别走批处理和流处理链路。

Flink 更强调统一的流批处理模型:

有界流 + 无界流
  • 无界流适合实时计算
  • 有界流适合批处理
  • 两者可以共享大量 API、状态和执行机制

十、一次作业运行的完整过程

可以把 Flink 作业运行概括为:

1. 用户编写 DataStream 或 Table/SQL 程序 2. Client 生成作业图 3. Client 将作业提交给 JobManager 4. JobManager 申请资源并生成执行计划 5. TaskManager 启动各个 SubTask 6. Source 持续读取数据 7. 数据经过算子链和网络 Shuffle 流转 8. 算子维护本地状态 9. JobManager 定期触发 Checkpoint 10. 发生故障时从 Checkpoint 恢复 11. 结果通过 Sink 写入外部系统

一句话总结:JobManager 负责控制和调度,TaskManager 负责执行,Operator 描述计算逻辑,State 保存中间结果,Checkpoint 提供一致性容错,网络 Shuffle 负责并行实例之间的数据交换。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询