精选·Flink

Flink 的整体架构与核心角色:JobManager、TaskManager 各干什么

91学AI·2026/7/13·8 阅读

考察点

这是 Flink 面试的开场题,几乎必问。面试官不只想听你背组件名,而是想确认你是否理解「一个作业提交之后到底发生了什么」:谁负责调度、谁真正跑数据、资源怎么切分。追问通常会往 Task Slot 共享、算子链、高可用方向走。

参考答案

总体架构:主从结构

Flink 集群是典型的 master-worker 结构。客户端提交作业后,JobManager 是大脑,负责把 JobGraph 转换成 ExecutionGraph、调度 Task、协调 Checkpoint、故障恢复;TaskManager 是干活的,真正执行数据流算子、管理内存和网络缓冲、上报心跳。中间还有一层:Dispatcher 接收作业提交、拉起 JobManager(Application/Per-Job 模式下每个作业一个);ResourceManager 负责资源的分配与回收,对接 YARN、Kubernetes 这类资源框架;JobMaster 是 JobManager 里管单个作业的核心组件。

生产上最常用的是 Application Mode 跑在 Kubernetes 或 YARN 上:用户代码的 main 方法在 JobManager 侧执行,避免客户端成为瓶颈,每个作业有独立的 JobManager,故障隔离好。

Task Slot:资源切分的单位

TaskManager 的资源被切成若干 Slot,一个 Slot 是一份固定的资源份额(内存均分,CPU 不隔离)。比如一个 4 核 16G 的 TaskManager 配 4 个 Slot,每个 Slot 拿 4G 堆内存的均分份额。Slot 是调度的最小单位,一个 Task(subtask)跑在一个 Slot 里。

关键点在于 Slot 共享:默认情况下,同一个作业的不同算子的 subtask 可以共享一个 Slot。这意味着一个 Slot 里可能跑一条完整的流水线(source → map → keyBy → sink 各一个实例),好处是资源利用率高、Slot 数量直接等于并行度需求,不用为每个算子单独算资源。可以用 SlotSharingGroup 把某些算子隔离出去,比如一个特别吃内存的算子单独占 Slot。

算子链(Operator Chain)

上下游算子并行度相同、且是一对一转发(forward,没有 shuffle)时,Flink 会把它们链在一起,在同一个线程里以方法调用的方式传递数据,省掉序列化和网络开销。keyBy 之后链就断了。算子链在 Web UI 上体现为一个框,排查性能问题时要会看——两个本该分开的算子链在一起互相拖慢时,可以用 disableChaining()startNewChain() 拆开。

一次提交发生了什么

客户端把代码打成 JobGraph → 提交给 Dispatcher → JobMaster 生成 ExecutionGraph(每个算子按并行度展开成 subtask)→ 向 ResourceManager 申请 Slot → 部署到 TaskManager → 开始跑,JobManager 周期性触发 Checkpoint。这套链路讲顺了,这道题就过了。

可能的追问

  • 一个 Slot 能跑几个 Task?——Slot 是资源单位不是线程容器,通过共享机制一个 Slot 可以跑同一作业多个不同算子的 subtask,但同一算子的两个并行实例不会进同一 Slot。
  • JobManager 挂了怎么办?——需要配高可用(ZooKeeper 存元数据、选举 standby JobManager),JobManager 故障期间作业不挂,恢复后接着调度;Checkpoint 协调会中断,恢复后从最近完成的 Checkpoint 继续。
  • 并行度和 Slot 的关系?——作业需要的 Slot 数默认等于最大并行度(因为共享),比如并行度 8 就需要 8 个 Slot,可以分布在多个 TaskManager 上。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.