资讯动态

Ray Runtime Env 架构解析:从 RuntimeEnvAgent 到 Worker 启动的完整链路

发布时间:2026/9/21 1:20:05 来源:尧图企业网站定制
Ray Runtime Env 架构解析从 RuntimeEnvAgent 到 Worker 启动的完整链路【免费下载链接】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/rayRuntime Env运行时环境是 Ray 用于为 Job、Task 与 Actor 动态准备 Python 依赖、工作目录与运行环境的机制也是 Ray 作为 AI 计算引擎实现环境即配置的核心能力。本文以仓库内 python/ray/runtime_env/ARCHITECTURE.md 为骨架结合 Python 侧实现RuntimeEnvAgent、RuntimeEnv、RuntimeEnvPlugin与 C 侧实现Worker Pool、Agent Manager、GCS 引用计数完整梳理 runtime env 的架构设计、插件机制、创建/删除流程、Worker 进程启动链路以及缓存与垃圾回收策略帮助你从源码层面理解 Ray 依赖管理的工作原理与关键调优参数。架构总览谁在负责创建 Runtime EnvRuntime Env 的创建由运行在集群每个节点上的Dashboard Agent进程即RuntimeEnvAgent负责其核心实现位于 python/ray/_private/runtime_env/agent/runtime_env_agent.py。RuntimeEnvAgent是一个 RPC 服务器以 Dashboard Agent 的形式随节点启动向本节点的 Raylet 提供 runtime env 的创建与删除能力。它持有 GCS Client用于从 GCS 拉取用户上传的包、插件管理器RuntimeEnvPluginManager、环境级缓存_env_cache与 URI 引用表ReferenceTable等关键组件。与 Raylet 的 Fate SharingRuntimeEnvAgent与 Raylet 进程命运共享fate-share如果 Dashboard Agent 失败runtime env 创建就会失败Raylet 将无法为其 Worker 准备运行环境。文档明确解释了这一设计动机简化故障模型——两者要么都活着要么都挂掉无需处理Agent 挂了但 Raylet 还活着的中间态Runtime Env 是调度任务与 Actor 的核心组件与 Raylet 绑定可保证 Worker 启动所依赖的环境准备服务始终可用。产物磁盘文件 RuntimeEnvContext一次成功的 runtime env 创建产生两类产物磁盘上的一组文件例如安装好的 pip 包、conda 环境、从 GCS 下载解压的working_dir文件等内存中的RuntimeEnvContextPython 对象定义在 python/ray/_private/runtime_env/context.py。RuntimeEnvContext被序列化后传给 Raylet在 Raylet 为该 runtime env 启动新 Worker 进程时使用详见下文Worker 进程启动链路。从源码看该对象包含四个字段command_prefix启动命令前缀例如conda activate some_env之前要执行的前置命令列表env_vars需要注入的环境变量字典py_executable使用的 Python 解释器路径默认取当前sys.executableoverride_worker_entrypoint覆盖 Worker 入口脚本路径容器场景下宿主机与容器内default_worker.py路径可能不同java_jarsJava Worker 需要加入 classpath 的 jar 路径列表。插件机制所有 Runtime Env 选项都是插件文档明确指出runtime env 的所有选项working_dir、pip、conda等都实现为遵循 RayRuntimeEnvPlugin接口的插件。该接口定义在 python/ray/_private/runtime_env/plugin.py属于DeveloperAPI包含安装、删除、更新RuntimeEnvContext三个阶段的公开方法方法作用说明validate(runtime_env_dict)校验用户传入的该插件字段在安装 runtime env 时被调用校验失败抛ValueErrorget_uris(runtime_env)返回该插件涉及的所有 URI用于缓存与引用计数无 URI 时每次创建都不查缓存create(uri, runtime_env, context, logger)创建并安装 runtime env在 runtime env agent 的安装阶段被调用返回值表示该次安装占用的磁盘空间字节modify_context(uris, runtime_env, context, logger)修改 Worker 启动行为例如向启动命令前置cd dir、追加环境变量delete_uri(uri, logger)按 URI 删除 runtime env返回值表示回收的磁盘空间在 runtime_env_agent.py 中11 个内置插件被依次实例化并注册self._pip_plugin PipPlugin(self._runtime_env_dir) self._uv_plugin UvPlugin(self._runtime_env_dir) self._conda_plugin CondaPlugin(self._runtime_env_dir) self._py_modules_plugin PyModulesPlugin(self._runtime_env_dir, self._gcs_client) self._py_executable_plugin PyExecutablePlugin() self._java_jars_plugin JavaJarsPlugin(self._runtime_env_dir, self._gcs_client) self._working_dir_plugin WorkingDirPlugin(self._runtime_env_dir, self._gcs_client) self._container_plugin ContainerPlugin(temp_dir) self._nsight_plugin NsightPlugin(self._runtime_env_dir) self._rocprof_sys_plugin RocProfSysPlugin(self._runtime_env_dir) self._image_uri_plugin get_image_uri_plugin_cls()(temp_dir)这些插件的实现分布在 python/ray/_private/runtime_env/ 目录下如 working_dir.py、pip.py、conda.py 等每个插件拥有独立的 schema 文件例如 python/ray/runtime_env/schemas/pip_schema.json 与 working_dir_schema.json。插件优先级与第三方插件加载RuntimeEnvPluginManager负责加载插件并管理每个插件的 URI 缓存优先级取值区间为 0100见 constants.py 中的RAY_RUNTIME_ENV_PLUGIN_MIN_PRIORITY/MAX_PRIORITY默认优先级为 10数字越小越先被创建sorted_plugin_setup_contexts()按优先级升序排列第三方插件通过环境变量RAY_RUNTIME_ENV_PLUGINS传入 JSON 配置形如[{class: xxx.xxx_plugin, priority: 10}]由RuntimeEnvPluginManager在 agent 启动时动态加载plugin.py每个插件自动关联一个URICache缓存上限来自环境变量RAY_RUNTIME_ENV_插件名_CACHE_SIZE_GB默认 10 GB对应文档中的RAY_RUNTIME_ENV_WORKING_DIR_CACHE_SIZE_GB等。create_for_plugin_if_neededplugin.py封装了按 URI 查缓存→命中则复用并mark_used未命中则调用plugin.create并写入缓存→最后modify_context的标准流程。用户侧入口RuntimeEnv 与 RuntimeEnvConfig用户通过 python/ray/runtime_env/runtime_env.py 中的RuntimeEnvPublicAPI(stabilitystable)声明运行环境可用字段包括py_modules、py_executable、working_dir、conda、pip、uv、container、env_vars、config、worker_process_setup_hook、image_uri等。其构造器会执行关键校验conda、pip、uv三者不能同时指定源码在 runtime_env.py 直接抛ValueErrorcontainer只能单独使用或与config、env_vars组合runtime_env.py指定pip/conda时会自动注入_ray_commit以保证与集群 Ray 版本兼容。RuntimeEnvConfig提供三个配置项字段默认值说明setup_timeout_seconds600每次 runtime env 创建的安装超时秒-1表示禁用超时不允许设置为其他 ≤0 的值eager_installTrue是否在ray.init()时、Worker 被租用之前就在集群上预装 runtime envlog_files[]需要呈现在 Dashboard 上的该 runtime env 日志文件列表注意RuntimeEnvConfig不参与 runtime env 的 hash 计算见其类 docstring因此配置不同但选项相同的两个 runtime env 在缓存层面被视为同一个。Worker Pool按 Runtime Env Hash 复用 Worker 进程Raylet 的 Worker Poolsrc/ray/raylet/worker_pool.cc负责 Worker 进程的缓存与新建。其工作逻辑为调度任务时该任务的TaskSpec中包含其 runtime env specWorker Pool 将该 runtime env spec 的hash与所有正在运行的 Worker 的 hash 进行比对若存在 hash 相同的 Worker任务直接复用到该已有 Worker 进程否则启动一个新的 Worker 进程。这一按 hash 匹配的设计使得多个使用相同 runtime env 的任务/ Actor 可以共享同一个 Worker 进程避免重复安装依赖与重复建进程的开销。RuntimeEnv的serialize()方法会以sort_keysTrue的 JSON 序列化结果参与该 hash 计算runtime_env.py。创建与删除Raylet ↔ Agent 的 gRPC 链路RuntimeEnvAgent向 Raylet 暴露创建与删除 runtime env 的gRPC 端点核心方法为GetOrCreateRuntimeEnv与DeleteRuntimeEnvIfPossible均在 runtime_env_agent.py 中实现。Raylet 侧的Agent Managersrc/ray/raylet/agent_manager.cc运行在 Raylet 进程内管理与 agent 的连接并调用上述创建/删除端点Worker Pool 持有 Agent Manager 的引用在需要新 Worker 进程或 Worker 进程被移除时向其发送CreateRuntimeEnvIfNeeded与DeleteRuntimeEnvIfPossible请求。GetOrCreateRuntimeEnv的完整处理流程结合 runtime_env_agent.py反序列化request.serialized_runtime_env为RuntimeEnv对象增加引用ReferenceTable.increase_reference()加锁去重每个序列化 env 对应一个asyncio.Lock防止同一 env 被并发安装查环境级缓存_env_cache命中成功结果直接返回序列化的RuntimeEnvContext命中失败结果则回滚引用并返回错误按优先级创建插件先创建working_dir它特殊需先于其他插件存在其他插件可在其目录下工作随后按优先级依次执行其余插件runtime_env_agent.py带重试的超时控制_create_runtime_env_with_retry以setup_timeout_seconds作为单次尝试超时-1时禁用超时失败后重试超时错误信息会提示用户通过runtime_env{config: {setup_timeout_seconds: 1800}, ...}调大超时runtime_env_agent.py成功后将序列化 context 写入_env_cache并返回。DeleteRuntimeEnvIfPossible则执行ReferenceTable.decrease_reference()递减引用引用归零时由回调清理缓存。此外agent 还提供GetRuntimeEnvsInfo端点可查询本节点上各 runtime env 的引用计数、创建耗时与成功状态。Worker 进程启动链路从 setup_worker 到 execvp新的 Worker 进程由 Raylet 通过 python/ray/_private/services.py 启动实际入口是 python/ray/_private/workers/setup_worker.py其执行流程为setup_worker.py反序列化 Raylet 传入的RuntimeEnvContext调用其exec_worker方法exec_worker实现在 context.py依次完成通过update_envs将RuntimeEnvContext.env_vars注入进程环境变量按语言组装可执行命令Python 为exec py_executableJava 则构造-cpclasspath含 Ray jars 与java_jars字段指定的 jar若设置了override_worker_entrypoint替换 Worker 入口脚本路径容器场景必备将command_prefix如conda activate some_env的前置命令拼接到命令最前最终调用os.execvp(bash, ...)以 bash 执行拼接后的完整命令用 exec 替换当前进程完成 Worker 进程的启动context.py。这一设计保证了 Worker 进程在启动瞬间就处于目标 runtime env 中正确的解释器、环境变量、激活的 conda/venv 环境无需 Worker 运行后再切换环境。缓存与垃圾回收两级 GC 机制Runtime Env 的清理分为两级GCS 内部 KV 的引用计数回收与各节点本地磁盘的回收。GCS Internal KV 垃圾回收Head 节点这一级处理存储在 Head 节点内部 KV 中的文件典型代表是用户通过working_dir或py_modules上传的包引用按**包URI**维度追踪引用计数管理实现在 src/ray/common/runtime_env_manager.cc引用递增时机某个 driver 启动并使用该 URI某个 detached actor 启动并使用该 URI引用递减时机该 driver 退出、该 detached actor 退出引用计数归零时文件被删除。Ray Jobs API 与 Ray Client 的临时引用当用户指定本地目录作为working_dir或py_modules时Ray 会将其打包为 zip 并以上传 URI 的形式存入 GCS。但如前所述引用计数通常只在 driver或 detached actor启动时才递增——而 Ray Jobs API 与 Ray Client 的上传发生在 driver 启动之前。为防止文件在 driver 启动前被误回收Ray 在上传 URI 时添加一个特殊的临时引用temporary reference。该引用在可配置的超时后移除超时由 Head 节点上的环境变量RAY_RUNTIME_ENV_TEMPORARY_REFERENCE_EXPIRATION_S控制默认 600 秒。本地节点垃圾回收所有节点这一级处理存储在各节点磁盘上的文件例如已安装的 pip 包、从 GCS 或远程 URI 下载解压的working_dir文件引用由每个节点上的 runtime env agent 进程追踪各 agent 在内存中维护独立的引用表即上文提到的ReferenceTable见 runtime_env_agent.py引用在 runtime env 创建时递增、删除时递减文件并非引用归零就立即删除而是等到引用计数归零且缓存大小超过最大缓存上限两个条件同时满足时才被清理每个 runtime_env 字段拥有独立的缓存大小上限默认10 GB可通过各节点上的环境变量RAY_RUNTIME_ENV_字段_CACHE_SIZE_GB配置例如RAY_RUNTIME_ENV_WORKING_DIR_CACHE_SIZE_GB。该逻辑对应插件管理器为每个插件创建URICache时的实现plugin.pyfRAY_RUNTIME_ENV_{plugin.name}_CACHE_SIZE_GB.upper()未设置时默认取 10。另外ReferenceTable内部还维护了序列化 runtime env → 引用计数与URI → 引用计数两张表前者用于环境级缓存_env_cache的清理失败结果还会按BAD_RUNTIME_ENV_CACHE_TTL_SECONDS额外缓存一段时间见unused_runtime_env_processor后者用于各插件 URI 缓存的mark_unused标记。source_process为client_server的引用会被排除不计数以避免 Ray Client 场景下的 URI 泄漏。测试与验证Runtime Env 功能的测试集中在文件名匹配test_runtime_env*的测试文件中位于 python/ray/tests/ 目录命名基本自解释例如test_runtime_env.py核心功能与参数行为test_runtime_env_conda_and_pip.py 及_2~_5conda/pip 组合场景test_runtime_env_container.py容器插件test_runtime_env_env_vars.py环境变量注入test_runtime_env_failure.py失败与超时路径。若需在本地复现或调试 runtime env 相关行为可参考这些测试文件定位对应的插件与 agent 代码路径。关键参数速查参数位置/环境变量默认值说明setup_timeout_secondsruntime_env[config]600 秒单次安装超时-1禁用eager_installruntime_env[config]True是否在ray.init()时预装RAY_RUNTIME_ENV_TEMPORARY_REFERENCE_EXPIRATION_SHead 节点环境变量600 秒Jobs/Ray Client 临时引用的存活时间RAY_RUNTIME_ENV_字段_CACHE_SIZE_GB各节点环境变量10 GB每个字段的本地缓存上限RAY_RUNTIME_ENV_PLUGINS环境变量无第三方插件 JSON 配置以上配置均可结合 python/ray/runtime_env/runtime_env.py 与 python/ray/_private/runtime_env/constants.py 中的源码定义进一步核实与扩展。小结从架构上看Ray Runtime Env 是一个典型的Agent 驱动 插件化 双级引用计数系统RuntimeEnvAgent作为每节点上的执行引擎以插件方式统一处理各类环境选项RuntimeEnvContext作为跨进程传递的环境契约保证 Worker 在 exec 瞬间即处于目标环境Worker Pool 按 spec hash 复用进程GCS 与本地节点分别通过引用计数与容量阈值控制回收时机。理解这条链路有助于你在遇到环境安装超时、缓存膨胀、包意外被清理等问题时快速定位到对应的组件与配置项。【免费下载链接】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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价