Ray Runtime Env 架构解析:从 RuntimeEnvAgent 到 Worker 启动的完整链路
2026/9/21 1:20:04 网站建设 项目流程

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/ray

Runtime Env(运行时环境)是 Ray 用于为 Job、Task 与 Actor 动态准备 Python 依赖、工作目录与运行环境的机制,也是 Ray 作为 AI 计算引擎实现"环境即配置"的核心能力。本文以仓库内 python/ray/runtime_env/ARCHITECTURE.md 为骨架,结合 Python 侧实现(RuntimeEnvAgentRuntimeEnvRuntimeEnvPlugin)与 C++ 侧实现(Worker Pool、Agent Manager、GCS 引用计数),完整梳理 runtime env 的架构设计、插件机制、创建/删除流程、Worker 进程启动链路以及缓存与垃圾回收策略,帮助你从源码层面理解 Ray 依赖管理的工作原理与关键调优参数。

架构总览:谁在负责创建 Runtime Env

Runtime 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 Sharing

RuntimeEnvAgent与 Raylet 进程命运共享(fate-share):如果 Dashboard Agent 失败,runtime env 创建就会失败,Raylet 将无法为其 Worker 准备运行环境。文档明确解释了这一设计动机:

  • 简化故障模型——两者要么都活着,要么都挂掉,无需处理"Agent 挂了但 Raylet 还活着"的中间态;
  • Runtime Env 是调度任务与 Actor 的核心组件,与 Raylet 绑定可保证 Worker 启动所依赖的环境准备服务始终可用。

产物:磁盘文件 + RuntimeEnvContext

一次成功的 runtime env 创建产生两类产物:

  1. 磁盘上的一组文件:例如安装好的 pip 包、conda 环境、从 GCS 下载解压的working_dir文件等;
  2. 内存中的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.executable
  • override_worker_entrypoint:覆盖 Worker 入口脚本路径(容器场景下宿主机与容器内default_worker.py路径可能不同);
  • java_jars:Java Worker 需要加入 classpath 的 jar 路径列表。

插件机制:所有 Runtime Env 选项都是插件

文档明确指出:runtime env 的所有选项(working_dirpipconda等)都实现为遵循 RayRuntimeEnvPlugin接口的插件。该接口定义在 python/ray/_private/runtime_env/plugin.py,属于DeveloperAPI,包含安装、删除、更新RuntimeEnvContext三个阶段的公开方法:

方法作用说明
validate(runtime_env_dict)校验用户传入的该插件字段在安装 runtime env 时被调用,校验失败抛ValueError
get_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 缓存:

  • 优先级取值区间为 0~100(见 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_needed(plugin.py)封装了"按 URI 查缓存→命中则复用并mark_used,未命中则调用plugin.create并写入缓存→最后modify_context"的标准流程。

用户侧入口:RuntimeEnv 与 RuntimeEnvConfig

用户通过 python/ray/runtime_env/runtime_env.py 中的RuntimeEnv@PublicAPI(stability="stable"))声明运行环境,可用字段包括:py_modulespy_executableworking_dircondapipuvcontainerenv_varsconfigworker_process_setup_hookimage_uri等。其构造器会执行关键校验:

  • condapipuv三者不能同时指定(源码在 runtime_env.py 直接抛ValueError);
  • container只能单独使用,或与configenv_vars组合(runtime_env.py);
  • 指定pip/conda时会自动注入_ray_commit以保证与集群 Ray 版本兼容。

RuntimeEnvConfig提供三个配置项:

字段默认值说明
setup_timeout_seconds600每次 runtime env 创建的安装超时(秒),-1表示禁用超时,不允许设置为其他 ≤0 的值
eager_installTrue是否在ray.init()时、Worker 被租用之前就在集群上预装 runtime env
log_files[]需要呈现在 Dashboard 上的该 runtime env 日志文件列表

注意:RuntimeEnvConfig不参与 runtime env 的 hash 计算(见其类 docstring),因此配置不同但选项相同的两个 runtime env 在缓存层面被视为同一个。

Worker Pool:按 Runtime Env Hash 复用 Worker 进程

Raylet 的 Worker Pool(src/ray/raylet/worker_pool.cc)负责 Worker 进程的缓存与新建。其工作逻辑为:

  1. 调度任务时,该任务的TaskSpec中包含其 runtime env spec;
  2. Worker Pool 将该 runtime env spec 的hash与所有正在运行的 Worker 的 hash 进行比对;
  3. 若存在 hash 相同的 Worker,任务直接复用到该已有 Worker 进程;否则启动一个新的 Worker 进程。

这一按 hash 匹配的设计使得多个使用相同 runtime env 的任务/ Actor 可以共享同一个 Worker 进程,避免重复安装依赖与重复建进程的开销。RuntimeEnvserialize()方法会以sort_keys=True的 JSON 序列化结果参与该 hash 计算(runtime_env.py)。

创建与删除:Raylet ↔ Agent 的 gRPC 链路

RuntimeEnvAgent向 Raylet 暴露创建与删除 runtime env 的gRPC 端点,核心方法为GetOrCreateRuntimeEnvDeleteRuntimeEnvIfPossible(均在 runtime_env_agent.py 中实现)。

  • Raylet 侧的Agent Manager(src/ray/raylet/agent_manager.cc)运行在 Raylet 进程内,管理与 agent 的连接,并调用上述创建/删除端点;
  • Worker Pool 持有 Agent Manager 的引用,在需要新 Worker 进程或 Worker 进程被移除时,向其发送CreateRuntimeEnvIfNeededDeleteRuntimeEnvIfPossible请求。

GetOrCreateRuntimeEnv的完整处理流程(结合 runtime_env_agent.py):

  1. 反序列化request.serialized_runtime_envRuntimeEnv对象;
  2. 增加引用ReferenceTable.increase_reference()
  3. 加锁去重:每个序列化 env 对应一个asyncio.Lock,防止同一 env 被并发安装;
  4. 查环境级缓存_env_cache:命中成功结果直接返回序列化的RuntimeEnvContext;命中失败结果则回滚引用并返回错误;
  5. 按优先级创建插件:先创建working_dir(它特殊,需先于其他插件存在,其他插件可在其目录下工作),随后按优先级依次执行其余插件(runtime_env_agent.py);
  6. 带重试的超时控制_create_runtime_env_with_retrysetup_timeout_seconds作为单次尝试超时(-1时禁用超时),失败后重试,超时错误信息会提示用户通过runtime_env={"config": {"setup_timeout_seconds": 1800}, ...}调大超时(runtime_env_agent.py);
  7. 成功后将序列化 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,其执行流程为:

  1. setup_worker.py反序列化 Raylet 传入的RuntimeEnvContext
  2. 调用其exec_worker方法;
  3. exec_worker(实现在 context.py)依次完成:
    • 通过update_envsRuntimeEnvContext.env_vars注入进程环境变量;
    • 按语言组装可执行命令:Python 为exec <py_executable>,Java 则构造-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_dirpy_modules上传的包:

  • 引用按**包(URI)**维度追踪,引用计数管理实现在 src/ray/common/runtime_env_manager.cc;
  • 引用递增时机:某个 driver 启动并使用该 URI;某个 detached actor 启动并使用该 URI;
  • 引用递减时机:该 driver 退出、该 detached actor 退出;
  • 引用计数归零时,文件被删除

Ray Jobs API 与 Ray Client 的"临时引用"

当用户指定本地目录作为working_dirpy_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.py:f"RAY_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~_5:conda/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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询