Ray 编程式集群扩容:深入解析 ray.autoscaler.sdk.request_resources
Ray 编程式集群扩容深入解析 ray.autoscaler.sdk.request_resources【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray导读在 Ray 分布式集群中自动扩缩容Autoscaler默认根据任务调度需求动态调整节点规模但存在扩容速度上限与调度延迟。ray.autoscaler.sdk.request_resources()提供了一条编程式通道在 Ray 程序内部直接向 Autoscaler 下达资源请求指令让集群立即扩容到指定规模CPU 数量或特定资源 Bundle 形状绕过常规扩容速度约束。阅读本文后你将掌握该 API 的完整签名与参数语义、持久化行为、底层实现链路V1 内部 KV 通道与 V2 GCS RPC 两套机制以及如何在真实作业提交前预扩容集群的实战方案。本文依据仓库文档 reference.rst 展开并对照 sdk.py 等源码实现进行深度印证。一、Programmatic Cluster Scaling什么是编程式集群扩容Ray 的 Autoscaler 在后台周期性运行依据集群负载运行中的任务、待调度的 Placement Group 等决定是否新增或回收节点。这种被动式扩缩容在大部分场景下足够但在以下场景中会带来不可接受的延迟即将提交一批大规模并行任务希望任务到达时节点已经就绪需要特定资源形状如 GPU CPU 组合的节点希望提前预留希望突破upscaling_speed等常规扩容速率限制实现立即扩容。ray.autoscaler.sdk.request_resources()正是为解决这类问题而生。正如 reference.rst 所描述的在一个 Ray 程序内部你可以通过一次request_resources()调用命令 Autoscaler 将集群扩缩容到期望的规模集群会立即尝试扩容以容纳所请求的资源绕过正常的扩容速度限制bypassing normal upscaling speed constraints。该 API 属于 Ray 的 DeveloperAPI开发者 API位于 sdk.py 中并通过 sdk/init.py 从ray.autoscaler.sdk命名空间公开导出。二、API 签名与参数详解request_resources()的完整签名如下见 sdk.pyDeveloperAPI def request_resources( num_cpus: Optional[int] None, bundles: Optional[List[dict]] None, bundle_label_selectors: Optional[List[dict]] None, ) - None:2.1 num_cpus按 CPU 总量请求类型Optional[int]默认None。语义请求集群扩容确保有这么多 CPU 资源可用。该请求是持久的persistent直到下一次request_resources()调用覆盖它为止。2.2 bundles按资源形状请求类型Optional[List[dict]]默认None。语义请求集群扩容确保这一组资源形状resource shapes能被容纳。同样具有持久性直到被下一次调用覆盖。每个 bundle 是一个资源字典键为资源名字符串值为数值int 或 float。2.3 bundle_label_selectors按节点标签筛选类型Optional[List[dict]]默认None。语义一组标签选择器label selectors按索引与bundles列表一一对应第 i 个选择器作用于第 i 个 bundle。没有标签要求的 bundle 对应空字典{}。标签选择器由零个或多个键值对组成键是标签名值形如in(A100)、!in(A100)等操作符表达式。约束一旦提供该参数bundles必须同时提供且两者长度必须相等源码在 sdk.py 中显式校验。三、核心行为语义3.1 立即扩容绕过扩容速度限制与 Autoscaler 按负载自动决策不同request_resources()的请求会被视为集群的资源需求约束resource constraintAutoscaler 会立即尝试满足不受常规upscaling_speed的逐步放大限制。3.2 考虑现有资源使用而非简单累加这是最容易误解的一点。函数 docstring 给出了明确示例见 sdk.py假设你调用request_resources(num_cpus100)当前已有 45 个运行中的任务每个任务占用 1 个 CPU。那么集群会新增节点使得最多 100 个任务可以并发运行而不会新增到能容纳 145 个任务。也就是说请求的是目标可用资源总量Autoscaler 在计算需要新增的节点时会扣除已有的资源占用。3.3 结果是提示而非精确保证docstring 同样强调该调用对 Autoscaler 而言只是一个提示hint。最终集群实际大小可能略大于或略小于预期具体取决于内部的装箱bin packing算法以及max_workers最大节点数限制。3.4 持久性后一次调用覆盖前一次无论是num_cpus还是bundles形式的请求都会一直生效直到下一次request_resources()调用发出新请求将其覆盖。测试用例 test_sdk.py 中明确验证了再次请求会覆盖之前的请求这一行为。四、完整代码示例以下示例均来自函数 docstring见 sdk.py可直接在 Ray 程序中运行from ray.autoscaler.sdk import request_resources # 请求 1000 个 CPU 可用 request_resources(num_cpus1000) # 请求 64 个 CPU同时容纳一个 1-GPU / 4-CPU 的任务 request_resources(num_cpus64, bundles[{GPU: 1, CPU: 4}]) # 等价于请求 num_cpus3三个 CPU1 的 bundle request_resources(bundles[{CPU: 1}, {CPU: 1}, {CPU: 1}]) # 两个 num_cpus1 的 bundle # 第一个要求节点标签 accelerator-type in(A100) # 第二个要求节点标签 market-type 为 spot request_resources( bundles[{CPU: 1}, {CPU: 1}], bundle_label_selectors[ {accelerator-type: in(A100)}, {market-type: spot}, ], )典型的使用场景是在集群的 Head 节点上、提交大量ray.remote任务之前调用该函数确保资源迅速就绪源码注释见 commands.pyThis function is to be called e.g. on a node before submitting a bunch of ray.remote calls to ensure that resources rapidly become available.。五、参数校验规则源码级在 sdk.py 的request_resources()实现中参数会先经过严格的类型校验然后再下发给底层命令层参数校验规则违反时的异常num_cpus必须为intNone允许TypeError: num_cpus should be of type int.bundles必须是List每个元素必须是Dict键必须为str值必须为int或floatbool会被显式拒绝避免{CPU: True}静默等价于{CPU: 1}TypeError: each bundle should be a Dict.等bundle_label_selectors提供时bundles必须同时提供且两者长度相等每个选择器必须是字符串键值对字典且键值均为str选择器本身需通过标签选择器语法校验ValueError系列这些校验在 v2/tests/test_sdk.py 的对应测试用例中得到验证包括对合法/非法标签选择器如缺失 token 导致的认证错误的处理。六、底层实现V1 内部 KV 通道 与 V2 GCS RPCrequest_resources()的用户态实现位于 commands.py 的request_resources()函数约 L193-L249它会根据 Autoscaler 版本走两条完全不同的链路6.1 前置检查if not ray.is_initialized(): raise RuntimeError(Ray is not initialized yet)调用前必须已初始化 Ray即运行在ray.init()之后、或已加入集群的环境中否则抛出RuntimeError。6.2 请求归一化无论用户传入num_cpus还是bundles都会被归一化为统一的内部结构{resources: {...}, label_selector: {...}}num_cpusN被展开为 N 个{resources: {CPU: 1}, label_selector: {}}每个bundles[i]对应一个{resources: bundles[i], label_selector: bundle_label_selectors[i] or {}}。6.3 V1 Autoscaler写入内部 KV 通道在 V1 路径下归一化后的请求被转换为纯资源字典列表通过_internal_kv_put写入名为AUTOSCALER_RESOURCE_REQUEST_CHANNEL的内部 KV 通道overwriteTrue即覆盖写入_internal_kv_put( AUTOSCALER_RESOURCE_REQUEST_CHANNEL, json.dumps(to_request_v1), overwriteTrue, )在 Head 节点的 Autoscaler Monitor 进程中update_resource_requests()方法见 monitor.py 约 L371-L383周期性从该 KV 通道读取请求反序列化后调用self.load_metrics.set_resource_requests(resource_request)注入负载指标LoadMetrics。这些资源请求随后与真实任务负载一起参与扩容决策使 Autoscaler 将请求的容量视为必须满足的集群资源约束。写入的请求会一直保留在 KV 中这正是持久化直到下次覆盖语义的来源。6.4 V2 AutoscalerGCS RPC 直连当集群使用 V2 Autoscaler 时代码会切换到新的格式通过 GCS RPC 直接提交资源约束gcs_address internal_kv_get_gcs_client().address request_cluster_resources(gcs_address, to_request)其底层实现在 v2/sdk.py 的request_cluster_resources()中将{resources, label_selector}字典归一化为ResourceRequest命名元组按资源形状与标签选择器聚合使用Counter对请求去重并统计数量例如 100 个{CPU: 1}请求会被聚合成 1 个 bundle count100通过GcsClient(gcs_address).request_cluster_resource_constraint(bundles, label_selectors, counts, timeout_stimeout)将资源约束异步传递给 GCSGCS 再将约束转发给 Autoscaler由 Autoscaler 尝试供给请求中的最小 bundle 集合。文档中说明见 v2/sdk.py 的 docstring如果集群当前已拥有to_request中的资源则该调用为 no-op后续通过该 API 提交的请求会覆盖之前的请求。6.5 版本判定两条路径的切换由 v2/utils.py 中的is_autoscaler_v2()判断也就是说调用方无需关心内部实现差异ray.autoscaler.sdk.request_resources对外保持统一接口。七、测试用例对行为的印证仓库中 test_sdk.py 的测试直接印证了上文描述的行为test_request_cluster_resources_basic请求[{CPU: 1}]后通过get_cluster_resource_state轮询断言集群资源约束中出现{CPU: 1}且 count1随后请求[{CPU: 2, GPU: 1}, {CPU: 1}]验证新请求覆盖旧请求再发送 100 个{CPU: 1}请求断言被按形状聚合成 count100的单一约束。test_request_cluster_resources_with_label_selectors对[{CPU: 1}, {GPU: 1, CPU: 2}]分别施加{region: us-west1}与{accelerator-type: !in(A100)}标签选择器验证标签选择器随 bundle 一并进入集群资源约束状态。八、实战建议与注意事项8.1 与min_workers的职责边界min_workers是集群配置YAML中声明的最小工作节点数属于静态配置在集群创建时即固定而request_resources()是运行时编程式的资源需求声明适合无法预先静态配置、需要按作业动态决定容量的场景。两者可以组合使用用min_workers保证基础容量用request_resources()应对突发的大规模作业。8.2 配合max_workers上限request_resources()只是给 Autoscaler 的提示最终扩容仍受集群配置中max_workers的限制。若请求的资源量超过上限集群只会扩容到max_workers允许的规模因此请合理设置max_workers避免容量永远无法满足请求。8.3 与 Placement Group 的取舍如果需求是为特定任务确保节点上能调度一组资源也可以考虑使用 Placement Groupray.util.placement_group。区别在于Placement Group 由调度器GCS负责具体资源预留与放置而request_resources()是纯容量层面的预扩容提示不绑定具体任务。批量任务提交前的预热场景更适合后者。8.4 调用位置建议在 Head 节点或已初始化 Ray 的驱动进程中提交大量ray.remote任务之前调用以隐藏节点启动Node Provider 创建实例、安装运行环境、加入集群的耗时一旦资源就绪任务即可立即调度。同时记住请求是持久的任务结束后若不再需要额外容量应再次调用request_resources()如传入空请求或依赖 Autoscaler 的空闲节点回收逻辑idle_timeout_minutes来释放。结语ray.autoscaler.sdk.request_resources()是 Ray 集群编程式扩容的核心入口以极简的三参数接口num_cpus/bundles/bundle_label_selectors覆盖了总量扩容形状扩容标签定向扩容三类诉求。其底层在 V1 与 V2 Autoscaler 中分别走内部 KV 通道与 GCS RPC但对外保持一致的持久化、覆盖式、考虑存量负载的语义。理解这些行为细节能帮助你在真实集群中精准控制扩缩容节奏为大规模作业提前铺好容量。【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考