欢迎光临
我们一直在努力

Flink运行时架构详解

Flink运行时架构详解:JobManager、TaskManager、Slot与并行度

前言

学Flink,运行时架构是必须搞清楚的基础。只有真正理解了JobManager、TaskManager、Slot、并行度这些概念之间的关系,后面学任务提交、状态管理、Checkpoint才不会一头雾水。本文基于尚硅谷Flink1.17教程,系统梳理Flink运行时架构的核心知识点。


一、系统架构

Flink集群由两类进程组成:JobManager(主进程)和TaskManager(工作进程)。

1.1 JobManager

JobManager是集群中任务管理和调度的核心,控制应用的执行。每个应用都有一个唯一对应的JobManager。

JobManager内部包含三个组件:

(1)JobMaster

JobMaster是JobManager中最核心的组件,与具体的Job一一对应。多个Job同时运行时,每个Job都有自己的JobMaster。

JobMaster的工作流程:

  • 接收提交的应用(JobGraph)
  • 将JobGraph转换为执行图(ExecutionGraph)——包含所有可以并发执行的任务
  • 向ResourceManager申请执行所需的资源(Slot)
  • 获取到足够资源后,将执行图分发给TaskManager执行
  • 运行过程中负责Checkpoint等中央协调操作
  • 注意:早期版本的Flink中没有JobMaster的概念,当时的JobManager实际上就是现在的JobMaster。

    (2)ResourceManager

    ResourceManager负责集群中资源的分配和管理,整个集群中只有一个。

    这里的"资源"主要指TaskManager的任务槽(Task Slots)。每一个Task都需要分配到一个Slot上执行。

    注意:要区分Flink内置的ResourceManager和YARN等资源管理平台的ResourceManager,它们是不同层面的概念。

    (3)Dispatcher

    Dispatcher的职责:

    • 提供REST接口,用于提交应用
    • 为每个新提交的作业启动一个新的JobMaster组件
    • 启动Web UI,展示和监控作业执行信息

    Dispatcher并不是必需组件,在不同的部署模式下可能被忽略。

    1.2 TaskManager

    TaskManager是Flink的工作进程,负责数据流的具体计算。

    • 集群中至少需要一个TaskManager
    • 每个TaskManager包含一定数量的任务槽(Task Slots)
    • Slot数量限制了TaskManager能并行处理的任务数量

    TaskManager的工作流程:

  • 启动后向ResourceManager注册自己的Slot
  • 收到ResourceManager指令后,将Slot提供给JobMaster调用
  • JobMaster分配任务后,TaskManager执行具体计算
  • 执行过程中可以缓冲数据,也可以与其他TaskManager交换数据

  • 二、核心概念

    2.1 并行度(Parallelism)

    什么是并行度

    当数据量很大时,可以把一个算子"复制"多份到多个节点上并行处理。一个算子的子任务(subtask)个数就称为该算子的并行度(Parallelism)。

    一个流程序的并行度 = 所有算子中最大的并行度。

    举例:

    source(并行度2)→ map(并行度2)→ window(并行度2)→ sink(并行度1)
    程序并行度 = 2

    在这里插入图片描述

    并行度的设置(优先级从高到低)

    方式一:代码中对单个算子设置(优先级最高)

    stream.map(word -> Tuple2.of(word, 1L)).setParallelism(2);

    方式二:代码中全局设置

    env.setParallelism(2);

    不推荐硬编码全局并行度,会导致无法动态扩容。

    方式三:提交作业时通过-p参数设置

    bin/flink run -p 2 -c com.atguigu.wc.SocketStreamWordCount ./FlinkTutorial-1.0-SNAPSHOT.jar

    方式四:配置文件中设置(优先级最低)

    # flink-conf.yaml
    parallelism.default: 2

    对整个集群所有作业有效,初始值为1。开发环境中没有配置文件,默认并行度 = 当前机器CPU核心数。

    注意:keyBy不是算子,无法设置并行度。


    2.2 算子链(Operator Chain)

    算子间的数据传输方式

    一对一(One-to-one / Forwarding):数据流维护分区和元素顺序,不需要重新分区。map、filter、flatMap等算子都是一对一关系,类似Spark的窄依赖。

    重分区(Redistributing):数据流分区发生改变,比如keyBy、window之后的操作,类似Spark的Shuffle。

    什么是算子链

    满足条件:并行度相同 + 一对一的数据传输关系

    满足条件的相邻算子可以合并成一个Task,由一个线程执行,这就是算子链(Operator Chain)。

    优势:

    • 减少线程间的切换开销
    • 减少基于缓冲区的数据交换
    • 降低延迟,提升吞吐量

    Flink默认会按算子链原则自动合并,也可以手动控制:

    // 禁用当前算子的链接
    .map(word -> Tuple2.of(word, 1L)).disableChaining();

    // 从当前算子开始新的算子链
    .map(word -> Tuple2.of(word, 1L)).startNewChain();

    算子间数据传输
    在这里插入图片描述
    合并算子链
    在这里插入图片描述


    2.3 任务槽(Task Slots)

    什么是任务槽

    每个TaskManager是一个JVM进程,可以启动多个线程并行执行子任务。为了控制并发量,对每个任务运行所占用的资源进行明确划分,这就是任务槽(Task Slot)。

    每个Slot表示TaskManager计算资源的一个固定大小的子集,用于独立执行一个子任务。

    注意:Slot目前只隔离内存,不隔离CPU。

    任务槽数量设置

    # flink-conf.yaml
    taskmanager.numberOfTaskSlots: 8

    默认为1。建议配置为机器的CPU核心数,避免任务间CPU竞争。

    Slot共享

    默认情况下,同一个作业的不同算子的子任务可以共享同一个Slot。

    好处:

  • 资源密集型和非密集型任务可以在同一个Slot中自行分配资源占用比例,使负载平均分配
  • 即使某个TaskManager宕机,其他节点不受影响,作业可以继续执行(保存完整的作业管道)
  • 手动指定Slot共享组(只有同组的子任务才共享Slot):

    .map(word -> Tuple2.of(word, 1L)).slotSharingGroup("group1");

    使用Slot共享组后,所需总Slot数 = 各共享组最大并行度之和。
    在这里插入图片描述


    2.4 任务槽与并行度的关系

    这是容易混淆的概念,一定要区分清楚:

    概念类型含义配置参数
    任务槽(Slot) 静态概念 TaskManager具有的并发执行能力(上限) taskmanager.numberOfTaskSlots
    并行度(Parallelism) 动态概念 程序运行时实际使用的并发能力 parallelism.default

    举例:3个TaskManager,每个TM有3个Slot,共9个Slot。

    • 程序为:source → flatmap → reduce → sink
    • source和flatmap满足算子链条件,合并为一个Task
    • 最终3个Task节点,若并行度设为9,则需要9个Slot

    关系总结:Slot是资源的上限,并行度是实际使用量。并行度不能超过可用Slot总数。


    三、作业提交流程与图的转换

    3.1 四层图的转换

    一个Flink作业从代码到实际执行,要经历四次图的转换:

    用户代码

    逻辑流图(StreamGraph)
    ↓ 算子链优化
    作业图(JobGraph)
    ↓ 按并行度拆分
    执行图(ExecutionGraph)
    ↓ 分发到TaskManager
    物理图(Physical Graph)

    在这里插入图片描述
    在这里插入图片描述

    ①逻辑流图(StreamGraph)

    根据用户DataStream API代码生成的初始DAG图,表示程序的拓扑结构。一般在客户端生成。

    ②作业图(JobGraph)

    StreamGraph经优化后的结果,是提交给JobManager的数据结构。主要优化:将符合条件的算子合并成算子链,减少数据交换消耗。一般也在客户端生成,作业提交时传给JobMaster。

    在Flink Web UI中点击作业可以看到对应的作业图。

    ③执行图(ExecutionGraph)

    JobMaster收到JobGraph后生成,是调度层最核心的数据结构。与JobGraph的最大区别:按照并行度拆分了并行子任务,并明确了任务间数据传输方式。

    ④物理图(Physical Graph)

    JobMaster将执行图分发给TaskManager后,TaskManager部署任务形成的实际执行"图"。这不是一个具体的数据结构,而是实际执行层面的概念。主要在执行图基础上进一步确定数据存放位置和收发的具体方式。


    四、总结

    用一张思维导图来回顾本文的核心内容:

    Flink运行时架构
    ├── 系统架构
    │ ├── JobManager
    │ │ ├── JobMaster(与Job一一对应,生成ExecutionGraph,协调Checkpoint)
    │ │ ├── ResourceManager(管理Slot资源)
    │ │ └── Dispatcher(REST接口,启动JobMaster,Web UI)
    │ └── TaskManager(执行具体计算,提供Slot)

    ├── 核心概念
    │ ├── 并行度(算子子任务数量,动态概念)
    │ ├── 算子链(并行度相同+一对一 → 合并为一个Task,减少开销)
    │ ├── 任务槽(TaskManager资源划分单元,静态概念,只隔离内存)
    │ └── Slot共享(同作业不同算子可共享Slot)

    └── 图的转换
    └── StreamGraph → JobGraph → ExecutionGraph → Physical Graph

    几个容易混淆的点再强调一下:

  • JobManager ≠ JobMaster:JobManager是进程,JobMaster是其中负责单个Job的组件
  • Slot隔离内存,不隔离CPU:设置Slot数量时建议对应CPU核心数
  • 并行度不能超过Slot总数:Slot是上限,并行度是实际使用
  • 算子链的条件:并行度相同 + 一对一传输,缺一不可
  • 赞(0)
    未经允许不得转载:171主机测评 » Flink运行时架构详解
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址