Flink 2.3.0 从理论到实践 —— 第 2 章 Flink 运行时架构

发布时间:2026/10/10 6:17:35
Flink 2.3.0 从理论到实践 —— 第 2 章 Flink 运行时架构 Flink 2.3.0 从理论到实践 —— 第 2 章 Flink 运行时架构课程定位本章深入 Flink 2.3.0 的内部运作机制理解 JobManager / TaskManager 的协作、作业从提交到执行的完整链路、Task / SubTask / Slot 的资源模型以及 Flink 2.x 引入的自适应调度与可插拔运行时。掌握架构是后续调优与排错的基础。版本基线Flink 2.3.0 JDK 17章节导读2.1 架构总览两层进程模型2.2 核心组件详解2.3 作业执行链路从提交到物理执行2.4 任务与资源模型2.5 执行模式流 / 批 / 自动2.6 高可用HA2.7 Flink 2.x 架构演进2.8 本章小结与下章预告2.1 架构总览两层进程模型Flink 采用两层进程模型一个JobManager集群控制面 多个TaskManager数据面。┌──────────────────────────────────────────────────────────────────────┐ │ Flink 集群架构 │ └──────────────────────────────────────────────────────────────────────┘ ┌──────────────────────────────────────────────────────────────────┐ │ JobManager控制面 │ │ ┌───────────┐ ┌────────────┐ ┌───────────────┐ ┌───────────┐ │ │ │Dispatcher│ │JobMaster │ │ResourceManager│ │Web UI/REST│ │ │ │接收作业 │ │协调单个作业 │ │管理TM/资源 │ │监控入口 │ │ │ └───────────┘ └────────────┘ └───────────────┘ └───────────┘ │ └───────────────────────────────┬──────────────────────────────────┘ │ 心跳 任务调度 ┌───────────────────────┼───────────────────────┐ ▼ ▼ ▼ ┌───────────────┐ ┌───────────────┐ ┌───────────────┐ │ TaskManager 1 │ │ TaskManager 2 │ │ TaskManager N │ │ ┌───────────┐ │ │ ┌───────────┐ │ │ ┌───────────┐ │ │ │ Slot 1 │ │ │ │ Slot 1 │ │ │ │ Slot 1 │ │ │ │ Slot 2 │ │ │ │ Slot 2 │ │ │ │ Slot 2 │ │ │ │ ... │ │ │ │ ... │ │ │ │ ... │ │ │ └───────────┘ │ │ └───────────┘ │ │ └───────────┘ │ │ 本地状态网络 │ │ 本地状态网络 │ │ 本地状态网络 │ └───────────────┘ └───────────────┘ └───────────────┘两个核心原则原则说明JM 无状态逻辑上TM 有状态物理上JM 负责调度与元数据TM 承载数据计算与本地状态JM 与 TM 解耦JM 故障不影响正在运行的 TM 计算可从 Checkpoint 恢复2.2 核心组件详解2.2.1 JobManagerJMJobManager 是 Flink 集群的控制面负责作业的接收、调度、协调与监控。它内部包含以下组件组件职责Dispatcher接收客户端提交的作业启动 JobMaster提供 REST API 与 Web UIJobMaster每个作业一个 JobMaster负责该作业的调度、Checkpoint 协调、故障恢复ResourceManager管理 TaskManager 的注册与资源分配处理 Slot 请求与回收Web UI / REST提供作业监控、集群状态、历史查询等可视化入口关键点JobMaster 是每作业一个不是全局一个。多个作业共享同一个 JobManager 进程但各自有独立的 JobMaster。2.2.2 TaskManagerTMTaskManager 是 Flink 集群的数据面Worker 进程负责执行具体的计算任务。组件职责Task SlotTM 的资源单位一个 Slot 对应一组计算资源CPU 内存切片Task / SubTask实际运行在 Slot 中的计算线程Network BufferTM 间数据传输的缓冲区基于 NettyState Backend本地状态存储HashMap / ForStFlink 2.x 变化TM 的内存模型在 2.x 中进一步细化引入了更精细的托管内存Managed Memory划分用于状态、网络缓冲、批处理等。2.2.3 客户端Client客户端不是集群的一部分负责将用户程序编译成JobGraph通过 Dispatcher 提交到集群接收作业执行结果可选用户代码 → Client 编译 → JobGraph → Dispatcher → JobMaster → 调度到 TM2.3 作业执行链路从提交到物理执行一个 Flink 作业从用户提交到实际执行要经历四层图结构的转化。2.3.1 四层图结构用户代码(StreamGraph) │ Client 编译 优化 ▼ JobGraph(可提交的作业图) │ JobMaster 生成 ▼ ExecutionGraph(并行执行图) │ 调度器部署 ▼ 物理执行(Task 在 TM 的 Slot 中运行)阶段产物生成者说明1. StreamGraph流图Client用户 API 调用生成的最初 DAG算子间以边连接2. JobGraph作业图Client 优化算子链化Operator Chain 中间结果分区策略可提交3. ExecutionGraph执行图JobMasterJobGraph 的并行化版本每个算子展开为并行实例4. 物理执行Task调度器ExecutionGraph 部署到 TM 的 Slot 中实际运行2.3.2 算子链化Operator Chain为了减少线程切换与网络开销Flink 会把能链在一起的算子合并成一个 Task在同一个线程中执行。未链化(3 个 Task,3 次线程切换): Source → Map → Filter → Sink 链化后(1 个 Task,0 次线程切换): [Source → Map → Filter → Sink] ← 同一个线程算子链化的条件必须同时满足条件说明上下游并行度相同并行度不同无法链化上下游在同一个 Slot Sharing Group默认都在default组数据转发策略为 Forward即forward分区非rebalance/keyBy没有禁用链化未调用disableChaining()上下游算子数量未超限受chain.length限制实践算子链化是 Flink 性能优化的重要手段。但在某些场景如需要独立监控某个算子下可以通过disableChaining()或slotSharingGroup()手动拆分。2.3.3 调度过程JobMaster 中的Scheduler负责将 ExecutionGraph 中的 Task 调度到 TM 的 Slot 中JobMaster(Scheduler) │ 1. 请求 Slot ▼ ResourceManager │ 2. 分配 TM Slot ▼ TaskManager │ 3. 启动 Task 线程 ▼ Task 在 Slot 中运行2.4 任务与资源模型这是 Flink 中最容易混淆的概念。我们用一个具体例子来说明。2.4.1 核心概念定义概念定义Operator算子用户定义的计算逻辑单元如map、keyBy、windowTask算子链化后的物理执行单元一个 Task 包含一个或多个算子SubTaskTask 的并行实例Task 有并行度 N 就有 N 个 SubTaskSlotTaskManager 的资源单位一个 Slot 可以运行多个 SubTask不同 Task 的Parallelism并行度一个算子/Task 的并行实例数量2.4.2 关系图解假设一个作业有 3 个算子Source、Map、Sink并行度为 2且能链化为一个 TaskTaskManager2 个 Slot ┌────────────────────────────────────┐ │ Slot 0 │ Slot 1 │ │ ┌────────────┐ │ ┌────────────┐ │ │ │ SubTask 0 │ │ │ SubTask 1 │ │ │ │ [Source-0] │ │ │ [Source-1] │ │ │ │ [Map-0] │ │ │ [Map-1] │ │ │ │ [Sink-0] │ │ │ [Sink-1] │ │ │ └────────────┘ │ └────────────┘ │ └────────────────────────────────────┘ Task 数 1链化后 SubTask 数 2并行度 Slot 数 2每个 SubTask 占 1 个 Slot2.4.3 Slot Sharing槽共享Slot Sharing是 Flink 的关键设计同一个作业的不同 Task 的 SubTask 可以共享同一个 Slot。假设 3 个 TaskA、B、C并行度 22 个 Slot Slot 0 Slot 1 ┌──────────────┐ ┌──────────────┐ │ A-SubTask-0 │ │ A-SubTask-1 │ │ B-SubTask-0 │ │ B-SubTask-1 │ │ C-SubTask-0 │ │ C-SubTask-1 │ └──────────────┘ └──────────────┘ 3 个 Task × 2 并行度 6 个 SubTask 共享到 2 个 Slot 中每个 Slot 跑 3 个 SubTaskSlot Sharing 的好处好处说明资源利用率高轻算子和重算子共享 Slot避免轻算子占着 Slot 不干活Slot 数 最大并行度只需保证 Slot 数 ≥ 作业的最大并行度即可负载均衡每个 Slot 运行完整的算子管线负载更均匀Slot Sharing Group默认所有算子在default组可共享可通过slotSharingGroup(name)为算子指定组不同组之间不能共享常用于隔离重型算子如大状态的 Join到独立 Slot结论作业所需的最少 Slot 数 所有 Slot Sharing Group 中的最大并行度。2.4.4 并行度的来源与优先级并行度的设置有四个来源按优先级从高到低优先级 1: 算子级 .setParallelism(n) ← 最高 优先级 2: 环境级 env.setParallelism(n) 优先级 3: 客户端命令 -p n / flink run -p n 优先级 4: 集群默认 parallelism.default ← 最低注意来自实测经验Flink SQL 中SET table.exec.default-parallelism n设置的是 Table API 的默认并行度它与 DataStream 的env.setParallelism()是两套机制。Kafka Source 的并行度由 Kafka 分区数决定不受table.exec.default-parallelism控制。2.4.5 资源配置示例假设一个 TM 配置taskmanager.numberOfTaskSlots: 4taskmanager.memory.process.size: 8g配置项值说明taskmanager.numberOfTaskSlots4每个 TM 4 个 Slottaskmanager.memory.process.size8gTM 总内存含 JVM、托管、网络taskmanager.memory.task.heap.size~3gTask 堆内存taskmanager.memory.managed.size~4g托管内存状态、批处理Slot 与内存的关系Slot 不独占内存内存是 TM 级别共享的。Slot 主要是线程调度与资源隔离的逻辑单位不是严格的内存隔离。2.5 执行模式流 / 批 / 自动Flink 2.x 支持三种执行模式通过execution.runtime-mode配置。2.5.1 三种模式对比模式execution.runtime-mode特点适用场景流模式streaming持续运行处理无界流实时计算批模式batch处理有界数据跑完即退出离线 ETL自动模式automatic根据 Source 是否有界自动选择流批一体2.5.2 流模式 vs 批模式的执行差异流模式(streaming): Source ──► Map ──► KeyBy ──► Window ──► Sink (持续运行, 数据逐条流动, 增量计算) 批模式(batch): Source ──► Map ──► Shuffle ──► Reduce ──► Sink (阶段执行, 上一阶段完成后才开始下一阶段, 全量计算)维度流模式批模式数据流动逐条流式按阶段Stage批量Shuffle 策略Pipeline流水线Blocking阻塞状态持续维护Checkpoint 快照阶段结束即释放容错Checkpoint 恢复重算整个阶段调度所有 Task 同时调度按阶段依次调度内存状态 网络缓冲排序、哈希表2.5.3 自动模式Source 有界? ──是──► batch 模式 │ 否 ▼ streaming 模式适用同一套 SQL 既能跑实时又能跑离线时使用。例如SELECT * FROM kafka_table流vsSELECT * FROM hive_table批设为automatic即可自动切换。2.6 高可用HA生产环境中JobManager 的单点故障会导致整个集群不可用。Flink 通过Standby JobManager实现高可用。2.6.1 HA 架构┌────────────────────────────────────────────────────────────┐ │ ZooKeeper / Kubernetes │ │ (用于 Leader 选举 元数据持久化) │ └──────────────────┬─────────────────────┬───────────────────┘ │ │ ┌────────▼────────┐ ┌────────▼────────┐ │ JobManager │ │ JobManager │ │ (Leader) │ │ (Standby) │ └─────────────────┘ └─────────────────┘ │ ┌────────▼───────────────┐ │ TaskManager 集群 │ └─────────────────────────┘2.6.2 HA 机制机制说明Leader 选举多个 JM 竞争选出一个 Leader其余 Standby元数据持久化JobGraph、Checkpoint 元数据存到 ZooKeeper/K8s 持久化存储故障切换Leader 故障时Standby 接管从持久化存储恢复作业2.6.3 两种 HA 方案方案依赖适用场景ZooKeeper HAZooKeeper 集群YARN / Standalone 部署Kubernetes HAK8s API ServerKubernetes 部署Flink 2.x 趋势Kubernetes HA 方案逐渐成为主流不依赖外部 ZooKeeper运维更简单。2.7 Flink 2.x 架构演进Flink 2.0 对运行时架构做了重大重构引入了两个关键能力。2.7.1 Adaptive Scheduler自适应调度器传统调度器在作业启动时静态分配并行度与 Slot运行中无法调整。Adaptive Scheduler 支持运行时动态调整维度传统调度器Adaptive Scheduler并行度启动时固定可运行时增减Slot 分配静态动态申请/释放扩缩容需重启在线调整资源利用按峰值规划按需弹性传统: 3 TM × 4 Slot 12 Slot 全部预占 自适应: 闲时 2 TM × 2 Slot 4 Slot,忙时自动扩到 12 Slot2.7.2 Pluggable Runtime可插拔运行时Flink 2.x 将运行时抽象为可插拔接口允许接入不同的执行引擎Flink API / SQL 层 │ ▼ ┌─────────────────────────────────┐ │ Pluggable Runtime Interface │ ├──────────────┬──────────────────┤ │ 传统执行引擎 │ 其他执行引擎 │ │ (Stream/Batch)│ (如 Flink-1.x │ │ │ 兼容模式等) │ └──────────────┴──────────────────┘意义为 Flink 未来支持多种执行后端如自研引擎、异构计算打下基础是架构层的长期投资。2.7.3 其他 2.x 架构变化变化说明ForSt 状态后端RocksDB 的替代实现性能更优2.x 默认推荐内存模型细化托管内存划分更精细支持状态、网络、批处理独立配置网络栈优化基于 Netty 的网络栈升级提升 TM 间吞吐2.8 本章小结与下章预告本章小结┌────────────────────────────────────────────────────────────────┐ │ 第 2 章 要点回顾 │ └────────────────────────────────────────────────────────────────┘ ✓ 两层进程模型: JobManager(控制面) TaskManager(数据面) JM 内部: Dispatcher / JobMaster / ResourceManager / Web UI ✓ 作业四层图: StreamGraph → JobGraph(算子链化) → ExecutionGraph(并行化) → 物理执行 算子链化条件: 同并行度 同 Slot Group forward 未禁用 ✓ 资源模型: Operator → Task(链化后) → SubTask(并行实例) → Slot(资源单位) Slot Sharing: 同作业不同 Task 的 SubTask 可共享 Slot 最少 Slot 数 最大 Slot Sharing Group 的并行度 ✓ 并行度优先级: 算子级 环境级 客户端命令 集群默认 Kafka Source 并行度 Kafka 分区数 ✓ 执行模式: streaming(持续) / batch(阶段) / automatic(自动判断) ✓ 高可用: ZooKeeper HA / Kubernetes HA Leader 选举 元数据持久化 Standby 接管 ✓ Flink 2.x: Adaptive Scheduler(在线扩缩容) Pluggable Runtime(可插拔执行引擎) ForSt 状态后端(替代 RocksDB)下章预告第 3 章 环境准备与集群部署动手搭建 Flink 2.3.0 运行环境。涵盖 JDK 17 与依赖准备、本地模式安装、Standalone 集群、YARN 三种模式Session / Per-Job / Application、Kubernetes Operator 部署、flink-conf.yaml核心配置参数、Web UI 与 REST API 基础使用。官方参考资料Flink 架构官方文档https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/flink-architecture/Task Slot 与资源https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/flink-architecture/#task-slots-and-resources执行模式https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/execution_mode/高可用配置https://nightlies.apache.org/flink/flink-docs-stable/docs/deployment/ha/overview/

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询