Article

数据计算 Ray

更新于:2026-07-13

第一章:Ray 概述与核心概念

1.1 什么是 Ray?

概念名称说明注意事项
Ray一个开源的分布式计算框架,旨在为AI和通用计算提供简单、通用且高性能的并行与分布式执行能力。由UC Berkeley RISELab开发,支持Python和Java,特别适用于机器学习、强化学习、超参数调优等场景。Ray 不仅是一个任务调度系统,更是一个构建分布式应用的通用基础设施平台。
分布式计算框架允许程序在多台机器上并行执行,以提升性能、吞吐量和可扩展性。Ray 提供了低延迟的任务调度和高效的对象共享机制。初学者需理解”分布式”与”并行”的区别:分布式强调跨节点资源协调,而并行强调任务同时执行。
核心特性:轻量级任务模型支持每秒数百万级轻量级任务调度,延迟极低(微秒级),适合细粒度并行。与传统批处理框架不同,Ray 的任务是异步、动态生成的,适用于动态控制流场景。
核心特性:Actor 模型支持原生支持有状态的分布式Actor,便于管理持久化状态和封装行为。Actor 是有状态的,其生命周期独立于调用者,需注意资源释放问题。
核心特性:透明的对象共享使用共享内存(如Plasma Object Store)实现零拷贝数据共享,减少序列化开销。大对象存储在Object Store中,小对象可能直接通过IPC传输。

1.2 Ray 的设计目标与适用场景

概念名称说明注意事项
设计目标:通用性(Generality)能够表达任意计算模式,包括无状态任务、有状态Actor、流水线、图结构等。Ray 并非专为某类应用设计,而是作为底层运行时支撑多种上层库(如Tune、Serve)。
设计目标:简单性(Simplicity)API 简洁,开发者只需添加 @ray.remote.remote() 即可将函数或类转为分布式执行。尽管使用简单,但理解其背后机制对性能优化至关重要。
设计目标:高性能(Performance)实现低延迟任务调度(<1ms)、高吞吐(>1M tasks/s),支持大规模集群。性能受网络、序列化、对象存储等因素影响,需合理配置资源。
适用场景:机器学习训练与调优支持分布式超参搜索(Ray Tune)、模型训练(Ray Train)、强化学习(RLlib)。非常适合需要大量试验的ML工作流。
适用场景:模型服务部署通过 Ray Serve 实现高并发、弹性伸缩的模型在线推理服务。可结合Kubernetes进行生产级部署。
适用场景:数据处理流水线使用 Ray Data 构建高效的数据预处理和ETL流程。类似Spark但更灵活,尤其适合Python生态。
适用场景:科学计算与仿真支持复杂模拟、蒙特卡洛方法、并行搜索等计算密集型任务。可替代multiprocessing或joblib进行跨节点扩展。

1.3 Ray 的核心抽象:任务(Task)与 Actor

概念名称说明注意事项
任务(Task)无状态的远程函数调用单元。通过 @ray.remote 装饰普通函数后,调用其 .remote() 方法即创建一个任务。任务是不可变的,每次调用都会产生新的执行实例;返回值为 ObjectRef
Actor有状态的远程对象实例。通过 @ray.remote 装饰类后,创建其实例即生成一个Actor,可在分布式环境中保持状态。每个Actor运行在独立的进程中,拥有自己的内存空间和状态;方法调用也通过 .remote() 异步执行。
任务 vs Actor 对比任务用于无状态计算(如数据转换、模型推理),Actor用于有状态服务(如缓存、数据库连接、环境模拟器)。选择使用任务还是Actor取决于是否需要维护跨调用的状态。
远程执行(.remote())所有 Ray 分布式调用均通过 .remote() 触发,返回 ObjectRef 而非实际结果。.remote() 是异步非阻塞的,必须配合 ray.get()ray.wait() 获取结果。
ObjectRef表示远程任务或Actor方法调用结果的引用句柄,类似于Future或Promise。可在网络间传递,实现任务依赖链;不持有实际数据,只指向Object Store中的对象。

1.4 Ray 的架构组成(Driver、Worker、Object Store、GCS 等)

组件名称说明注意事项
Driver(驱动程序)用户编写的Python脚本入口点,负责启动Ray应用、定义任务/Actor并触发执行。通常运行在集群的一个节点上,是任务提交的源头。
Worker执行远程任务或Actor方法的实际进程。分为两种类型:任务Worker(执行无状态函数)和Actor Worker(运行Actor实例的方法)。Worker由Ray运行时自动管理,开发者无需手动创建或销毁。
Object Store(对象存储)基于共享内存(如Apache Arrow Plasma)的高性能数据存储系统,用于存放任务输入输出对象。支持零拷贝读取,极大提升数据共享效率;对象保留在内存中直到被垃圾回收。
Global Control Store (GCS)中央协调服务(基于Redis或内置GCS Server),存储全局元信息(如Actor位置、资源分配、作业状态)。在大规模集群中建议使用专用GCS服务器以提高可用性和性能。
Raylet(已逐步淘汰,由Core取代)旧版架构中每个节点上的本地调度器和服务代理,负责管理本地资源、Worker和对象。新版本Ray(>=2.0)已用统一的Ray Core组件替代Raylet。
Ray Core(新版核心)当前Ray的核心运行时,整合了调度、资源管理、对象管理等功能,提供统一API。支持更高效的跨节点通信和弹性扩缩容。
Dashboard Agent运行在每个节点上的监控代理,收集指标并提供Web界面访问。默认启用,可通过浏览器访问 http://<head-node>:8265 查看集群状态。

1.5 Ray 与其他分布式框架的对比

对比维度RayApache SparkDaskmultiprocessing
编程模型任务 + Actor(通用计算模型)RDD/DataFrame(批处理为主)Task Graph + Futures(类似Ray)进程/线程并发(单机)
语言支持Python, Java(主要Python)Scala, Python, Java, RPython(为主)Python
执行模式细粒度任务调度,低延迟(μs级)批处理为主,延迟较高(ms~s级)支持流式和批处理,延迟较低单机多进程,无网络通信
状态管理原生支持有状态Actor无原生状态管理,依赖外部存储支持有限状态(如Dask Actor)共享内存或队列
适用场景AI/ML、强化学习、服务部署、通用分布式应用大数据ETL、SQL分析、批处理数据科学、并行计算、替代multiprocessing单机CPU密集型任务
扩展性支持数千节点集群成熟的大规模集群支持一般支持数百节点仅限单机
易用性Python API简洁,.remote() 即分布式DSL丰富(DataFrame API),但学习曲线较陡与Pandas/Numpy兼容性好标准库,易上手但难跨主机
内置AI工具Ray Tune, Serve, Train, RLlib 等完整生态MLlib(较基础),集成第三方工具Dask-ML, cuML等
数据共享机制共享内存(Object Store),零拷贝序列化+网络传输或磁盘序列化+TCP/IP或IPCPipe, Queue, Shared Memory
注意事项更适合动态、控制流复杂的AI应用更适合结构化数据批处理更适合数据科学工作流无法跨机器扩展

第二章:Ray 入门与环境搭建

2.1 安装 Ray(单机模式)

操作名称操作细节注意事项
使用 pip 安装 Ray执行命令:pip install ray 或安装带额外依赖的版本:pip install "ray[default]"推荐使用虚拟环境(如 venv 或 conda)避免依赖冲突;确保 Python 版本为 3.7+。
安装特定模块(可选)可选择性安装:ray[tune](超参数调优)、ray[serve](模型服务)、ray[train](分布式训练)、ray[dashboard](启用可视化仪表盘)若后续使用相关功能,建议提前安装对应模块,避免运行时报错。
验证安装成功在 Python 中执行:import ray; print(ray.version)若无报错并能输出版本号,则表示安装成功。
系统依赖(Linux/macOS)通常无需额外配置;Windows 用户需注意:支持 Windows 10/11,某些功能(如 Object Store)性能可能略低于 Unix 系统。Windows 上建议使用 WSL2 获得更佳体验。
安装问题排查常见问题:权限不足(使用 --user 参数)、网络超时(更换 pip 源,如清华、阿里云镜像)、依赖冲突(使用虚拟环境隔离)推荐使用 pip install -U pip 升级 pip 后再安装 Ray。

2.2 启动与关闭 Ray 集群(本地 & 集群)

操作名称操作细节注意事项
启动本地单节点集群在 Python 脚本中调用:ray.init()默认自动启动所有必要组件(GCS、Object Store、Dashboard 等);适用于开发和测试。
指定资源启动(可选)ray.init(num_cpus=4, num_gpus=1)可手动限制资源使用,用于模拟资源受限环境或调试。
启动无 Dashboard 模式ray.init(include_dashboard=False)减少内存占用,适合生产环境或资源紧张场景。
以配置文件启动集群(多节点)编写 cluster.yaml 配置文件,并使用 CLI:ray up cluster.yaml用于云环境或 Kubernetes 部署;需提前安装 ray[cluster]
连接已有集群ray.init(address="auto") 或指定地址:ray.init(address="ray://<head-node>:10001")必须确保网络可达且端口开放;常用于多客户端连接同一集群。
关闭本地集群在代码中调用:ray.shutdown()显式释放所有资源,防止内存泄漏;每个 ray.init() 应对应一个 ray.shutdown()
命令行关闭集群使用命令:ray stop终止当前节点上所有 Ray 进程,包括后台守护进程。
注意事项:ray.init() 多次调用ray.init() 在同一进程中只能成功调用一次;重复调用会抛出异常。若需重置状态,先调用 ray.shutdown() 再重新 ray.init()

