Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 11

ResourceManager & 调度深度解析

理解 ResourceManager、Slot、资源声明和调度策略如何把计算计划落到实际集群资源。

阅读时间
60 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

预计阅读时间: 60 分钟 前置阅读: doc-01(Job 调度) 下一次阅读: doc-10(TM 内存)


ResourceManager & 调度深度解析 图 01

1. ResourceManager 架构

1.1 三种部署模式

// ResourceManager.java — 抽象基类
public abstract class ResourceManager<T extends ResourceID>
        extends FencedRpcEndpoint<ResourceManagerId> {

    // 已注册的 TaskManager
    protected final Map<ResourceID, TaskExecutorRegistration> taskExecutors;

    // 已注册的 JobMaster
    protected final Map<JobID, JobMasterRegistration> jobMasterRegistrations;

    // Slot 管理器 (跟踪集群所有 Slot)
    private final SlotManager slotManager;

    // 资源分配器 (向 YARN/K8s 申请新 Container/Pod)
    private final ResourceAllocator resourceAllocator;
}

部署模式对应的实现:

Standalone → StandaloneResourceManager    (静态 TM,不动态分配)
YARN       → YarnResourceManager         (向 YARN 申请 Container)
K8s        → KubernetesResourceManager   (向 K8s 申请 Pod)

1.2 ResourceManager 核心流程

// TM 注册
public void registerTaskExecutor(TaskExecutorRegistration registration) {
    // 1. 验证 TM 身份
    // 2. 注册到 SlotManager (初始化 Slot 列表)
    // 3. 存储到 taskExecutors
    // 4. 启动心跳
}

// JM 声明资源需求
public void declareRequiredResources(JobMasterId jmId, ResourceRequirements requirements) {
    // 1. 验证 JobMasterId
    // 2. slotManager.processResourceRequirements(requirements)
    //    → 匹配可用 Slot (SlotManager 内部)
    //    → 不足 → 触发 ResourceAllocator.allocateResources()
}

// TM 通知 Slot 可用
public void notifySlotAvailable(InstanceID instanceId, SlotID slotId, AllocationID) {
    // → slotManager.freeSlot(slotId, allocationId)
    // → SlotManager 将 Slot 分配给等待中的 Job
}

2. Slot 模型与 SlotSharingGroup

2.1 Slot 的物理含义

TaskManager JVM
│
├── TaskSlot 0 (SlotID: "slot-0", ResourceID: "tm-1")
│   ├── SubTask A-0 (Source, parallelism_idx=0)
│   ├── SubTask B-0 (Map, parallelism_idx=0)
│   └── SubTask C-0 (Sink, parallelism_idx=0)
│
├── TaskSlot 1 (SlotID: "slot-1", ResourceID: "tm-1")
│   ├── SubTask A-1 (Source, parallelism_idx=1)
│   ├── SubTask B-1 (Map, parallelism_idx=1)
│   └── SubTask C-1 (Sink, parallelism_idx=1)
│
└── TaskSlot 2 ...

默认 SlotSharingGroup = "default" 时:

  • 同 Parallelism Index 的不同算子在同一个 Slot 中
  • Slot 数 = max(所有算子的并行度)
  • 最小资源需求 = max(单个算子资源需求),而非 sum

2.2 Slot 分配流程

JobMaster.SlotPool (需求端)
  │
  │ requestNewAllocatedSlot(resourceProfile)
  ▼
ResourceManager.declareRequiredResources(requirements)
  │
  │ SlotManager.processResourceRequirements()
  │   ├── 已有空闲 Slot → 直接分配
  │   └── 不够 → ResourceAllocator.allocateResources()
  │
  ▼
YARN: new Container → launch TaskManager
K8s:  new Pod → launch TaskManager
  │
  │ TM 启动后向 RM 注册 → offerSlots
  ▼
JobMaster.SlotPool.newSlotsAreAvailable()
  │
  │ 匹配 PendingRequests → 满足请求
  ▼
ExecutionDeployer.deploy(Execution)
  │
  ▼
TaskExecutor.submitTask(TaskDeploymentDescriptor)

3. YARN 模式深度分析

3.1 架构

YARN Cluster:
├── ResourceManager (YARN)
│   ├── ApplicationMaster (Flink JobManager)  ← Container 1
│   │   └── YarnResourceManager
│   ├── TaskManager-1 (Flink TM)               ← Container 2
│   ├── TaskManager-2 (Flink TM)               ← Container 3
│   └── ...

3.2 YarnResourceManager

// YarnResourceManager.java — YARN 模式资源管理
public class YarnResourceManager extends ActiveResourceManager<YarnWorkerNode> {

    // 向 YARN 申请 Container
    protected void startNewWorker(ResourceProfile resourceProfile) {
        // 1. 构造 ContainerRequest
        AMRMClient.ContainerRequest request = new ContainerRequest(
            capability,  // 内存 + vCores
            null,        // 不限主机
            null,        // 不限机架
            Priority.newInstance(1)
        );

        // 2. 发送到 YARN RM: allocate()
        amrmClient.addContainerRequest(request);

        // 3. 异步等待 YARN 分配 Container
        // 4. 在分配的 Container 中启动 TaskManager
    }

