欢迎光临
我们一直在努力

Flink运行时架构:JobManager与TaskManager

一、前言

在掌握了Flink的部署方式之后,我们需要进一步了解Flink作业提交后内部是如何运转的。Flink作为分布式流处理引擎,其运行时架构设计直接决定了作业的执行效率、容错能力和资源利用率。

理解运行时架构,对于以下场景至关重要:

  • 性能调优:知道Slot如何分配,才能合理设置并行度
  • 故障排查:明白Checkpoint机制,才能快速定位恢复问题
  • 面试准备:JobManager/TaskManager职责、提交流程是高频考点

本文将围绕三大核心主题展开:

  • 系统架构:JobManager与TaskManager的职责分工
  • 核心概念:并行度、算子链、任务槽的深层原理
  • 提交流程:从代码到物理执行的完整转换链路

  • 二、系统架构总览

    Flink运行时架构采用经典的Master-Worker模式,主要由两类进程构成:

    在这里插入图片描述

    2.1 JobManager(作业管理器)

    JobManager是整个Flink集群的控制中枢,每个作业对应一个JobManager实例(严格来说是JobMaster)。它包含三个核心组件:

    2.1.1 JobMaster

    JobMaster是JobManager中最核心的组件,负责处理单个作业的全部生命周期:

    职责说明
    作业调度 将JobGraph转换为ExecutionGraph,调度任务执行
    资源申请 向ResourceManager申请Task Slot资源
    Checkpoint协调 触发、协调分布式快照的保存与恢复
    故障恢复 作业失败时,从最近的Checkpoint重新启动

    历史演变:早期Flink版本中,JobManager的概念范围较小,实际指的就是现在的JobMaster。从1.x版本开始,为了支持多作业同时运行,才明确区分了JobManager(进程)和JobMaster(作业实例)的概念。

    2.1.2 ResourceManager

    ResourceManager负责集群资源的统一管理和分配:

    • Slot资源管理:维护集群中所有可用Slot的注册信息
    • 资源分配:根据JobMaster的请求,分配空闲Slot给作业使用
    • 资源回收:作业完成后,回收Slot供其他作业使用
    • 外部对接:在YARN/K8S模式下,与外部资源管理器交互

    ⚠️ 注意区分:Flink内置的ResourceManager与YARN的ResourceManager是不同的组件,前者管理Flink的Slot资源,后者管理整个集群的容器资源。

    2.1.3 Dispatcher

    Dispatcher提供了Flink对外的REST API接口和Web UI:

    • 接收作业提交请求,为每个新作业启动一个JobMaster
    • 提供Web UI展示作业执行状态、Metrics指标
    • 在部分部署模式下(如YARN应用模式)可能被省略

    2.2 TaskManager(任务管理器)

    TaskManager是Flink的工作进程,负责实际的数据计算:

    核心职责:

  • 执行计算任务:接收JobMaster分配的任务,执行具体的算子逻辑
  • 数据缓冲与交换:在任务之间缓存数据,与其他TaskManager交换数据
  • Slot资源注册:启动后向ResourceManager注册自己的Slot
  • 状态维护:在本地维护算子的状态数据
  • 关键特性:

    • 每个TaskManager是一个独立的JVM进程
    • 包含多个Task Slot,Slot数量由配置决定
    • 可以缓冲数据,支持反压(Backpressure)机制

    三、核心概念详解

    3.1 并行度(Parallelism)

    3.1.1 什么是并行度?

    当数据量巨大时,Flink可以将一个算子"复制"多份到不同节点并行执行。一个算子被拆分成的并行子任务数量,就是该算子的并行度。

    // 设置全局并行度
    env.setParallelism(2);

    // 为单个算子设置并行度
    stream.map(new MyMapper()).setParallelism(4);

    3.1.2 并行度设置优先级

    Flink提供了四级并行度配置,优先级从高到低:

    在这里插入图片描述

    优先级设置方式影响范围推荐用法
    1 算子.setParallelism() 单个算子 针对热点算子单独调优
    2 env.setParallelism() 全局默认 开发测试环境
    3 flink run -p 2 提交时全局 生产环境推荐
    4 parallelism.default: 2 集群默认 集群基础配置

    💡 最佳实践:生产环境不要在代码中硬编码并行度,应通过提交参数动态指定,便于根据数据量灵活调整。

    3.1.3 并行度与Slot的关系

    在这里插入图片描述

    核心公式:

    所需Slot数量 = 所有算子中最大的并行度(考虑Slot共享后)

    以WordCount为例,假设:

    • Source并行度 = 2
    • FlatMap并行度 = 2(与Source合并算子链)
    • KeyBy/Sum并行度 = 2
    • Sink并行度 = 1

    则整个作业需要 2个Slot(取最大并行度2),而非5个。

    3.2 算子链(Operator Chain)

    3.2.1 数据传输模式

    算子之间的数据传输有两种模式:

    模式特点对应算子Spark类比
    One-to-One 分区不变,顺序保持 map/filter/flatMap 窄依赖
    Redistributing 重新分区,可能乱序 keyBy/rebalance Shuffle
    3.2.2 算子链合并条件

    Flink会自动将满足以下条件的算子合并为一个Task:

  • 数据传输模式为One-to-One
  • 并行度相同
  • 合并的好处:

    • 减少线程切换开销
    • 减少序列化/反序列化开销
    • 基于缓冲区批量传输数据,提升吞吐量

    // 禁用算子链(调试用)
    stream.map(...).disableChaining();

    // 从当前算子开始新链
    stream.map(...).startNewChain();

    3.3 任务槽(Task Slot)

    3.3.1 Slot的本质

    Slot是TaskManager上资源的最小分配单元,表示TaskManager拥有的计算资源子集(主要是内存)。

    # flink-conf.yaml
    # 每个TaskManager的Slot数量,默认1
    taskmanager.numberOfTaskSlots: 4

    ⚠️ 注意:Slot目前只隔离内存,不涉及CPU隔离。建议将Slot数配置为CPU核心数,避免CPU竞争。

    3.3.2 Slot共享机制

    在这里插入图片描述

    Flink默认允许不同算子的子任务共享Slot,只要它们属于同一作业。这种设计的优势:

  • 资源均衡:资源密集型和非密集型任务共存,避免"忙闲不均"
  • 管道完整性:即使某个TaskManager宕机,其他节点仍能维持完整处理管道
  • 提高利用率:减少Slot空闲等待时间
  • Slot共享组:

    // 手动指定Slot共享组,实现资源隔离
    stream.map(...).slotSharingGroup("group1");
    stream.filter(...).slotSharingGroup("group2");

    不同共享组的任务必须分配到不同Slot,总Slot需求 = 各组最大并行度之和。


    四、作业提交流程

    4.1 四张图的转换

    Flink作业从代码到物理执行,会经历四次图转换:

    在这里插入图片描述

    4.1.1 StreamGraph(逻辑流图)
    • 生成位置:客户端
    • 特点:根据DataStream API代码直接生成,节点=算子
    • 作用:表示程序的拓扑结构
    4.1.2 JobGraph(作业图)
    • 生成位置:客户端
    • 优化:合并符合条件的算子链(Operator Chain)
    • 特点:节点=Task,边=数据交换方式
    • 提交方式:通过REST API或命令行提交给JobManager
    4.1.3 ExecutionGraph(执行图)
    • 生成位置:JobMaster
    • 转换:按并行度拆分为并行子任务
    • 特点:节点=ExecutionVertex(子任务),是调度核心数据结构
    • 作用:明确任务间数据传输方式,指导任务调度
    4.1.4 Physical Graph(物理流图)
    • 生成位置:TaskManager
    • 特点:具体执行层面的图,非正式数据结构
    • 作用:确定数据存放位置和收发方式

    4.2 Standalone会话模式提交流程

    1. 用户提交作业 → Client生成StreamGraph → 优化为JobGraph
    2. Client将JobGraph提交给Dispatcher
    3. Dispatcher启动JobMaster,将JobGraph传递给JobMaster
    4. JobMaster生成ExecutionGraph,向ResourceManager申请Slot
    5. ResourceManager分配Slot,JobMaster将任务分发到TaskManager
    6. TaskManager执行任务,定期汇报心跳和状态

    4.3 YARN应用模式提交流程

    YARN应用模式与Standalone的主要区别在于:

    • main()方法在JobManager执行:客户端只需发送作业jar包
    • 动态资源申请:TaskManager按需启动,非预先分配
    • YARN负责容器管理:Flink的ResourceManager与YARN RM交互

    五、代码实战:观察并行度与Slot

    5.1 实验代码

    以下代码演示如何观察算子链和并行度的实际效果:

    import org.apache.flink.api.common.eventtime.WatermarkStrategy;
    import org.apache.flink.api.common.functions.MapFunction;
    import org.apache.flink.api.java.tuple.Tuple2;
    import org.apache.flink.streaming.api.datastream.DataStreamSource;
    import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

    public class ParallelismDemo {
    public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env =
    StreamExecutionEnvironment.getExecutionEnvironment();

    // 设置全局并行度为2
    env.setParallelism(2);

    // 读取数据源
    DataStreamSource<String> source = env.socketTextStream("hadoop102", 7777);

    // Map算子:与Source形成算子链(One-to-One + 相同并行度)
    SingleOutputStreamOperator<Tuple2<String, Integer>> mapped = source
    .map(new MapFunction<String, Tuple2<String, Integer>>() {
    @Override
    public Tuple2<String, Integer> map(String value) {
    return Tuple2.of(value, 1);
    }
    });

    // keyBy会改变分区,触发重分区,算子链在此处断开
    // sum聚合算子
    SingleOutputStreamOperator<Tuple2<String, Integer>> result = mapped
    .keyBy(value -> value.f0)
    .sum(1);

    // Sink单独设置并行度为1
    result.print().setParallelism(1);

    env.execute("Parallelism Demo");
    }
    }

    5.2 Web UI观察要点

    提交作业后,打开Flink Web UI(http://jobmanager:8081),重点关注:

  • JobGraph视图:观察算子链合并情况
  • Task Managers:查看Slot分配和使用情况
  • Checkpoints:观察Barrier对齐过程

  • 六、常见问题与注意事项

    6.1 Slot不足导致作业无法启动

    现象:作业提交后一直处于CREATED状态,无法进入RUNNING。

    原因:申请的Slot总数超过了集群可用Slot数。

    解决:

    # 查看可用Slot数量
    # 在Web UI Overview页面查看"Task Managers"和"Available Slots"

    # 方案1:增加TaskManager数量或每个TM的Slot数
    # 方案2:降低作业并行度
    flink run -p 2 your-job.jar # 将全局并行度设为2

    6.2 算子链断开导致性能下降

    现象:作业吞吐量低,CPU利用率不高。

    原因:算子之间频繁序列化/反序列化,网络传输开销大。

    排查:在Web UI的JobGraph中查看算子链是否被意外断开。

    常见断开原因:

    • 并行度不一致
    • 显式调用了startNewChain()或disableChaining()
    • 算子之间是Redistributing模式(如keyBy)

    6.3 状态过大导致OOM

    现象:TaskManager频繁崩溃,日志显示内存溢出。

    原因:状态数据持续增长,超过TaskManager可用内存。

    解决:

    // 1. 使用RocksDB状态后端,将状态存储到磁盘
    env.setStateBackend(new EmbeddedRocksDBStateBackend());

    // 2. 配置状态TTL,自动清理过期状态
    StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(24))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .build();


    总结

    本文深入解析了Flink运行时架构的核心内容:

  • JobManager作为控制中枢,通过JobMaster、ResourceManager、Dispatcher三个组件协同管理作业生命周期
  • TaskManager作为工作进程,在Slot中执行具体的计算任务
  • 并行度、算子链、Slot三者紧密配合,决定了作业的执行效率和资源利用率
  • 四张图的转换体现了Flink从逻辑定义到物理执行的完整链路
  • 理解这些底层原理,将帮助你在实际工作中更好地进行性能调优和故障排查。下一篇我们将深入Flink的DataStream API,从执行环境到Source算子,开启Flink编程的核心篇章。

    赞(0)
    未经允许不得转载:171主机测评 » Flink运行时架构:JobManager与TaskManager
    分享到: 更多 (0)

    评论 抢沙发

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