Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 11

Resource scheduling

Learn how ResourceManager, slots, resource declarations, and scheduling map plans to cluster capacity.

Reading
60 min
Track
Flink internals
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 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) 内完成。


下一步