news 2026/10/12 2:15:21

Flink 资源申请流程源码级别跟踪:从SlotPool到TaskExecutor的完整调用链

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink 资源申请流程源码级别跟踪:从SlotPool到TaskExecutor的完整调用链

前面讲了 SlotManager 的启动流程,这篇聚焦资源申请的完整链路。你有没有想过,JobMaster 需要 Slot 时,一次 requestSlot 调用经过了哪些组件?SlotPool 到 ResourceManager 到 SlotManager 再到 TaskExecutor,每个环节做了什么?没有空闲 Slot 时请求是如何等待和超时的?这篇从源码级别跟踪资源申请的完整调用链。


一、资源申请整体架构

下面这张图是 Flink 资源申请整体架构,包括四个核心组件和资源申请流程。

资源申请涉及四个核心组件,各司其职:

  • SlotPool:运行在 JobMaster 内部,是作业级 Slot 池。负责向 ResourceManager 申请 Slot,管理已分配的 Slot,提供 Slot 给调度器使用。
  • ResourceManager:集群级资源管理器。协调 Slot 分配,管理 TaskManager 注册和心跳,将 Slot 请求转发给 SlotManager 处理。
  • SlotManager:运行在 ResourceManager 内部,是 Slot 分配的实际执行者。管理所有 TaskManagerSlot 的状态,查找空闲 Slot,处理 Slot 分配和回收。
  • TaskExecutor:运行在 TaskManager 内部。接收 Slot 分配请求,管理本地 SlotTable,部署和执行 Task。

四个组件通过 RPC 调用协作,构成完整的资源申请链路。


二、资源申请完整流程时序

下面这张图是资源申请完整流程时序图,包括6步消息流和源码级别调用链。

一次完整的资源申请分为6步:

第1步:SlotPool 发起请求。SlotPool 调用requestSlot(jobId, allocationId, resourceProfile, targetAddress),通过 ResourceManagerGateway 向 RM 发送 RPC 请求。传入作业ID、分配ID、资源需求等参数。

第2步:RM 转发给 SlotManager。ResourceManager 收到请求后,调用slotManager.requestSlot(jobId, allocationId, resourceProfile),将请求转发给 SlotManager 处理。

第3步:SlotManager 查找空闲 Slot。SlotManager 遍历所有 TaskManagerSlot,调用findFreeSlot(resourceProfile)查找 FREE 状态且满足资源需求的 Slot。

第4步:通知 TaskExecutor 分配。找到 Slot 后标记为 PENDING,通过 RM 向 TaskExecutor 发送requestSlot(slotId, allocationId, jobId)请求,TaskExecutor 在本地 SlotTable 中分配 Slot。

第5步:TaskExecutor 确认分配。TaskExecutor 分配成功后返回 Acknowledge,SlotManager 收到确认后调用slot.markAllocated(),将 Slot 状态从 PENDING 改为 ALLOCATED。

第6步:RM 返回 SlotPool 成功。ResourceManager 向 SlotPool 返回SlotRequestSuccess(allocationId, slotId),SlotPool 接收 Slot,加入已分配集合,通知等待的 Execution 部署 Task。


三、源码级别调用链跟踪

下面这张图是源码级别跟踪,包括四个核心类的关键方法、requestSlot 参数详解、最佳实践和总结。

3.1 SlotPoolImpl.requestSlot