    // YARN Container → Flink TaskManager
    protected void onContainerAllocated(Container container) {
        // 1. 构造 TaskManager 启动命令
        String tmCommand = TaskManagerRunner.composeStartCommand();

        // 2. 通过 ContainerLauncher 启动 TM
        containerLauncher.launch(
            ContainerLaunchContext.newInstance(tmCommand, env, ...)
        );

        // 3. 等待 TM 启动并向 Flink RM 注册
    }
}

4. K8s 模式深度分析

4.1 Session vs Application Mode

Session Mode:
  - 1 个常驻 JM Pod + 1 个常驻 TM Deployment (cluster)
  - 多个 Job 共享同一个 FLink 集群
  - 适合: 多个短 Job、快速提交

Application Mode:
  - 每个 Job = 1 个完整的 Flink 集群 (JM+TM)
  - Job 完成后集群销毁
  - 适合: 长 Job、资源隔离需求

4.2 KubernetesResourceManager

// KubernetesResourceManager.java — K8s 模式
public class KubernetesResourceManager extends ActiveResourceManager<KubernetesWorkerNode> {

    // 创建 TM Pod
    protected void startNewWorker(ResourceProfile resourceProfile) {
        // 1. 根据 TaskManager 配置生成 Pod Spec
        Pod pod = new PodBuilder()
            .withNewMetadata()
                .withName("flink-tm-" + UUID.randomUUID())
                .withLabels(Map.of("app", "flink", "component", "taskmanager"))
            .endMetadata()
            .withNewSpec()
                .addNewContainer()
                    .withName("flink-taskmanager")
                    .withImage(tmImage)
                    .withResources(new ResourceRequirementsBuilder()
                        .withRequests(Map.of("cpu", cpuQuantity, "memory", memoryQuantity))
                        .build())
                .endContainer()
            .endSpec()
            .build();

        // 2. 创建 Pod
        kubernetesClient.pods().create(pod);

        // 3. 等待 Pod Ready
    }
}

5. 调度策略

策略行为适用场景
Eager所有 Vertex 同时分配 Slot小 Job、资源充足
Lazy-From-Sources从 Source 开始,逐步向下游分配大 Job、节省资源
AdaptiveSource 决定下游并行度(1.15+ 默认)流批统一、自动优化平行度

AdaptiveScheduler 的核心逻辑

// AdaptiveScheduler.java (不在默认代码路径,1.15+ 可选)
// 核心思想:
// 1. 从 Source 开始,运行 Source 的 Subtask
// 2. 检测 Source 的数据分区数 → 推断最优的并行度
// 3. 根据 #partitions 确定下游并行度 (如果上游并行度=N,下游可以是 N×factor)
// 4. 不需要用户手动指定所有算子的并行度

6. 源码导航(完整版)

文件关键类/方法职责
flink-runtime/.../resourcemanager/ResourceManager.javaregisterTaskExecutor(), declareRequiredResources()RM 基类
flink-runtime/.../resourcemanager/ActiveResourceManager.javastartNewWorker(), onWorkerRegistered()主动 RM 基类
flink-runtime/.../resourcemanager/slotmanager/SlotManager.javaprocessResourceRequirements(), freeSlot()Slot 管理
flink-runtime/.../jobmaster/slotpool/DeclarativeSlotPoolBridge.javarequestNewAllocatedSlot(), newSlotsAreAvailable()JM 端 SlotPool
flink-runtime/.../taskexecutor/slot/TaskSlot.javaSlot 实现
flink-yarn/.../YarnResourceManager.javastartNewWorker()YARN 模式
flink-kubernetes/.../KubernetesResourceManager.javastartNewWorker()K8s 模式

7. 常见问题 / 面试题

Q1: SlotSharingGroup 如何影响资源利用率?

A: 默认所有算子在一个 SlotSharingGroup → Slot 数 = max(并行度) → 资源利用率高。但可以通过显式设置不同的 SlotSharingGroup 来隔离资源。例如,给 Source 和 Window 分配不同的 SlotSharingGroup → Source 和 Window 不会在同一 Slot 中 → 保证各自有独立的资源。

Q2: YARN vs K8s 部署选哪个?

A: YARN: Hadoop 生态成熟,和 HDFS/Hive/Spark 天然集成,多租户隔离靠 YARN Queue。K8s: 容器原生,弹性扩缩容方便,适合云原生的环境。2024+ 新项目普遍推荐 K8s Application 模式。

Q3: TM 向 RM 注册时带了什么信息?

A: TaskExecutorRegistration 包含:(1) ResourceID (唯一标识 TM);(2) Slot 列表 (每个 Slot 的 SlotID + ResourceProfile);(3) TM 的总硬件规格;(4) DataPort (Shuffle 服务端口)。RM 用这些信息初始化 SlotManager 的 Slot 记录。

Q4: Slot 被释放后如何重新分配?

A: TM.notifySlotAvailable()RM.slotManager.freeSlot()SlotManager 将该 Slot 入队 (free slots queue) → 检查是否有等待中的 ResourceRequirement → 匹配后分配给对应 Job 的 SlotPool。整个过程在 TM 心跳间隔 (默认 10s) 内完成。


下一步