
Ray Core API 完全指南从ray.init到跨语言调用的分布式编程核心接口【免费下载链接】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 是一个面向 AI 计算场景的分布式运行时其核心价值在于把远程函数调用Task远程对象Object有状态服务Actor等分布式原语以接近本地编程的语法暴露给开发者。本文以仓库文档 doc/source/ray-core/api/core.rst 列出的公开 API 清单为骨架逐层拆解 Ray Core 的六大 API 家族——Core 生命周期、Tasks、Actors、Objects、Runtime Context 与 Cross Language并结合仓库源码如 python/ray/_private/worker.py、python/ray/actor.py、python/ray/remote_function.py说明每个接口的真实行为、默认值与边界条件。读完本文你将能熟练使用这些接口搭建自己的分布式任务、Actor 服务与跨语言调用程序并理解其底层机制。一、Core API集群生命周期与作业级配置Core API 解决的是如何启动、连接、关闭一个 Ray 集群/作业这一最基础的问题对应文档中的ray.init、ray.shutdown、ray.is_initialized、ray.job_config.JobConfig与ray.LoggingConfig。这些函数在 python/ray/init.py 中统一导出并经过AUTO_INIT_APIS/NON_AUTO_INIT_APIS两套集合的区分get、put、wait、get_actor、get_runtime_context等属于可自动触发ray.init()的 API而init、remote、shutdown等则不会自动初始化。1.1ray.init()启动或连接集群ray.init()是 Ray 程序的入口其完整签名定义在 python/ray/_private/worker.py#L1439-L1465。它的职责是连接已存在的 Ray 集群或启动一个新的本地集群并连接之。参数众多这里按用途分组说明参数默认值作用addressNone集群地址解析规则见下文num_cpus虚拟核数分配给每个 raylet 的 CPU 数num_gpus探测到的 GPU 数分配给每个 raylet 的 GPU 数resourcesNone自定义资源名称到数量的字典labelsNone[实验性] 节点的键值标签object_store_memory系统内存 30%上限 shm 大小与 200G对象存储内存字节ignore_reinit_errorFalse重复调用init时是否抑制报错不会重启集群include_dashboardNone依赖存在则启动是否启动 Dashboarddashboard_host/dashboard_portlocalhost / 8265Dashboard 绑定地址与端口8265 被占用时自动找空闲端口job_configNone作业配置JobConfigconfigure_loggingTrue是否允许在此配置日志logging_level/logging_formatlogging.INFO/ 内置格式驱动进程 ray logger 的级别与格式logging_configNone[实验性] 应用于驱动与所有 worker 的日志配置LoggingConfignamespaceNone命名空间逻辑上对作业与命名 Actor 进行分组runtime_envNone作业级运行时环境依赖包、环境变量等enable_resource_isolationFalse通过 cgroupv2 为系统进程预留 CPU/内存需 cgroupv2 且 raylet 有读写权限system_reserved_cpumin(3.0, max(1.0, 0.05 * 核数))为系统进程预留的 CPU 核数可小数system_reserved_memorymin(10GB, max(500MB, 0.10 * 可用内存))为系统进程预留的内存字节proxy_server_urlNone[实验性] 覆盖 Dashboard 后端请求的服务地址address的解析优先级来自init源码 docstring提供具体地址如localhost:port→ 直接连接以ray://前缀连接远程集群需ray[client]不提供地址 → 先查环境变量RAY_ADDRESS再查最近一次启动的集群记录于/tmp/ray/ray_current_cluster都没有则启动新的本地实例传auto→ 与上述流程相同但找不到现有集群时报ConnectionError而非自启传local→ 强制启动新本地实例即使已有本地集群存在。此外还支持大量**kwargs隐藏参数如object_spilling_directory对象落盘目录默认节点 session 目录、_enable_object_reconstruction通过重放任务重建丢失对象、_plasma_directory、_metrics_export_port、_node_name等它们不稳定可能随版本变化。import ray # 最常见的用法自动探测或启动本地集群 ray.init() # 连接现有本地集群不存在则报错 # ray.init(addressauto) # 连接远程集群 # ray.init(addressray://123.45.67.89:10001) # 带资源与日志配置的本地启动 ray.init( num_cpus4, num_gpus1, resources{custom_resource: 2.0}, namespacemy_namespace, logging_level20, )1.2ray.shutdown()与ray.is_initialized()ray.shutdown(*, wait_for_processes: bool False)python/ray/_private/worker.py#L2070断开与集群的连接若该进程是本地集群的启动者还会一并关闭所有相关进程。ray.is_initialized() - boolpython/ray/_private/worker.py#L2483返回当前进程是否已初始化。它是编写库代码时的常用守卫if not ray.is_initialized(): ray.init()1.3ray.job_config.JobConfig与ray.LoggingConfigJobConfig用于封装作业级配置例如为作业设置runtime_env、namespace、worker_process_setup_hook等通过ray.init(job_config...)传入作用于该作业下的所有 worker。LoggingConfig定义于 python/ray/_private/ray_logging/logging_config.py#L68是 [实验性] 的日志配置对象可同时应用于驱动进程与当前作业的所有 worker 进程的根 logger。它支持encoding配置且init内部会读取环境变量RAY_LOGGING_CONFIG_ENCODING作为默认编码并执行_apply()见 python/ray/_private/worker.py#L1631-L1636。开发者可借此统一结构化日志格式便于采集与检索。二、Tasks远程函数原语Tasks 部分对应文档中的ray.remote、ray.remote_function.RemoteFunction.options与ray.cancel。远程函数是 Ray 分布式执行的最小单元用ray.remote装饰的普通函数调用其.remote()方法即可在集群中异步执行。2.1ray.remote把函数/类变为远程对象ray.remote的若干重载定义在 python/ray/_private/worker.py#L3686-L3785其行为因修饰对象不同而分为两种装饰普通函数 → 返回RemoteFunction定义于 python/ray/remote_function.py#L41调用fn.remote(...)提交任务立即返回ObjectRef装饰类 → 返回ActorClass用于创建 Actor见第三节。import ray ray.init() ray.remote def add(a, b): return a b ref add.remote(1, 2) # 异步提交不阻塞 result ray.get(ref) # 阻塞获取结果 print(result) # 3任务在调度时会把参数序列化后放入对象存储任务结果同样以ObjectRef形式写回对象存储由ray.get取回。2.2RemoteFunction.options()任务的调度选项每次提交任务时的资源与调度策略通过RemoteFunction.options(**task_options)配置见 python/ray/remote_function.py#L202。常用选项包括num_cpus/num_gpus单个任务占用的 CPU/GPU 资源resources自定义资源需求num_returns返回多少个ObjectRef配合多返回值任务max_retries任务失败后的重试次数-1表示无限重试scheduling_strategy调度策略如PlacementGroupSchedulingStrategy、SPREAD、PACK等runtime_env任务级运行时环境name任务名用于 Dashboard 观察。ray.remote def heavy_task(x): return x * x ref heavy_task.options(num_cpus2, num_gpus1, max_retries3).remote(10) print(ray.get(ref)) # 1002.3ray.cancel()取消未完成的任务ray.cancel(ref, *, forceFalse, recursiveTrue, ignore_unfinished_tasksFalse)见 python/ray/_private/worker.py#L3505用于取消一个尚未完成的任务forceFalse时向任务所在 worker 发送取消请求任务在下一个中断点退出forceTrue时直接强杀执行该任务的进程recursiveTrue时同时取消该任务派生的子任务取消后对结果的ray.get会抛出TaskCancelledError。三、Actors有状态服务的核心抽象Actors 对应文档中的ray.remote类装饰、ray.actor.ActorClass、ActorClass.options、ray.actor.ActorMethod、ray.actor.ActorHandle、ray.actor.ActorClassInheritanceException、ray.actor.exit_actor、ray.method、ray.get_actor与ray.kill。Actor 与 Task 的本质区别在于Actor 的实例化状态保存在集群中的某个 worker 进程内后续方法调用共享该状态。核心类型定义在 python/ray/actor.py。3.1 定义与创建ActorClass与ActorClass.options用ray.remote装饰类即得到ActorClasspython/ray/actor.py#L1553通过cls.remote(*args)创建 Actor 实例并返回ActorHandleActorClass.options(...)文档中的ray.actor.ActorClass.options在创建前指定 Actor 级资源与配置ray.remote(num_cpus1) class Counter: def __init__(self, start0): self.value start def increment(self): self.value 1 return self.value # 带额外资源的创建 actor Counter.options(num_cpus2, max_restarts3).remote(start10) print(ray.get(actor.increment.remote())) # 11options中常用的 Actor 专属参数还包括max_restarts异常退出后的最大重启次数、max_task_retries、max_pending_calls、concurrency_groups、lifetimedetached表示脱离驱动进程存活等。3.2 方法调用ActorMethod与ray.method通过ActorHandle访问的每个方法是一个ActorMethodpython/ray/actor.py#L848其.remote(*args)异步提交方法调用并返回ObjectRef。ray.methodpython/ray/actor.py#L454 附近的多重重载用于给 Actor 方法附加选项例如num_returns、max_retries、concurrency_groupclass Worker: ray.method(num_returns2) def split(self, s): half len(s) // 2 return s[:half], s[half:] w Worker.remote() ref_a, ref_b w.split.remote(raycore) print(ray.get([ref_a, ref_b])) # [ray, core]3.3 句柄与查询ActorHandle、get_actor、kill、exit_actorActorHandlepython/ray/actor.py#L2274是驱动进程中代表远程 Actor 的句柄可被序列化后在任务、其他 Actor 或客户端之间传递从而实现分布式协作。ray.get_actor(name, namespaceNone)python/ray/_private/worker.py#L3426按名字获取已注册的 Actor 句柄是命名 Actor模式跨作业/跨驱动访问的关键接口结合init(namespace...)可实现作业间的服务发现。ray.kill(actor, *, no_restartTrue)python/ray/_private/worker.py#L3461强制终止 Actor 进程no_restartTrue默认时禁止其自动重启。ray.actor.exit_actor()python/ray/actor.py#L2944在 Actor 方法内部调用主动优雅退出当前 Actor适合处理完队列后自我结束的场景。ActorClassInheritanceExceptionpython/ray/actor.py#L1509是一个TypeError子类当用户违反 Actor 类继承规则例如继承非ray.remote类、或跨层级重复装饰时抛出用于引导正确的 Actor 继承写法。ray.remote class Logger: def log(self, msg): print(msg) ray.init(namespaceapp) Logger.options(namemain_logger, lifetimedetached).remote() # 其他驱动进程中 logger ray.get_actor(main_logger, namespaceapp) ray.get(logger.log.remote(hello from another job)) ray.kill(logger)四、Objects分布式对象存储与异步等待Objects 部分对应文档中的ray.get、ray.wait、ray.put、ray.util.as_completed与ray.util.map_unordered。ObjectRef是 Ray 分布式对象的句柄可跨节点、跨进程传递实际数据存放在集群的对象存储Plasma 内存/磁盘中。4.1ray.put与ray.get显式存入与阻塞取出ray.put(value) - ObjectRefpython/ray/_private/worker.py#L3034将对象显式写入对象存储返回引用只要引用存在对象不会被驱逐。仓库提示相关反模式可参考 ray-core/patterns/return-ray-put.rst函数内先ray.put再返回会破坏零拷贝优化。ray.get(object_refs, *, timeoutNone)python/ray/_private/worker.py#L2881-L2943阻塞直至对象可用传入单个ObjectRef返回单个对象传入列表返回列表且顺序与输入一致timeoutNone无限阻塞timeout0表示对象可用则立即返回否则抛GetTimeoutError其他数值为最大等待秒数任一任务抛异常时立即抛出该异常不再等待其余结果在 async Actor 中调用会产生阻塞事件循环的告警建议改用await object_ref或asyncio.gather(*refs)源码在 python/ray/_private/worker.py#L2947-L2960不允许传入ObjectRefGenerator。4.2ray.wait非全量等待ray.wait(ray_waitables, *, num_returns1, timeoutNone, fetch_localTrue)python/ray/_private/worker.py#L3090只等待其中一部分就绪即返回两个列表就绪列表与未就绪列表。支持ObjectRef与ObjectRefGenerator混合输入num_returns需要就绪的最小数量timeout最大等待秒数二者先到先返回fetch_localTrue时等待对象下载到本节点才算就绪对 generator 则等待下一个对象下载完毕False时只要集群内任意节点可用即返回。refs [task.remote(i) for i in range(10)] ready, pending ray.wait(refs, num_returns3, timeout5.0) print(fready{len(ready)}, pending{len(pending)})4.3as_completed与map_unordered结果流式消费ray.util.as_completed(ray_futures, *, num_returns1, timeoutNone)python/ray/util/helpers.py#L57按完成顺序产出ObjectRef的迭代器适合谁先完成先处理谁的流水线场景可传入num_returns让单个任务多返回值也按完成粒度产出。ray.util.map_unordered(fn, args, *, num_returns1, ...)python/ray/util/helpers.py#L135对参数集并发执行远程函数并无序地逐个产出结果天然解决ray.get全量等待造成的阻塞与内存堆积问题。两者配合可在不丢失完成顺序信息的前提下实现大规模并行的流式处理。from ray.util import as_completed, map_unordered ray.remote def f(x): import time; time.sleep(x % 3) return x # 按完成顺序消费 for ref in as_completed([f.remote(i) for i in range(10)]): print(done:, ray.get(ref)) # 无序流式映射 for result in map_unordered(f, list(range(10))): print(got:, result)相关最佳实践可进一步阅读 ray-core/patterns/ray-get-loop.rst 与 ray-core/patterns/unnecessary-ray-get.rst。五、Runtime Context任务/作业自身的运行时信息Runtime Context 部分对应文档中的ray.runtime_context.get_runtime_context、ray.runtime_context.RuntimeContext与ray.get_gpu_ids实现位于 python/ray/runtime_context.py。5.1RuntimeContext与get_runtime_context()get_runtime_context()python/ray/runtime_context.py#L641返回当前作业/任务的RuntimeContext单例python/ray/runtime_context.py#L17提供如下常用信息job_id/worker_id/node_id当前作业、worker、节点标识namespace当前作业所属命名空间get_assigned_resources()当前 worker 被调度到的资源集合current_actor/current_task_idActor 环境下的自身句柄与任务 IDruntime_env当前作业的运行时环境配置。典型用途是在任务内部自适应资源如根据 CPU 数调整并行度、打印定位日志import ray ctx ray.runtime_context.get_runtime_context() print(ctx.job_id, ctx.node_id, ctx.namespace) ray.remote def whoami(): c ray.runtime_context.get_runtime_context() return {task: c.get_current_task_id(), resources: c.get_assigned_resources()}5.2ray.get_gpu_ids()GPU 可见性查询ray.get_gpu_ids()python/ray/_private/worker.py#L1178返回当前 worker 可见的 GPU ID 列表。只有当任务/Actor 通过num_gpus声明了 GPU 需求时Ray 才会为其分配 GPU 并通过该接口暴露因此它常与CUDA_VISIBLE_DEVICES配合使用ray.remote(num_gpus1) def train(): gpu_ids ray.get_gpu_ids() assert len(gpu_ids) 1 # 使用 gpu_ids[0] 绑定 CUDA 设备 return gpu_ids六、Cross Language跨语言调用 JavaCross Language 部分对应文档中的ray.cross_language.java_function与ray.cross_language.java_actor_class它们由 python/ray/init.py#L140 从 python/ray/cross_language.py 导入。二者的完整用法可参阅配套指南 doc/source/ray-core/cross-language.rst其基本形态是import ray from ray.cross_language import java_function, java_actor_class ray.init() # 调用 Java 端注册的静态函数 f java_function(io.ray.examples.HelloWorld, hello) result ray.get(f.remote(ray)) # 创建 Java Actor 并调用其方法 Counter java_actor_class(io.ray.examples.Counter) actor Counter.remote(0) value ray.get(actor.increment.remote())机制上Ray 将 Python 侧的调用翻译为对 Java worker 中对应类/方法的远程调用返回值经对象存储跨语言序列化回传。需要注意的是跨语言调用要求集群中同时存在 Java 与 Python worker集群启动时配置好ray start --java相关参数并遵循 java/README.md 中约定的类与方法导出规则。七、总结如何系统学习这些 API本文覆盖的接口仅是 Ray Core 公开 API 的主骨架。从文档结构看doc/source/ray-core/api/ 目录下还配套有 exceptions.rstGetTimeoutError、TaskCancelledError、ActorAlreadyExistsError等异常体系、scheduling.rst调度器接口、utility.rst工具函数与 cli.rstray start/ray stop等 CLI它们与本文的六大 API 家族共同构成完整的 Ray Core 编程面。动手实践建议按以下路径推进生命周期用ray.init()ray.is_initialized()建立稳定的程序骨架Task用ray.remoteoptions()理解资源声明与重试语义Actor用ActorClass.optionsget_actorkill掌握有状态服务与命名服务发现Object用ray.put/get/wait与as_completed/map_unordered写出不阻塞、可流式的并行代码进阶用 Runtime Context 感知环境用 Cross Language 打通 Java/Python 混合集群。每一步都能在仓库的 python/ray 源码与 doc/source/ray-core 文档中找到对应的实现与示例代码例如 doc/source/ray-core/doc_code/tasks.py、doc/source/ray-core/doc_code/actors.py 与 doc/source/ray-core/doc_code/obj_ref.py 提供了可直接运行的最小样例。【免费下载链接】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),仅供参考