publicclassSlotPoolImplimplementsSlotPool{privatefinalMap<AllocationID,PendingRequest>pendingRequests;privatefinalMap<AllocationID,AllocatedSlot>allocatedSlots;@OverridepublicCompletableFuture<SlotRequestSuccess>requestSlot(JobIDjobId,AllocationIDallocationId,ResourceProfileresourceProfile,StringtargetAddress){// 1. 创建 PendingRequest,记录请求信息PendingRequestpendingRequest=newPendingRequest(allocationId,jobId,resourceProfile,targetAddress);pendingRequests.put(allocationId,pendingRequest);// 2. 通过 RM Gateway 发送 RPC 请求returnresourceManagerGateway.requestSlot(jobId,allocationId,resourceProfile,targetAddress,resourceManagerId,timeout);}}

SlotPool 是请求的发起方,创建 PendingRequest 记录请求状态,然后通过 RPC 调用 RM。

3.2 ResourceManagerImpl.requestSlot

publicclassResourceManagerImpl<WorkerTypeextendsResourceIDRetrievable>extendsFencedRpcEndpoint<ResourceManagerId>implementsResourceManager<WorkerType>{privatefinalSlotManagerslotManager;@OverridepublicCompletableFuture<Acknowledge>requestSlot(JobIDjobId,AllocationIDallocationId,ResourceProfileresourceProfile,StringtargetAddress,ResourceManagerIdresourceManagerId){// 验证 fencing tokenvalidateRunsInMainThread();// 转发给 SlotManager 处理slotManager.requestSlot(jobId,allocationId,resourceProfile);returnCompletableFuture.completedFuture(Acknowledge.get());}}

RM 是协调者,验证请求合法性后直接转发给 SlotManager,不做实际分配逻辑。

3.3 SlotManagerImpl.requestSlot

publicclassSlotManagerImplimplementsSlotManager{privatefinalMap<ResourceID,Map<SlotID,TaskManagerSlot>>taskManagerSlots;privatefinalMap<AllocationID,PendingSlotRequest>pendingSlotRequests;@OverridepublicCompletableFuture<Acknowledge>requestSlot(JobIDjobId,AllocationIDallocationId,ResourceProfileresourceProfile){// 1. 查找 FREE SlotTaskManagerSlotfreeSlot=findFreeSlot(resourceProfile);if(freeSlot!=null){// 2. 标记 PENDINGfreeSlot.assignPendingAllocation(allocationId,jobId);// 3. 创建 PendingSlotRequestPendingSlotRequestpendingRequest=newPendingSlotRequest(allocationId,jobId,resourceProfile,freeSlot.getSlotId());pendingSlotRequests.put(allocationId,pendingRequest);// 4. 通过 RM 通知 TaskExecutor 分配TaskExecutorGatewaytaskExecutorGateway=taskManagerConnections.get(freeSlot.getTaskManagerId());returntaskExecutorGateway.requestSlot(freeSlot.getSlotId(),jobId,allocationId,targetAddress,timeout);}else{// 无 FREE Slot,创建 PendingSlotRequest 等待PendingSlotRequestpendingRequest=newPendingSlotRequest(allocationId,jobId,resourceProfile,null);pendingSlotRequests.put(allocationId,pendingRequest);// 通知 RM 申请新资源notifyResourceShortage(resourceProfile);returnCompletableFuture.completedFuture(Acknowledge.get());}}}

SlotManager 是分配的核心执行者。找到空闲 Slot 时标记 PENDING 并通知 TM;无空闲时创建等待请求并通知 RM 申请新容器。

3.4 TaskExecutor.requestSlot

publicclassTaskExecutorextendsTaskExecutorFencedRpcEndpoint{privatefinalSlotTableslotTable;@OverridepublicCompletableFuture<Acknowledge>requestSlot(SlotIDslotId,JobIDjobId,AllocationIDallocationId,StringtargetAddress){// 1. 在本地 SlotTable 中分配 SlotslotTable.addSlot(slotId,jobId,allocationId,resourceManagerId);// 2. 返回确认returnCompletableFuture.completedFuture(Acknowledge.get());}}

TaskExecutor 是 Slot 的实际提供者,在本地 SlotTable 中分配 Slot 资源,返回确认。

3.5 TaskManagerSlot 状态转换

publicclassTaskManagerSlot{privateSlotStatestate;// FREE / PENDING / ALLOCATED / RELEASED// FREE → PENDINGpublicvoidassignPendingAllocation(AllocationIDallocationId,JobIDjobId){Preconditions.checkState(state==SlotState.FREE);this.state=SlotState.PENDING;this.allocationId=allocationId;this.jobId=jobId;}// PENDING → ALLOCATEDpublicvoidmarkAllocated(){Preconditions.checkState(state==SlotState.PENDING);this.state=SlotState.ALLOCATED;}}

Slot 状态转换有严格的前置条件检查,防止非法状态转换。


四、无空闲 Slot 时的等待机制

当 SlotManager 找不到 FREE Slot 时,不会立即失败,而是创建 PendingSlotRequest 等待:

privatevoidcheckTimeoutSlotRequests(){longcurrentTime=System.currentTimeMillis();longtimeout=configuration.getSlotRequestTimeout().toMillis();Iterator<Map.Entry<AllocationID,PendingSlotRequest>>iterator=pendingSlotRequests.entrySet().iterator();while(iterator.hasNext()){PendingSlotRequestrequest=iterator.next().getValue();if(currentTime-request.getCreationTime()>timeout){// 超时,取消请求iterator.remove();notifySlotRequestTimeout(request);}}// 继续下一次检查if(running){scheduledExecutor.schedule(this::checkTimeoutSlotRequests,timeout,TimeUnit.MILLISECONDS);}}

等待机制的关键点:

  • PendingSlotRequest 记录创建时间,用于超时判断
  • 定期检查超时,默认5分钟,超时后取消请求并通知 JobMaster
  • 新 TaskManager 注册或 Slot 释放时,尝试为等待中的请求分配 Slot
  • 同时通知 ResourceManager 通过 Driver 申请新容器,启动新 TaskManager

五、关键配置与最佳实践

关键配置

# Slot 请求超时时间slotmanager.request-timeout:5min# TaskManager 超时时间slotmanager.taskmanager-timeout:30s# 每个 TaskManager 的 Slot 数taskmanager.numberOfTaskSlots:1# 调度器类型jobmanager.scheduler:adaptive# 故障恢复策略jobmanager.execution.failover-strategy:region

最佳实践

  1. 合理配置 Slot 数:根据 TM 资源配置 numberOfTaskSlots,每个 Slot 配置足够的 CPU 和内存。
  2. 监控 PendingSlotRequest:关注等待中的 Slot 请求队列长度,队列堆积说明资源不足。
  3. 使用 Slot Sharing:默认开启,多个算子共享一个 Slot,提高资源利用率。
  4. 配置合适超时:大作业可调大 request-timeout 到 10-15 分钟,避免频繁超时。
  5. 使用 adaptive 调度器:支持动态调整并行度,根据可用 Slot 数自动调整。
  6. 避免 Slot 泄漏:确保作业正常完成时 Slot 被正确释放,定期检查 Slot 状态。

六、总结

Flink 资源申请流程是一个四组件协作的过程:SlotPool 发起请求 → ResourceManager 转发协调 → SlotManager 执行分配 → TaskExecutor 确认分配。核心方法 requestSlot() 贯穿四个组件,通过 RPC 调用链完成 Slot 分配。Slot 状态机(FREE → PENDING → ALLOCATED → RELEASED)保证分配的原子性和一致性。无空闲 Slot 时通过 PendingSlotRequest 等待,超时机制防止请求无限等待。理解这个流程有助于排查 Slot 请求超时、分配失败、资源不足等常见问题。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!