2.3 验证安装与基本连接测试

操作名称操作细节注意事项
编写测试脚本 test_ray.py示例代码:这是最小可运行示例,验证基本功能是否正常。
import ray
ray.init()

@ray.remote
def hello():
    return "Hello from Ray!"

result = ray.get(hello.remote())
print(result)
ray.shutdown()
操作名称操作细节注意事项
运行测试脚本执行命令:python test_ray.py应输出 “Hello from Ray!”,表示任务成功执行。
检查进程是否启动使用系统命令查看 Ray 相关进程:ps aux | grep ray应能看到多个 Ray 进程(如 redis-server、raylet、dashboard agent 等)。
检查端口占用默认端口包括:6379(Redis/GCS)、8265(Dashboard)、10001(Ray Client Server)若端口被占用,可在 ray.init() 中通过参数指定其他端口。
测试多任务并行示例代码:
import ray
ray.init()

@ray.remote
def task(n):
    return n * n

tasks = [task.remote(i) for i in range(4)]
print(ray.get(tasks))
ray.shutdown()

应输出 [0, 1, 4, 9],验证并行执行能力。 | | 注意事项:首次运行延迟 | 首次启动可能耗时较长(5~10秒),因需初始化多个后台服务。 | 后续运行将显著加快。 |

2.4 Ray Dashboard 的启用与访问

