ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

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

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 分配成功后返回 AcknowledgeSlotManager 收到确认后调用slot.markAllocated()将 Slot 状态从 PENDING 改为 ALLOCATED。第6步RM 返回 SlotPool 成功。ResourceManager 向 SlotPool 返回SlotRequestSuccess(allocationId, slotId)SlotPool 接收 Slot加入已分配集合通知等待的 Execution 部署 Task。三、源码级别调用链跟踪下面这张图是源码级别跟踪包括四个核心类的关键方法、requestSlot 参数详解、最佳实践和总结。3.1 SlotPoolImpl.requestSlotpublicclassSlotPoolImplimplementsSlotPool{privatefinalMapAllocationID,PendingRequestpendingRequests;privatefinalMapAllocationID,AllocatedSlotallocatedSlots;OverridepublicCompletableFutureSlotRequestSuccessrequestSlot(JobIDjobId,AllocationIDallocationId,ResourceProfileresourceProfile,StringtargetAddress){// 1. 创建 PendingRequest记录请求信息PendingRequestpendingRequestnewPendingRequest(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.requestSlotpublicclassResourceManagerImplWorkerTypeextendsResourceIDRetrievableextendsFencedRpcEndpointResourceManagerIdimplementsResourceManagerWorkerType{privatefinalSlotManagerslotManager;OverridepublicCompletableFutureAcknowledgerequestSlot(JobIDjobId,AllocationIDallocationId,ResourceProfileresourceProfile,StringtargetAddress,ResourceManagerIdresourceManagerId){// 验证 fencing tokenvalidateRunsInMainThread();// 转发给 SlotManager 处理slotManager.requestSlot(jobId,allocationId,resourceProfile);returnCompletableFuture.completedFuture(Acknowledge.get());}}RM 是协调者验证请求合法性后直接转发给 SlotManager不做实际分配逻辑。3.3 SlotManagerImpl.requestSlotpublicclassSlotManagerImplimplementsSlotManager{privatefinalMapResourceID,MapSlotID,TaskManagerSlottaskManagerSlots;privatefinalMapAllocationID,PendingSlotRequestpendingSlotRequests;OverridepublicCompletableFutureAcknowledgerequestSlot(JobIDjobId,AllocationIDallocationId,ResourceProfileresourceProfile){// 1. 查找 FREE SlotTaskManagerSlotfreeSlotfindFreeSlot(resourceProfile);if(freeSlot!null){// 2. 标记 PENDINGfreeSlot.assignPendingAllocation(allocationId,jobId);// 3. 创建 PendingSlotRequestPendingSlotRequestpendingRequestnewPendingSlotRequest(allocationId,jobId,resourceProfile,freeSlot.getSlotId());pendingSlotRequests.put(allocationId,pendingRequest);// 4. 通过 RM 通知 TaskExecutor 分配TaskExecutorGatewaytaskExecutorGatewaytaskManagerConnections.get(freeSlot.getTaskManagerId());returntaskExecutorGateway.requestSlot(freeSlot.getSlotId(),jobId,allocationId,targetAddress,timeout);}else{// 无 FREE Slot创建 PendingSlotRequest 等待PendingSlotRequestpendingRequestnewPendingSlotRequest(allocationId,jobId,resourceProfile,null);pendingSlotRequests.put(allocationId,pendingRequest);// 通知 RM 申请新资源notifyResourceShortage(resourceProfile);returnCompletableFuture.completedFuture(Acknowledge.get());}}}SlotManager 是分配的核心执行者。找到空闲 Slot 时标记 PENDING 并通知 TM无空闲时创建等待请求并通知 RM 申请新容器。3.4 TaskExecutor.requestSlotpublicclassTaskExecutorextendsTaskExecutorFencedRpcEndpoint{privatefinalSlotTableslotTable;OverridepublicCompletableFutureAcknowledgerequestSlot(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(stateSlotState.FREE);this.stateSlotState.PENDING;this.allocationIdallocationId;this.jobIdjobId;}// PENDING → ALLOCATEDpublicvoidmarkAllocated(){Preconditions.checkState(stateSlotState.PENDING);this.stateSlotState.ALLOCATED;}}Slot 状态转换有严格的前置条件检查防止非法状态转换。四、无空闲 Slot 时的等待机制当 SlotManager 找不到 FREE Slot 时不会立即失败而是创建 PendingSlotRequest 等待privatevoidcheckTimeoutSlotRequests(){longcurrentTimeSystem.currentTimeMillis();longtimeoutconfiguration.getSlotRequestTimeout().toMillis();IteratorMap.EntryAllocationID,PendingSlotRequestiteratorpendingSlotRequests.entrySet().iterator();while(iterator.hasNext()){PendingSlotRequestrequestiterator.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最佳实践合理配置 Slot 数根据 TM 资源配置 numberOfTaskSlots每个 Slot 配置足够的 CPU 和内存。监控 PendingSlotRequest关注等待中的 Slot 请求队列长度队列堆积说明资源不足。使用 Slot Sharing默认开启多个算子共享一个 Slot提高资源利用率。配置合适超时大作业可调大 request-timeout 到 10-15 分钟避免频繁超时。使用 adaptive 调度器支持动态调整并行度根据可用 Slot 数自动调整。避免 Slot 泄漏确保作业正常完成时 Slot 被正确释放定期检查 Slot 状态。六、总结Flink 资源申请流程是一个四组件协作的过程SlotPool 发起请求 → ResourceManager 转发协调 → SlotManager 执行分配 → TaskExecutor 确认分配。核心方法 requestSlot() 贯穿四个组件通过 RPC 调用链完成 Slot 分配。Slot 状态机FREE → PENDING → ALLOCATED → RELEASED保证分配的原子性和一致性。无空闲 Slot 时通过 PendingSlotRequest 等待超时机制防止请求无限等待。理解这个流程有助于排查 Slot 请求超时、分配失败、资源不足等常见问题。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进