
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、节省资源 |
| Adaptive | Source 决定下游并行度(1.15+ 默认) | 流批统一、自动优化平行度 |
AdaptiveScheduler 的核心逻辑
// AdaptiveScheduler.java (不在默认代码路径,1.15+ 可选)
// 核心思想:
// 1. 从 Source 开始,运行 Source 的 Subtask
// 2. 检测 Source 的数据分区数 → 推断最优的并行度
// 3. 根据 #partitions 确定下游并行度 (如果上游并行度=N,下游可以是 N×factor)
// 4. 不需要用户手动指定所有算子的并行度
6. 源码导航(完整版)
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-runtime/.../resourcemanager/ResourceManager.java | registerTaskExecutor(), declareRequiredResources() | RM 基类 |
flink-runtime/.../resourcemanager/ActiveResourceManager.java | startNewWorker(), onWorkerRegistered() | 主动 RM 基类 |
flink-runtime/.../resourcemanager/slotmanager/SlotManager.java | processResourceRequirements(), freeSlot() | Slot 管理 |
flink-runtime/.../jobmaster/slotpool/DeclarativeSlotPoolBridge.java | requestNewAllocatedSlot(), newSlotsAreAvailable() | JM 端 SlotPool |
flink-runtime/.../taskexecutor/slot/TaskSlot.java | — | Slot 实现 |
flink-yarn/.../YarnResourceManager.java | startNewWorker() | YARN 模式 |
flink-kubernetes/.../KubernetesResourceManager.java | startNewWorker() | 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) 内完成。
下一步
- TM 内存: doc-10-taskmanager-memory.md
- Connector 框架: doc-12-connector-framework.md