操作名称操作细节注意事项
默认启用 Dashboard调用 ray.init() 时自动启用(除非设置 include_dashboard=FalseDashboard 提供集群状态、任务监控、Actor 信息等可视化视图。
查看 Dashboard 地址启动 Ray 后,日志中会输出类似:Dashboard URL: http://127.0.0.1:8265记录该地址,用于浏览器访问。
浏览器访问 Dashboard打开浏览器,访问:http://localhost:8265若在远程服务器运行,需通过 SSH 端口转发或公网 IP 访问。
SSH 端口转发访问远程 Dashboard本地执行:ssh -L 8265:localhost:8265 user@server_ip将远程服务器的 8265 端口映射到本地,实现安全访问。
Dashboard 主要功能页包括:Cluster(节点资源使用情况)、Jobs(作业运行状态)、Actors(Actor 实例列表与状态)、Logs(各组件日志查看)、Metrics(性能指标图表)是调试和性能分析的重要工具。
自定义 Dashboard 端口启动时指定:ray.init(dashboard_port=8266)避免与本地其他服务端口冲突。
Dashboard 无法访问排查常见原因:防火墙阻止端口、include_dashboard=False、端口被占用、远程服务器未绑定正确 IP检查启动日志确认 Dashboard 是否成功启动。
生产环境建议可关闭 Dashboard 或限制访问权限,防止信息泄露。也可通过反向代理(如 Nginx)增加认证层。

第三章:任务并行化(Remote Functions)

3.1 定义远程函数(@ray.remote)

方法/语法用途代码示例注意事项
@ray.remote将普通 Python 函数转换为可在 Ray 集群中分布式执行的远程任务。@ray.remote
def my_function(x):
return x ** 2
必须先调用 ray.init() 才能使用 @ray.remote;装饰器不接受参数时可直接使用 @ray.remote
@ray.remote(num_returns=n)指定该函数返回多个对象引用(ObjectRef),Ray 会将其拆分为 n 个独立的返回值。@ray.remote(num_returns=2)
def divide_data(data):
mid = len(data)//2
return data[:mid], data[mid:]
left_ref, right_ref = divide_data.remote(lst)
当函数逻辑上产生多个独立输出时使用,避免打包成元组后解包开销。
@ray.remote(max_calls=k)设置函数执行 k 次后自动重启 Worker 进程,防止内存泄漏累积。@ray.remote(max_calls=1000)
def leaky_func():
# 可能缓慢增长内存
...
适用于可能因闭包或缓存导致内存缓慢增长的函数;设为 1 表示每次调用后重启。
@ray.remote(resources={'custom': 1})为任务指定自定义资源需求(如 GPU、TPU 或用户定义资源)。@ray.remote(resources={"gpu": 1})
def train_model():
...
资源名称和数量需与集群配置一致;可用于实现异构资源调度。
@ray.remote(num_cpus=n, num_gpus=m)显式声明任务所需的 CPU 和 GPU 数量。@ray.remote(num_cpus=2, num_gpus=1)
def gpu_task():
import torch
device = "cuda" if torch.cuda.is_available() else "cpu"
...
默认值为 num_cpus=1, num_gpus=0;若无 GPU 则 num_gpus 应省略或设为 0。

3.2 调用远程函数(.remote())

方法/语法用途代码示例注意事项
func.remote(*args, **kwargs)异步提交远程任务并立即返回 ObjectRef,不阻塞主程序执行。@ray.remote
def add(a, b):
return a + b
result_ref = add.remote(2, 3)
返回的是 ObjectRef,不是实际结果;必须通过 ray.get() 获取结果。
参数传递:基本类型支持 int、float、str、bool 等内置类型自动序列化传输。@ray.remote
def greet(name):
return f"Hello {name}"
ref = greet.remote("Alice")
序列化开销小,推荐用于轻量级数据。
参数传递:复杂对象支持 list、dict、numpy array、pandas DataFrame 等。import numpy as np
@ray.remote
def process_array(arr):
return arr.sum()
arr = np.random.rand(1000)
ref = process_array.remote(arr)
大对象会存入 Object Store,仅传递引用,提升效率。
参数传递:ObjectRef可将其他任务的结果引用作为输入,形成任务依赖图。ref1 = task_a.remote()
ref2 = task_b.remote(ref1) # 依赖 ref1
Ray 自动解析依赖关系,并在依赖完成后再调度后续任务。
并行调用多个任务可连续调用 .remote() 创建多个并行任务。refs = [worker.remote(i) for i in range(10)]所有任务几乎同时被调度,实现数据并行处理。

3.3 获取任务结果(ray.get)

方法/语法用途代码示例注意事项
ray.get(ObjectRef)同步阻塞等待任务完成,并获取其实际返回值。result = ray.get(result_ref)若任务未完成,当前线程将阻塞直到结果可用;适合等待关键结果。
ray.get(list of ObjectRefs)批量获取多个任务的结果,按顺序返回列表。results = ray.get([f1.remote(), f2.remote()])阻塞直到所有任务完成;返回顺序与输入 ObjectRef 顺序一致。
超时控制:ray.get(timeout=…)设置最大等待时间(秒),超时抛出 GetTimeoutErrortry:
result = ray.get(f.remote(), timeout=5)
except ray.exceptions.GetTimeoutError:
print("Task timed out")
timeout=None 表示无限等待;有助于避免程序永久挂起。
异常传播若远程任务内部抛出异常,ray.get() 会重新抛出相同异常。@ray.remote
def may_fail():
raise ValueError("oops")
try:
ray.get(may_fail.remote())
except ValueError as e:
print(e)
异常信息包含完整堆栈跟踪,便于调试远程错误。
性能建议避免频繁调用 ray.get() 在循环中等待单个结果。错误方式:for ref in refs: print(ray.get(ref))
正确方式:print(ray.get(refs))
批量获取更高效,减少调度开销。

3.4 等待任务部分完成(ray.wait)

方法/语法用途代码示例注意事项
ray.wait(object_refs, num_returns=1, timeout=None)等待至少 num_returns 个任务完成,返回已完成和未完成的 ObjectRef 元组。ready, remaining = ray.wait(refs, num_returns=2)不必等待全部任务完成,适合”谁先到谁先处理”场景。
num_returns=n指定希望等待完成的任务数量,默认为 1。ready, _ = ray.wait(refs, num_returns=5)可用于实现批量处理或流水线消费。
timeout=seconds设置最大等待时间,超时则返回已就绪的部分结果。ready, remaining = ray.wait(refs, timeout=1.0)timeout=None 表示无限等待;浮点数支持毫秒级精度。
非阻塞检查:timeout=0实现轮询机制,检查当前有哪些任务已完成。ready, remaining = ray.wait(refs, timeout=0)可用于实现自定义调度逻辑或状态监控。
处理就绪结果ready 中的引用调用 ray.get() 获取实际值。for r in ready: print(ray.get(r))ray.wait 只返回引用,仍需 ray.get 获取数据。
使用场景:流式处理一边生成任务,一边消费已完成结果,降低内存压力。while unstarted_tasks:
new_ref = launch_next_task()
refs.append(new_ref)
ready, refs = ray.wait(refs, min_returns=1, timeout=0)
for r in ready:
process(ray.get(r))
实现类似生成器的流式计算模式。

3.5 任务依赖与数据传递机制

概念名称说明注意事项
数据传递:值传递(Value Passing)所有参数在调用 .remote() 时被序列化并传入任务。基本类型和小对象直接复制;大对象通过共享内存(Object Store)传递引用。
数据传递:引用传递(ObjectRef 传递)将一个任务的 ObjectRef 作为参数传给另一个任务,形成显式依赖。Ray 自动确保依赖任务完成后才调度下游任务;是构建 DAG 的基础。
零拷贝共享(Zero-copy Sharing)使用 Apache Arrow 格式存储对象,不同进程可通过共享内存直接访问,无需复制。特别适用于 numpy arrays、pandas DataFrames 等;大幅提升大数据处理效率。
序列化机制默认使用 Pickle,对 Arrow 支持类型使用高效序列化。用户可注册自定义序列化器以优化特定类型。
任务依赖图(DAG)通过 ObjectRef 传递隐式构建任务间的依赖关系图。Ray 运行时根据 DAG 自动调度,无需手动管理依赖。
共享变量问题多个任务不能直接共享可变全局变量(因隔离执行)。应使用 Actor 或 Ray 内置的 ray.put() 存储共享只读数据。
ray.put(value)将一个值显式放入 Object Store,返回其 ObjectRef,供多个任务共享。data_id = ray.put(large_dataset)
task1.remote(data_id)
task2.remote(data_id)
ray.get(object_ref) 触发依赖ray.get 是同步操作,会阻塞直到对应任务完成。在任务内部使用 ray.get 会增加该任务的执行时间,并可能造成级联阻塞。

第四章:状态管理与 Actor 模型

4.1 定义远程 Actor(@ray.remote class)

方法/语法用途代码示例注意事项
@ray.remote 装饰类将普通 Python 类转换为可在 Ray 集群中分布式部署的远程 Actor。见下方示例类中的 __init__ 方法和所有实例方法均可远程调用;必须先调用 ray.init()
@ray.remote
class Counter:
    def __init__(self):
        self.value = 0

    def increment(self):
        self.value += 1
        return self.value
方法/语法用途代码示例注意事项
@ray.remote(class_args=[...])在创建 Actor 实例时传递固定参数给 __init__ 方法(较少使用)。@ray.remote(class_args=["Alice"])
class Greeter:
def __init__(self, name):
self.name = name

g = Greeter.remote() # 自动传入 "Alice"
一般推荐在 .remote() 调用时传参,更灵活;class_args 用于配置固定的类级参数。
@ray.remote(num_cpus=n, num_gpus=m)指定 Actor 实例运行所需的 CPU 和 GPU 资源。@ray.remote(num_cpus=2, num_gpus=1)
class ModelServer:
def __init__(self, model_path):
self.model = load_model(model_path)
资源在 Actor 生命周期内独占;调度器会为其分配满足条件的节点。
@ray.remote(resources={'custom': 1})指定 Actor 对自定义资源的需求(如 “ssd”, “tpu” 等)。@ray.remote(resources={"gpu": 1, "fast_network": 1})
class HighSpeedActor:
...
自定义资源需在集群启动时声明;可用于实现高级调度策略。
@ray.remote(max_concurrency=n)设置 Actor 方法可并发执行的最大数量(启用异步 Actor 模式)。@ray.remote(max_concurrency=100)
class AsyncWorker:
async def process(self, data):
await some_async_io()
return result
必须配合 async/await 使用;默认 max_concurrency=1(即方法串行执行)。
@ray.remote(lifetime='detached')创建脱离 Driver 生命周期的持久化 Actor,即使 Driver 退出仍存活。@ray.remote(lifetime="detached", name="global_counter")
class Counter:
def __init__(self):
self.val = 0

actor = Counter.remote()
需通过 ray.get_actor("global_counter") 获取引用;必须手动调用 ray.kill() 销毁。

4.2 创建 Actor 实例(.remote())

方法/语法用途代码示例注意事项
ActorClass.remote(*args, **kwargs)异步创建远程 Actor 实例,返回 Actor 句柄(ActorHandle),不阻塞主程序。counter = Counter.remote()
worker = Worker.remote(init_data)
返回的是 Actor 句柄,可用于后续方法调用;__init__ 在后台异步执行。
传递构造参数将参数传递给被装饰类的 __init__ 方法。@ray.remote
class Adder:
def __init__(self, base):
self.base = base

a = Adder.remote(10)
参数传递机制与远程函数相同,支持基本类型、复杂对象和 ObjectRef。
检查 Actor 是否创建成功使用 ray.get(actor.__init__.remote()) 等待构造完成。try:
ray.get(counter.__init__.remote(), timeout=5)
except ray.exceptions.GetTimeoutError:
print("Actor init timeout")
__init__ 方法也可能失败(如抛出异常),可通过 ray.get 捕获异常。
创建多个 Actor 实例可并行创建多个独立的 Actor 实例。actors = [Counter.remote() for _ in range(5)]每个 Actor 拥有独立的状态和资源;适用于数据并行或服务分片。
注意事项:资源分配Actor 创建时会立即申请声明的资源(CPU/GPU等),直到 ray.kill() 或节点失效才释放。避免创建过多 Actor 导致资源耗尽;可通过 ray list actors 查看当前 Actor 状态。

4.3 调用 Actor 方法(.remote())

方法/语法用途代码示例注意事项
actor_handle.method.remote(*args, **kwargs)调用远程 Actor 的方法,返回 ObjectRef,实现异步非阻塞调用。result_ref = counter.increment.remote()
result = ray.get(result_ref)
所有方法调用必须使用 .remote(),即使方法本身是同步的。
调用 __call__ 方法(可选)若类定义了 __call__,也可通过 actor_handle.remote() 调用。@ray.remote
class CallableActor:
def __call__(self, x):
return x * 2

a = CallableActor.remote()
ref = a.remote(5)
语法糖,不常用。
并行调用多个 Actor 方法可对同一或不同 Actor 发起多个 .remote() 调用。refs = [a.increment.remote() for a in actors]所有调用并发提交,但每个 Actor 默认串行处理方法(除非启用 max_concurrency)。
传递 ObjectRef 作为参数可将任务或其他 Actor 的结果引用传入 Actor 方法。data_ref = ray.put(large_data)
actor.process.remote(data_ref)
Ray 自动解析依赖,确保 data_ref 就绪后再执行 process 方法。
异常处理若 Actor 方法抛出异常,ray.get() 会重新抛出。try:
ray.get(actor.fail.remote())
except ValueError as e:
print("Caught remote exception")
异常信息包含完整堆栈,便于调试。
注意事项:方法执行顺序默认情况下,同一 Actor 的方法按 .remote() 提交顺序串行执行。使用 max_concurrency > 1 可启用并发执行,但需自行处理状态同步问题。

4.4 Actor 的状态隔离与生命周期

概念名称说明注意事项
状态隔离每个 Actor 实例在独立的进程中运行,拥有私有的内存空间和状态,不受其他 Actor 或任务影响。不同 Actor 之间不能直接访问对方的内部变量;必须通过方法调用通信。
状态持久性Actor 的状态在其生命周期内持续存在,跨方法调用保持不变。适用于维护数据库连接、模型缓存、环境状态等有状态服务。
生命周期开始调用 ActorClass.remote() 时开始,__init__ 方法执行标志初始化完成。__init__ 失败,Actor 创建失败,句柄不可用。
生命周期结束以下情况结束:显式调用 ray.kill(actor_handle)、Driver 调用 ray.shutdown() 且非 detached Actor、节点崩溃或资源回收。非 detached Actor 与创建它的 Driver 强关联;Driver 退出则 Actor 被销毁。
detached Actor通过 lifetime='detached' 创建,独立于 Driver 生命周期,需手动销毁。适用于全局共享服务(如配置中心、计数器);必须通过名称注册和获取。
资源占用Actor 创建后即占用声明的资源(CPU/GPU),直到销毁才释放。应合理设计 Actor 数量,避免资源浪费或调度失败。
故障恢复默认不提供自动恢复;若节点失败,Actor 及其状态丢失。可结合检查点(Checkpoint)和外部存储实现高可用;Ray 本身不自动重启失败的 Actor。

4.5 Actor 与任务的交互模式

交互模式说明代码示例注意事项
任务调用 Actor 方法远程任务中调用 Actor 方法,实现任务与有状态服务的协作。@ray.remote
def worker_task(counter):
for _ in range(10):
val = ray.get(counter.increment.remote())
return val

c = Counter.remote()
ref = worker_task.remote(c)
Actor 句柄可作为参数传递给任务;任务依赖 Actor 的状态。
Actor 调用远程任务Actor 方法内部提交远程任务,实现异步计算委托。@ray.remote
class Orchestrator:
def run_pipeline(self):
refs = [process_data.remote(d) for d in self.data]
return ray.get(refs)

o = Orchestrator.remote()
result_ref = o.run_pipeline.remote()
Actor 可作为任务编排者(Orchestrator),管理子任务生命周期。
Actor 之间通信一个 Actor 调用另一个 Actor 的方法,实现 Actor 间协作。a1 = ActorA.remote()
a2 = ActorB.remote()
a1.process_and_send.remote(a2)
构建分布式 Actor 系统(如 Actor 模型、微服务)。
共享 Actor 实例多个任务或 Actor 共享同一个 Actor 句柄,访问其共享状态。global_counter = Counter.remote()
tasks = [worker.remote(global_counter) for _ in range(10)]
实现全局计数器、锁、缓存等共享服务。
传递 ObjectRef 交互任务和 Actor 之间通过 ObjectRef 传递数据,解耦生产者与消费者。data_ref = ray.put(large_data)
actor.consume.remote(data_ref)
task.process.remote(data_ref)
结合 Object Store 实现高效数据共享,避免重复传输。
注意事项:死锁风险若多个 Actor 循环等待彼此方法完成,可能造成死锁。避免在 Actor 方法中同步调用 ray.get() 等待其他 Actor 的方法;使用异步模式或消息队列解耦。A 调用 B.method 并等待,B 同时调用 A.method 并等待。

第五章:对象存储与数据共享

5.1 Ray 的对象模型与 Object Store

概念名称说明注意事项
对象模型(Object Model)Ray 中所有任务输入、输出和 Actor 状态均以”对象”形式存在,通过唯一的 ObjectRef 引用。对象是不可变的,一旦创建内容不可更改;更新需生成新对象。
Object Store(对象存储)基于共享内存的高性能存储系统(底层使用 Apache Arrow Plasma 或新核心存储),用于存放分布式对象。每个节点有自己的 Object Store;对象在本地节点以零拷贝方式共享。
内存映射(Memory Mapping)Object Store 使用共享内存段,多个进程可直接访问同一数据块,避免复制开销。适用于 numpy arrays、pandas DataFrames 等 Arrow 支持的数据结构。
对象位置透明性开发者无需关心对象物理位置;Ray 运行时自动管理对象的本地/远程访问。若对象在远程节点,会自动通过网络传输(序列化)到本地 Object Store。
对象分发机制当任务依赖远程对象时,Ray 自动触发对象传输(pull 模式),确保执行节点有数据。调度器倾向于将任务调度到对象所在节点(数据本地性优化)。
节点本地存储每个节点的 Object Store 仅存储该节点任务频繁访问的对象,非全局共享。集群中不同节点可能存储同一对象的多个副本。
扩展性Object Store 内存受限于节点物理内存;超限可能导致 OOM 或性能下降。建议监控内存使用,合理设计数据分片和生命周期。

5.2 对象引用(ObjectRef)详解

方法/语法用途代码示例注意事项
ray.ObjectRef表示远程对象的引用句柄,是任务或 Actor 方法调用的返回值类型。ref = my_func.remote()ObjectRef 不包含实际数据,只指向 Object Store 中的对象。
ray.get(ObjectRef)同步获取 ObjectRef 指向的实际对象值。result = ray.get(ref)若对象未就绪,阻塞等待;若对象在远程,自动拉取并反序列化。
ray.wait(ObjectRef, ...)等待一个或多个 ObjectRef 就绪,返回已完成和未完成的引用列表。ready, _ = ray.wait([ref1, ref2], num_returns=1)非阻塞或超时等待,适合流式处理和依赖管理。
传递 ObjectRef 作为参数将 ObjectRef 传给其他任务或 Actor,形成任务依赖。@ray.remote
def consumer(data_ref):
data = ray.get(data_ref)
...

consumer.remote(result_ref)
Ray 自动解析依赖,确保 data_ref 就绪后再调度 consumer。
ObjectRef 的序列化ObjectRef 可被序列化并传递到其他任务或 Actor。@ray.remote
def forward_ref(ref):
return ref

new_ref = ray.get(forward_ref.remote(old_ref))
实现引用转发、任务编排等高级模式。
ObjectRef 的唯一性每个任务或方法调用返回唯一的 ObjectRef,即使输入相同。ref1 = func.remote(1)
ref2 = func.remote(1)
assert ref1 != ref2
ObjectRef 是句柄,不是对象内容的哈希。
注意事项:悬挂引用(Dangling Ref)若对象已被垃圾回收,ray.get() 可能失败。避免长期持有不再需要的 ObjectRef;及时释放资源。

5.3 零拷贝共享与序列化机制

方法/语法用途代码示例注意事项
零拷贝共享(Zero-copy Sharing)对支持 Arrow 格式的对象(如 numpy array),在进程间通过共享内存直接访问,无需复制。import numpy as np
arr = np.ones((1000, 1000))

@ray.remote
def process(arr):
return arr.sum()

ref = process.remote(arr)
大幅减少内存占用和传输延迟;仅适用于 Arrow 支持的不可变类型。
默认序列化(Pickle)对不支持 Arrow 的对象使用 Python pickle 进行序列化和反序列化。class CustomObj:
def __init__(self, x):
self.x = x

@ray.remote
def use_custom(obj):
return obj.x
Pickle 通用但较慢,且可能不安全;大对象序列化开销大。
Arrow 支持类型包括:numpy.ndarray(特定 dtype)、pandas.DataFramepyarrow.Tablebytesstrlistdict(小对象)。df = pd.DataFrame({"a": [1,2,3]})
ray.put(df) # 零拷贝
推荐使用 Arrow 兼容格式以获得最佳性能。
自定义序列化器用户可注册特定类型的高效序列化函数。import pyarrow as pa
ray.util.register_serializer(
MyClass,
serializer=lambda obj: obj.to_bytes(),
deserializer=MyClass.from_bytes
)
用于优化自定义类的传输性能。
序列化性能对比零拷贝:μs 级延迟,无内存复制;Pickle:ms 级延迟,需内存复制。优先使用 numpy/pandas 等库处理数据;大对象应尽量避免频繁序列化。
注意事项:可变对象风险若多个任务共享同一可变对象(如 list),修改可能影响其他任务。应假设对象是不可变的;修改前应深拷贝;Ray 不强制不可变性,需开发者自行保证。

5.4 大对象处理与分片策略

操作名称操作细节注意事项
使用 ray.put() 显式存放大对象将大型只读数据(如模型、数据集)放入 Object Store,返回 ObjectRef。model_ref = ray.put(large_model)
tasks = [infer.remote(model_ref, data) for data in dataset]
数据分片(Data Sharding)将大数据集切分为多个块,每个任务处理一个分片。chunks = [data[i:i+size] for i in range(0, len(data), size)]
refs = [process_chunk.remote(chunk) for chunk in chunks]
广播大对象通过 ray.put() + 多个任务引用,实现大对象的高效广播。ref = ray.put(huge_lookup_table)
workers = [Worker.remote(ref) for _ in range(10)]
流式处理大对象结合 ray.wait 实现边生成边消费,避免内存堆积。while data_stream.has_next():
ref = process.remote(data_stream.next())
ready, _ = ray.wait([ref], timeout=0)
for r in ready:
save_result(ray.get(r))
大对象传输优化启用压缩或自定义序列化减少网络传输量。ray.init(_plasma_directory="/dev/shm") # 使用内存盘加速
注意事项:对象大小限制单个对象默认最大约 2GB(受 Plasma 限制);新核心存储支持更大对象。超大对象应主动分片;避免单任务处理超大输入。
监控对象内存使用 ray memory 命令查看当前对象存储使用情况。ray memory --stats

5.5 对象生命周期与垃圾回收

概念名称说明注意事项
对象创建调用 ray.put(value)task.remote() 时创建对象并存入本地 Object Store。对象在提交任务时即创建,即使任务未开始执行。
对象就绪任务完成或 ray.put 完成后,对象状态变为”就绪”,可被 ray.get 获取。ray.wait 可检测对象是否就绪。
引用计数(Reference Counting)Ray 使用分布式引用计数跟踪每个对象的 ObjectRef 数量。当所有引用被删除且无任务依赖时,对象可被回收。
本地引用删除删除变量(如 del ref)或超出作用域会减少引用计数。Python 的引用计数是本地的;Ray 运行时维护全局引用状态。
分布式垃圾回收当所有节点上的引用计数归零,Ray 自动从 Object Store 中删除对象。回收是异步的,可能存在短暂延迟。
Detached Actor 关联对象由 detached Actor 创建的对象,即使 Driver 退出,只要 Actor 存活,对象仍有效。需手动销毁 Actor 以释放其创建的所有对象。
防止过早回收若需长期持有对象,应保持至少一个 ObjectRef 不被删除。可将重要 ObjectRef 存入全局变量或 Actor 状态中。
手动释放:ray.delete()显式删除一个或多个 ObjectRef 及其关联对象(如果无其他引用)。ray.delete([ref1, ref2])
注意事项:循环引用若 ObjectRef 形成跨任务/Actor 的循环引用,可能延迟回收。尽量避免复杂引用图;使用 ray.delete 主动清理。

第六章:资源管理与调度

6.1 资源标注(CPU、GPU、自定义资源)

概念名称说明注意事项
CPU 资源标注每个节点启动时自动检测可用 CPU 核心数,并作为默认资源 CPU。可通过 ray start --num-cpus=N 手动指定;浮点数支持(如 0.5 表示半核)。
GPU 资源标注自动检测 NVIDIA GPU 数量,注册为 GPU 资源;需安装 CUDA 和 nvidia-driver。通过 ray start --num-gpus=M 覆盖自动检测;每个 GPU 默认为 1 个资源单位。
自定义资源(Custom Resources)用户可定义任意名称的资源(如 “ssd”, “tpu”, “accelerator”),用于高级调度。启动节点时使用:ray start --resources='{"ssd":1,"fast_network":1}' 或在 ray.init(resources={...}) 中声明。
资源单位所有资源以”资源槽”为单位,任务/Actor 声明需求后由调度器分配。资源是排他的,一旦分配在任务/Actor 存活期间不被抢占。
资源可见性所有资源信息注册到 GCS,全局可见,调度器据此做决策。可通过 ray.cluster_resources() 查看集群总资源。
动态资源更新(实验性)某些场景下可运行时修改节点资源(如 Kubernetes 弹性伸缩)。需底层平台支持;不推荐频繁变更。
注意事项:资源命名冲突避免使用保留资源名(如 CPU, GPU, memory)作为自定义资源。使用前缀或命名空间(如 "myorg/ssd":1)避免冲突。

6.2 任务与 Actor 的资源需求配置

方法/语法用途代码示例注意事项
@ray.remote(num_cpus=n)指定远程函数或 Actor 所需的 CPU 资源数量。@ray.remote(num_cpus=2)
def cpu_intensive_task():
...
默认 num_cpus=1;设为 0 表示不限制(但仍需最小调度开销)。
@ray.remote(num_gpus=m)指定所需 GPU 数量,常用于深度学习任务。@ray.remote(num_gpus=1)
def train_model():
import torch
device = torch.device("cuda")
...
仅当节点有 GPU 且资源充足时才会被调度;num_gpus 必须为整数。
@ray.remote(resources={'key': value})请求自定义资源,实现异构硬件调度。@ray.remote(resources={"ssd": 1})
def fast_storage_task():
...
支持小数资源,实现资源切片;多个任务可共享一个物理资源(如多任务共用一块 GPU)。
.remote() 中覆盖资源(不支持)无法在 .remote() 调用时动态修改资源需求。❌ 错误:func.remote(num_cpus=2)资源需求必须在 @ray.remote 装饰器中静态声明。
多资源组合配置可同时声明多种资源需求。@ray.remote(num_cpus=4, num_gpus=1, resources={"fast_network": 1})调度器需找到同时满足所有资源条件的节点。
注意事项:资源超分配风险若总需求超过物理资源,任务将排队等待。监控 ray queue 查看等待任务;合理设置资源避免死锁。

6.3 动态资源分配与弹性伸缩

操作名称操作细节注意事项
自动扩缩容(Auto-scaling)Ray 集群可配置为根据负载自动添加或移除工作节点。使用 cluster.yaml 配置 min_workersmax_workers
资源需求驱动扩缩容当任务/Actor 因资源不足排队时,Ray Autoscaler 触发扩容。扩容决策基于未满足的资源请求;缩容基于节点空闲时间。
使用 Ray Client 动态连接客户端可随时连接到运行中的集群提交任务,无需重启。ray.init(address="ray://head-node:10001")
K8s 上的弹性伸缩使用 Ray Operator for Kubernetes 实现 Pod 级别扩缩容。配置 RayCluster CRD,设置 head 和 worker 组的 replicas 范围。
手动扩缩容使用 CLI 命令增减节点。ray up -n 5 cluster.yaml(扩容到 5 个 worker)
ray down cluster.yaml(销毁集群)
监控资源使用使用 ray memoryray status 和 Dashboard 查看资源消耗。ray memory --stats
ray status
注意事项:扩缩容延迟云节点启动通常需 1~5 分钟,不适合毫秒级响应场景。对延迟敏感的应用应保持最小数量的常驻节点。

6.4 调度策略与亲和性控制

概念名称说明注意事项
数据本地性调度(Data Locality)调度器优先将任务调度到其输入数据所在的节点,减少网络传输。自动启用;对大对象依赖的任务效果显著。
FIFO 调度(默认)任务按提交顺序排队,先提交的先调度(在资源满足的前提下)。公平但可能不利于高优先级任务。
优先级调度(实验性)可为任务设置优先级,高优先级任务优先获取资源。需启用实验性功能;API 可能变更。
节点亲和性(Node Affinity)通过自定义资源实现”粘性”调度,将任务固定到特定类型节点。示例:标记 SSD 节点 --resources='{"ssd":1}',任务请求 resources={"ssd":1},确保任务只在 SSD 节点运行。
反亲和性(Anti-Affinity)通过资源分配间接实现,如确保两个 Actor 不共享同一 GPU。使用 num_gpus=1 并只有 1 个 GPU 时,两个 Actor 无法共存于同一节点。
Actor 位置固定一旦 Actor 创建,其位置固定,所有方法调用均发送到该节点。适用于状态服务;避免频繁迁移状态。
调度器性能Ray 使用分布式调度器,支持高吞吐(>1M tasks/s)。避免创建过多小任务(微任务),否则调度开销占比过高。
注意事项:资源碎片化小数资源或不规则资源请求可能导致碎片,降低利用率。合理设计资源单位(如统一为 0.1 单位),避免过度切片。

第七章:容错与持久化

7.1 任务重试与失败处理

方法/语法用途代码示例注意事项
@ray.remote(max_retries=n)指定远程函数在失败后自动重试的最大次数。@ray.remote(max_retries=3)
def flaky_task():
if random.random() < 0.5:
raise Exception("Transient error")
return "success"
max_retries=0 表示不重试(默认);-1 表示无限重试(谨慎使用)。
重试条件仅对系统级错误(如节点崩溃、网络中断)或用户异常进行重试;异常包括:RayTaskErrorRayActorError 等。不重试资源不足(如 OOMKilled)或代码逻辑错误(如 ValueError)。
重试间隔Ray 内部自动进行指数退避重试,无需手动控制。系统自动处理;开发者无需实现重试逻辑,由运行时保障。
异常捕获与处理使用 ray.get() 捕获任务异常,实现自定义恢复逻辑。try:
result = ray.get(task.remote())
except ray.exceptions.RayTaskError as e:
print(f"Task failed: {e}")
适用于非幂等任务或需特殊处理的场景。
幂等性要求建议远程函数设计为幂等(Idempotent),确保重试不会产生副作用。示例:避免在函数内写入全局文件或修改外部数据库状态。若非幂等,重试可能导致数据重复或不一致。
注意事项:Actor 创建失败ActorClass.remote() 失败通常不自动重试,需手动处理。可包装在循环中重试创建;Actor 初始化失败可能因资源不足或构造函数异常。

7.2 Actor 故障恢复机制

概念名称说明注意事项
默认无自动恢复Ray 不提供 Actor 的自动故障恢复;若节点崩溃,Actor 及其状态永久丢失。适用于瞬时服务或可重建状态的场景;生产环境需自行实现恢复机制。
状态丢失风险Actor 是有状态的,其内存中的变量在进程终止后消失。不适合存储关键业务状态,除非结合持久化。
手动重建 Actor应用层可监控 Actor 状态,失败后重新创建并恢复状态。使用 ray.kill() + Actor.remote() 重建。
使用 ray.get_actor() 获取 detached Actor对于 lifetime="detached" 的 Actor,可通过名称全局访问。try:
actor = ray.get_actor("my_service")
except ValueError:
actor = MyActor.options(name="my_service", lifetime="detached").remote()
健康检查可定期调用 Actor 方法验证其是否存活。try:
ray.get(actor.health_check.remote(), timeout=5)
except (ray.exceptions.RayTaskError, ray.exceptions.GetTimeoutError):
# 触发重建逻辑
注意事项:Actor 依赖链若一个 Actor 依赖另一个,后者失败会导致前者调用失败。建议实现断路器(Circuit Breaker)模式或降级逻辑。

7.3 检查点(Checkpointing)与状态持久化

方法/语法用途代码示例注意事项
手动保存检查点在 Actor 方法中定期将关键状态写入外部存储。见下方示例推荐使用 JSON、Pickle、Joblib 或数据库(SQLite/Redis)存储。
@ray.remote
class StatefulActor:
    def __init__(self, checkpoint_path):
        self.path = checkpoint_path
        self.load_checkpoint()

    def load_checkpoint(self):
        if os.path.exists(self.path):
            with open(self.path, "rb") as f:
                self.state = pickle.load(f)

    def save_checkpoint(self):
        with open(self.path, "wb") as f:
            pickle.dump(self.state, f)
方法/语法用途代码示例注意事项
自动定期检查点结合 threading.Timer 或异步循环实现周期性保存。async def periodic_save(self):
while True:
await asyncio.sleep(60)
self.save_checkpoint()
max_concurrency > 1 的异步 Actor 中使用 asyncio。
检查点触发时机建议在以下时机保存:关键状态变更后、Actor 正常关闭前、定期时间间隔。可通过 @ray.remote(on_kill=lambda: actor.save_checkpoint()) 实现关闭钩子(实验性)。on_kill 不保证调用(如节点硬崩溃)。
使用 Ray 内部机制(有限)Ray 不提供内置的自动检查点功能,需用户自行实现。所有持久化逻辑必须由应用层编码完成。
检查点存储位置建议使用共享存储(如 S3、NFS、云存储)以便跨节点访问。import boto3
s3 = boto3.client('s3')
s3.upload_file(local_path, bucket, key)
避免使用本地磁盘,防止节点故障导致数据丢失。
恢复策略Actor 初始化时优先从最新检查点恢复状态。def __init__(self):
self.counter = 0
self.load_checkpoint()
实现”启动即恢复”语义。
注意事项:原子性与一致性文件写入应保证原子性(如先写临时文件再重命名)。with open(tmp_path, "w") as f:
json.dump(data, f)
os.replace(tmp_path, final_path)
防止写入中途崩溃导致数据损坏。

7.4 Ray 的高可用(HA)集群配置

配置项说明注意事项
高可用头节点(HA Head Node)使用外部数据库(如 Redis 集群、MySQL)作为 GCS 后端,允许多个头节点。启动时使用:ray start --head --gcs-server-address=external_db:6379
多头节点部署支持多个头节点共享同一 GCS 存储,实现负载均衡和故障转移。需配置共享存储和网络互通。
工作节点自动注册工作节点自动发现并连接头节点,头节点故障后可重连新主节点。ray start --address=head:6379
使用 Kubernetes 部署通过 Ray Operator 部署,利用 K8s 的 Pod 自愈和滚动更新能力。配置 RayCluster CRD,设置 head 和 worker 的 replicas。
集群监控与告警集成 Prometheus + Grafana 监控集群状态,设置资源、任务失败等告警。Ray 暴露指标端点 /metrics
备份与恢复 GCS 数据定期备份 GCS 元数据(如 Actor、任务状态)。使用外部数据库的备份机制(如 Redis RDB/AOF)。
客户端重连机制Ray Client 支持自动重连到集群,连接中断后可恢复。ray.init("ray://cluster:10001", allow_reinit=True)
注意事项:网络稳定性高可用集群依赖稳定的网络;建议使用内网部署,避免公网抖动。配置合理的超时和重试参数。

第八章:Ray 集成模块(Ray Core 高级功能)

8.1 Ray Data:分布式数据处理

功能说明代码示例注意事项
创建 Dataset从多种数据源(本地文件、S3、Pandas、Arrow 等)创建分布式数据集。import ray

# 从文件加载
ds = ray.data.read_csv("s3://bucket/data/*.csv")

# 从 Pandas DataFrame
import pandas as pd
df = pd.DataFrame({"x": [1,2,3]})
ds = ray.data.from_pandas(df)
支持 CSV、JSON、Parquet、NumPy 等格式;自动分片。
数据转换(map, filter)对数据集进行并行转换操作。ds = ds.map(lambda row: {"x": row["x"] * 2})
ds = ds.filter(lambda row: row["x"] > 10)
mapfilter 函数在远程 Worker 上并行执行;返回新 Dataset。
聚合与统计支持内置聚合操作。mean = ds.mean("x")
max_val = ds.max("y")
grouped = ds.groupby("category").count()
自动并行化,结果返回给 Driver。
数据混洗(Shuffle)重分区或按键重分布数据。shuffled = ds.random_shuffle()
by_key = ds.sort("timestamp")
Shuffle 是昂贵操作,涉及跨节点数据传输;建议最小化使用。
批处理与迭代将数据以批的形式消费,用于模型训练或推理。for batch in ds.iter_batches(batch_size=2):
print(batch)
支持 pandas、numpy、arrow 等批格式。
与 Ray Dataset 集成可作为 Ray Train、Tune、Serve 的输入。train_ds, valid_ds = ds.train_test_split(test_size=0.2)实现端到端 ML 流水线。
注意事项Dataset 是惰性求值的(Lazy);适合 ETL、特征工程、预处理;不适合低延迟查询。使用 .materialize() 强制执行所有转换。

8.2 Ray Train:分布式模型训练

功能说明代码示例注意事项
集成主流框架支持 PyTorch、TensorFlow、XGBoost、LightGBM 等。from ray import train
import torch

def training_loop():
model = train.get_model()
for batch in train.get_dataset_shard("train"):
# 训练逻辑
train.report(loss=loss, accuracy=acc)
通过 train.report() 向 Driver 报告指标。
分布式训练器(Trainer)封装分布式训练逻辑,自动管理 Worker。from ray.train.torch import TorchTrainer
from ray.train import ScalingConfig

trainer = TorchTrainer(
training_loop,
scaling_config=ScalingConfig(num_workers=4),
datasets={"train": train_ds}
)

result = trainer.fit()
num_workers 控制并行训练进程数;支持 CPU/GPU。
数据分片(Dataset Shard)自动将 Dataset 分片给每个训练 Worker。train.get_dataset_shard("train")每个 Worker 获取数据的一个子集,实现数据并行。
检查点与恢复自动保存和恢复训练状态。train.save_checkpoint(epoch=epoch, model=model.state_dict())

checkpoint = result.checkpoint
支持故障后从检查点继续训练。
资源管理通过 ScalingConfig 指定资源。ScalingConfig(num_workers=4, use_gpu=True)可指定 resources_per_worker
与 Tune 集成可作为 Tune 的训练函数,实现分布式超参调优。TuneConfig(num_samples=10) + TorchTrainer实现大规模 HPO。
注意事项简化了分布式训练的样板代码;依赖 Ray Cluster 资源;需处理框架特定的分布式细节(如 DDP)。使用 TorchConfig(backend="ddp") 配置后端。

8.3 Ray Tune:超参数调优

功能说明代码示例注意事项
定义可调函数将训练过程包装为可调用函数,接收超参 config。def train_func(config):
lr = config["lr"]
model = Model(lr=lr)
for step in range(100):
loss = train_step(model)
tune.report(loss=loss, accuracy=acc)
必须使用 tune.report() 报告指标。
搜索算法支持网格搜索、随机搜索、贝叶斯优化(BayesOpt)、进化算法等。from ray.tune.search import BasicVariantGenerator

search_alg = BasicVariantGenerator()
高级算法如 OptunaSearch、BayesOptSearch 需额外安装。
调度器(Scheduler)控制试验的生命周期,支持早停(Early Stopping)。from ray.tune.schedulers import ASHAScheduler

scheduler = ASHAScheduler(metric="loss", mode="min")
ASHA、HyperBand 可加速收敛;PopulationBasedTraining 用于进化式调优。
启动调优实验提交调优任务,自动管理试验。from ray import tune

result = tune.run(
train_func,
config={
"lr": tune.loguniform(1e-4, 1e-1),
"batch_size": tune.choice([32, 64, 128])
},
num_samples=10,
search_alg=search_alg,
scheduler=scheduler,
metric="loss",
mode="min"
)
num_samples 控制试验总数;并行执行多个试验。
结果分析获取最佳超参和试验结果。best_result = result.get_best_result()
best_config = best_result.config
支持可视化(TensorBoard、分析 API)。
与 Train/RL 集成可调优 Ray Train 训练器或 RLlib 训练任务。tune.run(TorchTrainer, ...)实现端到端自动化调优。
注意事项大规模并行调优时需充足资源;选择合适的搜索算法和调度器;确保 tune.report() 频率合理。避免报告过多次,影响性能。

8.4 Ray Serve:模型服务部署

功能说明代码示例注意事项
定义可服务类(Deployment)将模型包装为可部署的类,支持 HTTP 和 gRPC 接口。from ray import serve
from fastapi import Request

@serve.deployment
class MyModel:
def __init__(self, model_path):
self.model = load_model(model_path)

async def __call__(self, request: Request):
data = await request.json()
return self.model.predict(data)
使用 @serve.deployment 装饰类。
部署与更新部署服务并支持滚动更新。app = MyModel.bind("path/to/model")
serve.run(app, name="my_model", route_prefix="/predict")

# 更新
serve.run(app_v2, name="my_model")
serve.run() 部署或更新服务;支持蓝绿部署。
流量路由与分片支持多副本、分片和负载均衡。@serve.deployment(num_replicas=4)
class ScaledModel:
...
num_replicas 控制副本数;自动负载均衡。
版本控制与金丝雀发布支持多版本并存和流量切分。serve.run(v1.bind(), name="model", traffic={"v1": 0.9, "v2": 0.1})实现灰度发布和 A/B 测试。
内置监控提供 Prometheus 指标(请求延迟、QPS 等)。访问 http://localhost:9999 查看指标。可集成 Grafana 可视化。
与 Ray 集成可调用其他 Ray 任务或 Actor。async def __call__(self, request):
preprocessed = await ray.get(preprocess.remote(data))
return self.model.predict(preprocessed)
构建复杂推理流水线。
注意事项适合在线推理服务;支持异步(async)提高吞吐;部署后通过 HTTP 端点访问。默认端口 8000;生产环境需配置 TLS 和认证。

8.5 Ray Workflows:长周期工作流管理

功能说明代码示例注意事项
定义工作流步骤使用 @workflow.step 装饰函数或类方法,使其可持久化。import ray
from ray import workflow

@workflow.step
def step1(x):
return x + 1

@workflow.step
def step2(y):
return y * 2
类似 @ray.remote,但支持检查点和恢复。
编排工作流通过调用 .step() 方法连接步骤,形成 DAG。result = step2.step(step1.step(1))
final = result.run("my_workflow")
run() 提交工作流;"my_workflow" 为唯一 ID。
持久化与恢复工作流状态自动持久化到外部存储(如 S3、文件系统)。workflow.init("s3://bucket/workflows")即使集群重启,也可通过 workflow.resume("my_workflow") 恢复。
容错性任意步骤失败后,可从最后一个检查点恢复,无需重跑整个流程。适用于可能运行数小时或数天的长周期任务。
并行执行支持并行执行独立步骤。branch1 = step1.step(1)
branch2 = step2.step(2)
final = combine.step(branch1, branch2)
自动并行化 DAG 中的独立分支。
检查点管理每个 workflow.step 完成后自动创建检查点。可通过配置控制检查点频率和存储位置。
注意事项已弃用:Ray Workflows 已被 Ray DAG 和 KubeRay 等方案取代。官方推荐迁移至新工作流系统;建议使用 Ray Pipelines 或 Prefect、Airflow 集成。

第九章:性能优化与调试

9.1 性能监控与指标采集

工具/方法用途使用方式注意事项
ray status查看集群整体资源使用情况(CPU、GPU、内存)和节点状态。在终端运行:ray status输出包括:总资源、已用资源、任务队列长度、对象存储内存等;是诊断资源瓶颈的第一步。
Prometheus + Grafana集成 Ray 暴露的 Prometheus 指标,实现可视化监控。Ray 自动暴露 /metrics 端点(默认端口 8080),配置 Prometheus 抓取并导入 Grafana 仪表盘监控关键指标如:
ray_task_scheduled_latency_seconds(任务调度延迟)
ray_worker_cpu_usage(Worker CPU 使用率)
ray_object_store_memory_used(对象存储内存占用)
ray memory分析 Object Store 中的对象内存占用,识别大对象或内存泄漏。ray memory --stats
ray memory --group-by=type
ray memory --group-by=stack
可按类型、堆栈跟踪分组;帮助定位创建大量小对象或未释放引用的代码位置。
自定义指标上报在任务或 Actor 中使用 ray.util.metrics 上报业务指标。from ray.util.metrics import Counter
req_counter = Counter("requests_total", description="Total requests")
@ray.remote
def handle_request():
req_counter.inc()
支持 Counter, Gauge, Histogram;可集成到现有监控体系。
Ray Client 日志启用详细日志查看连接、任务提交等客户端行为。ray.init(logging_level=logging.DEBUG)用于调试连接问题或任务提交失败。
注意事项• 建议在生产环境部署完整的监控栈(Prometheus + Grafana)避免过度采样影响性能;合理设置指标上报频率
• 定期运行 ray statusray memory 进行巡检

9.2 常见性能瓶颈分析

瓶颈类型表现诊断方法解决方案
CPU 瓶颈节点 CPU 使用率持续接近 100%,任务处理缓慢。ray status 查看 CPU 利用率、系统命令 top/htop• 优化算法效率
• 增加 num_cpus 并行度
• 将计算密集型任务拆分为更小粒度
GPU 瓶颈GPU 利用率低但任务排队,或显存不足(OOM)。nvidia-smi 查看 GPU 利用率和显存、ray status 查看 GPU 资源分配• 确保正确使用 CUDA 上下文
• 批量处理(Batching)提高利用率
• 使用混合精度训练
• 检查是否有显存泄漏
网络瓶颈大对象传输慢,跨节点通信延迟高。iftop/nethogs 查看网络流量、对比本地 vs 远程 ray.get() 时间• 使用零拷贝共享(Arrow 格式)
• 数据分片减少单次传输量
• 提升数据本地性(将任务调度到数据所在节点)
I/O 瓶颈读写磁盘或外部存储(S3)成为瓶颈。iotop 查看磁盘 I/O、Cloud Provider 监控 S3 吞吐• 使用高速存储(SSD、内存盘 _plasma_directory
• 异步 I/O 或预加载数据
• 缓存频繁访问的数据
调度瓶颈任务长时间排队,无法被调度执行。ray status 查看 “waiting” 任务数、检查资源是否碎片化• 调整任务/Actor 的 num_cpus/gpus 配置
• 减少小任务数量(批处理)
• 扩容集群
序列化瓶颈大对象或复杂对象序列化耗时过长。使用 cProfile 分析 ray.put()/.remote() 调用时间• 改用 Arrow 兼容格式(numpy, pandas)
• 实现自定义高效序列化器
• 避免传递大对象,改用 ray.put() + ObjectRef
注意事项性能问题是系统性的,需结合多种工具综合分析;避免”过早优化”应先测量再优化,确保解决的是真实瓶颈

9.3 内存与对象存储优化

优化策略说明实施方法注意事项
使用零拷贝数据结构对 numpy arrays、pandas DataFrames 等 Arrow 支持的类型,实现进程间零拷贝共享。确保数据为 Arrow 兼容格式:
import pyarrow as pa
arr = pa.array([1,2,3]) # 推荐
# 而非 list 或 dict
零拷贝仅适用于不可变数据;修改会触发复制。
避免大对象直接传递.remote() 参数中直接传递大对象会导致昂贵的序列化和传输。使用 ray.put() 显式存入 Object Store:
model_ref = ray.put(large_model)
task.remote(model_ref) # 传递引用
大对象建议 > 1MB 即采用此策略。
主动删除无用引用长期持有不再需要的 ObjectRef 会阻止垃圾回收,导致内存堆积。del ref # 删除引用
ray.delete([ref1, ref2]) # 强制删除
在循环或长生命周期任务中尤其重要。
配置 Plasma 存储路径将 Object Store 放在高速存储上提升性能。ray.init(
_plasma_directory="/dev/shm", # 使用内存盘
object_store_memory=10**10 # 限制大小
)
/dev/shm 是 tmpfs,速度极快;注意内存容量。
分片大对象超大对象(> 1GB)应主动分片处理。chunks = [data[i:i+size] for i in range(0, len(data), 10_000)]
refs = [process_chunk.remote(chunk) for chunk in chunks]
避免单个对象过大导致传输和 GC 压力。
监控对象内存定期检查 Object Store 使用情况。ray memory --stats --group-by=type识别 list, dict 等低效类型的大对象,优先转换为 Arrow 格式。
注意事项Object Store 内存独立于 Python 堆内存;del obj 不等于释放 Object Store 内存必须删除所有指向该对象的 ObjectRef 才能触发回收

9.4 调用链追踪与日志调试

方法说明使用方式注意事项
结构化日志记录在任务和 Actor 中使用标准 logging 模块输出结构化日志。import logging
logging.basicConfig(level=logging.INFO)
@ray.remote
def my_task():
logging.info("Starting task", extra={"task_id": ray.get_runtime_context().get_task_id()})
日志会聚合到 Driver 或重定向到文件;包含节点 IP、PID 信息。
分布式追踪(实验性)Ray 支持 OpenTelemetry 等追踪框架。需配置 OTLP 导出器,启用追踪上下文传播目前生态支持有限,主要用于高级场景。
获取运行时上下文获取当前任务/Actor 的元信息用于日志标记。ctx = ray.get_runtime_context()
print(ctx.get_actor_id())
print(ctx.get_job_id())
有助于关联同一作业或 Actor 的日志。
异常堆栈追踪Ray 自动捕获并传输远程异常的完整堆栈。try:
ray.get(failing_task.remote())
except Exception as e:
logging.exception("Task failed") # 包含完整远程堆栈
是调试远程错误的主要手段。
集中式日志收集使用 ELK (Elasticsearch, Logstash, Kibana) 或 Loki 收集 Ray 节点日志。配置 Filebeat/Fluentd 从 Ray 日志目录 (/tmp/ray/session_latest/logs) 收集实现跨节点日志搜索和分析。
使用 ray stack查看所有 Ray 进程的调用栈,诊断死锁或卡住任务。ray stack输出每个 Worker、Core Worker 的 Python 调用栈。
注意事项避免在高频任务中打印过多日志,影响性能建议使用不同日志级别(INFO, DEBUG)控制输出量

9.5 使用 Ray Dashboard 进行可视化分析

功能区域用途访问方式注意事项
Cluster 视图查看集群节点列表、资源使用率(CPU、GPU、内存)、运行时版本。浏览器访问 http://<head-node>:8265实时显示节点健康状态;红色表示资源紧张或节点失联。
Jobs 视图查看所有作业(Job)的状态、开始时间、运行时长、日志链接。Dashboard 主页的 “Jobs” 标签页每个 Driver 连接视为一个 Job;便于管理多用户环境。
Tasks 视图查看任务执行详情:名称、状态(RUNNING, FINISHED)、执行时间、资源消耗。“Tasks” 标签页可排序和过滤;点击任务查看详情(如输入参数、返回值大小)。
Actors 视图查看所有 Actor 实例的状态、类名、资源占用、创建时间。“Actors” 标签页可查看 detached Actor;对调试有状态服务非常有用。
Logs 视图聚合所有节点的日志,支持按节点、关键字搜索。“Logs” 标签页可实时流式查看日志,无需 SSH 登录各节点。
Metrics 视图内置 Grafana 面板,展示关键性能指标(任务延迟、对象存储、Worker 数量)。“Metrics” 标签页无需额外配置即可使用;是性能分析的核心工具。
Files 视图浏览集群节点上的文件系统(受限于权限)。“Files” 标签页可用于查看输出文件、检查点等。
注意事项• Dashboard 默认开启,端口 8265Dashboard 本身消耗少量资源,但对调试至关重要
• 生产环境应配置防火墙和认证(通过 ray.init(dashboard_host="0.0.0.0") 开放)
• 若无法访问,检查 head 节点防火墙和 ray start 日志

第十章:生产环境部署与最佳实践

10.1 Ray 集群部署方式(本地、Kubernetes、云平台)

部署方式说明工具/命令注意事项
本地单机部署用于开发、测试,启动一个包含所有组件的 Ray 实例。ray start --head --port=6379
# 或 Python API
import ray
ray.init()
默认模式;所有服务(GCS, Object Store)运行在同一进程或本地;不适合生产。
本地多节点集群手动在多个物理机或虚拟机上启动 Ray 节点,组成集群。Head 节点:
ray start --head --node-ip-address=<head-ip>
Worker 节点:
ray start --address=<head-ip>:6379
需手动管理节点生命周期;适用于小规模固定环境;网络配置复杂。
Kubernetes (K8s) 部署使用 KubeRay Operator 在 Kubernetes 上部署和管理 Ray 集群,实现弹性伸缩和高可用。定义 RayCluster CRD:
apiVersion: ray.io/v1
kind: RayCluster
metadata:
name: small-cluster
spec:
headGroupSpec:
serviceType: ClusterIP
template:
spec:
containers:
- name: ray-head
image: rayproject/ray:latest
workerGroupSpecs:
- replicas: 3
minReplicas: 1
maxReplicas: 10
template:
spec:
containers:
- name: ray-worker
image: rayproject/ray:latest
• 推荐生产环境使用
• 利用 K8s 的调度、自愈、扩缩容能力
• 需熟悉 K8s 生态
云平台 CLI 部署使用 Ray 自带的 ray up 命令在 AWS、GCP、Azure 等云平台快速创建集群。编写 cluster.yaml
provider:
type: aws
region: us-west-2
head_node_type: cpu_4xlarge
available_node_types:
cpu_4xlarge:
resources: {"CPU": 4}
gpu_p4:
resources: {"GPU": 1, "CPU": 8}
setup_commands:
- pip install -U ray[default]
然后执行:ray up cluster.yaml
• 快速搭建实验或批处理集群
• 支持自动扩缩容(Autoscaler)
• 成本需监控,避免忘记销毁
注意事项• 生产环境优先选择 K8s 或云平台部署避免混合部署模式导致管理混乱
• 确保节点间网络互通(开放端口 6379, 8265, 10001 等)
• 统一基础镜像和依赖环境

10.2 安全配置与访问控制

安全方面配置方法说明注意事项
Dashboard 访问控制设置 Dashboard 绑定地址和启用认证。ray.init(
dashboard_host="0.0.0.0", # 开放外部访问
# 认证需通过反向代理(如 Nginx)实现
)
默认仅限本地访问;生产环境应通过 TLS 反向代理添加用户名/密码认证。
客户端连接安全使用 Ray Client 并结合网络策略。ray.init("ray://head-node:10001")• 确保 10001 端口受防火墙保护
• 在 VPC 内部网络使用,避免暴露公网
数据加密启用传输加密(TLS)和静态加密。• 传输:目前 Ray 核心不直接支持 TLS,需通过隧道或代理
• 静态:对持久化数据(检查点、日志)使用加密存储(如 S3 SSE)
核心通信(GCS、Worker)目前为明文;敏感环境需网络层加密(如 IPSec)。
文件系统权限限制 Ray 进程对主机文件系统的访问。• 使用容器化部署(Docker/K8s)隔离文件系统
• 配置只读挂载非必要目录
防止恶意任务读取或篡改主机文件。
依赖安全扫描对 Docker 镜像和 Python 依赖进行漏洞扫描。使用 Trivy、Snyk 等工具:trivy image rayproject/ray:latest定期更新基础镜像和依赖库。
最小权限原则为 Ray Worker 分配最小必要权限。• K8s 中使用 RBAC 限制 Pod 权限
• 云平台中使用 IAM 角色限制访问范围
避免使用 root 用户运行 Ray 进程。
注意事项• Ray 目前安全功能相对基础,需结合外围设施加强建议将 Ray 部署在私有网络内
• 避免在公共网络直接暴露 Ray 集群

10.3 多租户与资源隔离

隔离机制实现方式说明注意事项
命名空间(Namespace)每个作业(Job)在逻辑上隔离,Actor 和任务名称可限定作用域。ray.init(namespace="team-a")
actor = MyActor.options(name="service", namespace="team-a").remote()
防止命名冲突;ray list actors 可按 namespace 过滤。
资源配额(Quota)通过集群资源总量和任务/Actor 资源请求实现软性隔离。• 集群总 CPU/GPU 有限
• 用户 A 的任务请求 num_cpus=4,可能因资源不足排队
Ray 无硬性配额管理,需用户自行规划资源使用。
Kubernetes Namespaces在 K8s 上为不同团队创建独立的 RayCluster 或使用单一集群但分 namespace。metadata:
namespace: team-a
name: ray-cluster-a
物理隔离更彻底;每个团队独占一组节点。
虚拟集群(Virtual Clusters)实验性功能,允许在共享物理集群上划分虚拟资源池。通过自定义资源和调度器插件实现需定制开发,尚未成熟。
作业级隔离每个 ray.init() 连接视为一个 Job,其创建的资源独立计费和监控。ray list jobs 查看所有作业便于多用户共享同一集群时进行审计和成本分摊。
容器化隔离使用 Docker/K8s 为不同租户的任务运行在独立容器中。每个 Worker Pod 可设置资源 limit/request实现 CPU、内存、网络的强隔离。
注意事项• 纯 Ray 单集群多租户是”尽力而为”的隔离防止”邻居噪声”(Noisy Neighbor)问题
• 关键业务建议采用物理或 K8s namespace 级隔离
• 监控各租户资源消耗

10.4 监控与告警集成

监控目标集成方案工具与配置注意事项
集群健康状态监控节点存活、资源使用率。ray status 输出解析
• Prometheus 指标:ray_node_statusray_cluster_resources_total
设置告警规则:节点失联、CPU/GPU 使用率 > 90% 持续 5 分钟。
任务与 Actor 状态跟踪任务延迟、失败率、Actor 数量。Prometheus 指标:
ray_task_state_count{state="RUNNING"}
ray_actor_state_count{state="ALIVE"}
ray_task_scheduled_latency_seconds_bucket
告警:任务积压过多(waiting 状态)、Actor 频繁崩溃重启。
对象存储内存防止 OOM 和内存泄漏。Prometheus 指标:ray_object_store_memory_usedray_object_store_memory_limit告警:Object Store 使用率 > 80%;结合 ray memory 定位大对象。
自定义业务指标上报模型推理延迟、QPS、错误率等。使用 ray.util.metrics.Counter, Gauge
latency_gauge = Gauge("inference_latency_ms")
latency_gauge.set(123.4)
将业务指标与系统指标关联分析。
日志聚合与分析集中收集和搜索日志。• ELK Stack (Elasticsearch, Logstash, Kibana)
• Grafana Loki + Promtail
• 云服务(AWS CloudWatch, GCP Stackdriver)
设置日志保留策略;索引关键字段(job_id, actor_id)。
告警通知异常时通知负责人。Prometheus Alertmanager 集成:
• 邮件
• Slack
• PagerDuty
• Webhook
配置合理的告警分级和静默时段,避免告警风暴。
注意事项• 建立完整的监控体系是生产化的基石监控本身也需监控(Monitor the Monitor)
• 定期审查告警规则有效性
• 使用 Dashboard 进行日常巡检

10.5 典型应用场景与架构模式

应用场景架构模式使用的 Ray 模块说明
大规模超参数调优(HPO)提交大量独立训练试验,并行探索超参空间。Ray Tune + Ray Train + Ray Data• Tune 管理试验生命周期
• Train 执行分布式训练
• Data 加载和预处理数据
• 典型于 AutoML、模型研发阶段
实时模型推理服务高并发、低延迟的在线预测 API。Ray Serve + (Ray Data / Ray AIR)• Serve 部署模型,提供 HTTP/gRPC 接口
• 支持多模型、A/B 测试、自动扩缩容
• 可集成预处理流水线
分布式 ETL 与特征工程处理海量原始数据,生成训练特征。Ray Data + Ray Tasks/Actors• Data 读取、转换、写入大规模数据集
• 自定义 Actor 实现复杂状态处理逻辑
• 输出 Parquet 文件供训练使用
强化学习(RL)训练分布式采集经验、并行训练、集中更新。RLlib + Ray Actors• RLlib 提供 PPO、DQN 等算法
• 自定义环境在远程 Actor 中运行
• 支持大规模分布式训练
长周期工作流编排执行多步骤、有依赖关系的批处理任务。Ray DAG / (原 Workflows) + Ray Tasks• 使用函数式组合构建 DAG
• 每步任务持久化结果
• 替代 Airflow 用于 Ray 内部任务编排
交互式数据分析在 Jupyter Notebook 中并行处理大数据。Ray Core + Ray Data• 无缝扩展 Pandas/Numpy 操作到分布式
• 适合数据科学家探索性分析
混合工作负载平台同一集群运行训练、推理、批处理等多种任务。Ray Core + Serve + Train + Tune• 统一基础设施,提高资源利用率
• 通过资源标注和调度实现隔离
• 典型于 MLOps 平台
注意事项• 根据场景选择合适的部署模式(如 Serve 适合在线,Tune 适合离线 HPO)• 架构设计时考虑容错和可观测性• 优先使用高级库(Data, Train, Serve)而非直接操作 Ray Core