Article
第一章: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 与其他分布式框架的对比
| 对比维度 | Ray | Apache Spark | Dask | multiprocessing |
|---|---|---|---|---|
| 编程模型 | 任务 + Actor(通用计算模型) | RDD/DataFrame(批处理为主) | Task Graph + Futures(类似Ray) | 进程/线程并发(单机) |
| 语言支持 | Python, Java(主要Python) | Scala, Python, Java, R | Python(为主) | 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或IPC | Pipe, 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=False) | Dashboard 提供集群状态、任务监控、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.remotedef 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.remotedef add(a, b): return a + bresult_ref = add.remote(2, 3) | 返回的是 ObjectRef,不是实际结果;必须通过 ray.get() 获取结果。 |
| 参数传递:基本类型 | 支持 int、float、str、bool 等内置类型自动序列化传输。 | @ray.remotedef greet(name): return f"Hello {name}"ref = greet.remote("Alice") | 序列化开销小,推荐用于轻量级数据。 |
| 参数传递:复杂对象 | 支持 list、dict、numpy array、pandas DataFrame 等。 | import numpy as np@ray.remotedef 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=…) | 设置最大等待时间(秒),超时抛出 GetTimeoutError。 | try: result = ray.get(f.remote(), timeout=5)except ray.exceptions.GetTimeoutError: print("Task timed out") | timeout=None 表示无限等待;有助于避免程序永久挂起。 |
| 异常传播 | 若远程任务内部抛出异常,ray.get() 会重新抛出相同异常。 | @ray.remotedef 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 = nameg = 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 = 0actor = 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.remoteclass Adder: def __init__(self, base): self.base = basea = 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.remoteclass CallableActor: def __call__(self, x): return x * 2a = 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.remotedef worker_task(counter): for _ in range(10): val = ray.get(counter.increment.remote()) return valc = Counter.remote()ref = worker_task.remote(c) | Actor 句柄可作为参数传递给任务;任务依赖 Actor 的状态。 |
| Actor 调用远程任务 | Actor 方法内部提交远程任务,实现异步计算委托。 | @ray.remoteclass 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.remotedef consumer(data_ref): data = ray.get(data_ref) ...consumer.remote(result_ref) | Ray 自动解析依赖,确保 data_ref 就绪后再调度 consumer。 |
| ObjectRef 的序列化 | ObjectRef 可被序列化并传递到其他任务或 Actor。 | @ray.remotedef forward_ref(ref): return refnew_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 nparr = np.ones((1000, 1000))@ray.remotedef process(arr): return arr.sum()ref = process.remote(arr) | 大幅减少内存占用和传输延迟;仅适用于 Arrow 支持的不可变类型。 |
| 默认序列化(Pickle) | 对不支持 Arrow 的对象使用 Python pickle 进行序列化和反序列化。 | class CustomObj: def __init__(self, x): self.x = x@ray.remotedef use_custom(obj): return obj.x | Pickle 通用但较慢,且可能不安全;大对象序列化开销大。 |
| Arrow 支持类型 | 包括:numpy.ndarray(特定 dtype)、pandas.DataFrame、pyarrow.Table、bytes、str、list、dict(小对象)。 | df = pd.DataFrame({"a": [1,2,3]})ray.put(df) # 零拷贝 | 推荐使用 Arrow 兼容格式以获得最佳性能。 |
| 自定义序列化器 | 用户可注册特定类型的高效序列化函数。 | import pyarrow as paray.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_workers 和 max_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 memory、ray status 和 Dashboard 查看资源消耗。 | ray memory --statsray 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 表示无限重试(谨慎使用)。 |
| 重试条件 | 仅对系统级错误(如节点崩溃、网络中断)或用户异常进行重试;异常包括:RayTaskError、RayActorError 等。 | 不重试资源不足(如 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 boto3s3 = 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 DataFrameimport pandas as pddf = 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) | map 和 filter 函数在远程 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 trainimport torchdef 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 TorchTrainerfrom ray.train import ScalingConfigtrainer = 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 BasicVariantGeneratorsearch_alg = BasicVariantGenerator() | 高级算法如 OptunaSearch、BayesOptSearch 需额外安装。 |
| 调度器(Scheduler) | 控制试验的生命周期,支持早停(Early Stopping)。 | from ray.tune.schedulers import ASHASchedulerscheduler = ASHAScheduler(metric="loss", mode="min") | ASHA、HyperBand 可加速收敛;PopulationBasedTraining 用于进化式调优。 |
| 启动调优实验 | 提交调优任务,自动管理试验。 | from ray import tuneresult = 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 servefrom fastapi import Request@serve.deploymentclass 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 rayfrom ray import workflow@workflow.stepdef step1(x): return x + 1@workflow.stepdef 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 --statsray memory --group-by=typeray memory --group-by=stack | 可按类型、堆栈跟踪分组;帮助定位创建大量小对象或未释放引用的代码位置。 |
| 自定义指标上报 | 在任务或 Actor 中使用 ray.util.metrics 上报业务指标。 | from ray.util.metrics import Counterreq_counter = Counter("requests_total", description="Total requests")@ray.remotedef handle_request(): req_counter.inc() | 支持 Counter, Gauge, Histogram;可集成到现有监控体系。 |
| Ray Client 日志 | 启用详细日志查看连接、任务提交等客户端行为。 | ray.init(logging_level=logging.DEBUG) | 用于调试连接问题或任务提交失败。 |
| 注意事项 | • 建议在生产环境部署完整的监控栈(Prometheus + Grafana) | 避免过度采样影响性能;合理设置指标上报频率 | |
• 定期运行 ray status 和 ray 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 paarr = 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 logginglogging.basicConfig(level=logging.INFO)@ray.remotedef 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 默认开启,端口 8265 | Dashboard 本身消耗少量资源,但对调试至关重要 | |
• 生产环境应配置防火墙和认证(通过 ray.init(dashboard_host="0.0.0.0") 开放) | |||
• 若无法访问,检查 head 节点防火墙和 ray start 日志 |
第十章:生产环境部署与最佳实践
10.1 Ray 集群部署方式(本地、Kubernetes、云平台)
| 部署方式 | 说明 | 工具/命令 | 注意事项 |
|---|---|---|---|
| 本地单机部署 | 用于开发、测试,启动一个包含所有组件的 Ray 实例。 | ray start --head --port=6379# 或 Python APIimport rayray.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/v1kind: RayClustermetadata: name: small-clusterspec: 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-2head_node_type: cpu_4xlargeavailable_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_status、ray_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_used、ray_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 |