Ray分布式Python运行时:一套API搞定单机到集群并行
发布时间:2026/9/26 16:54:45 作者:尧图编辑部 阅读量:1,286

先说结论如果你正在写 Python 代码且发现单机跑得慢、数据量大到内存顶不住、或者想在 GPU 集群上快速铺开一个训练/推理任务直接上 Ray 会比你去啃那套老旧的 MPI 或者 Spark 要舒服得多。Ray 不是一个服务框架也不是一个消息队列它是一层很薄的分布式运行时核心卖点就是“一套 API 把单机并发写到集群并行”。我最早接触 Ray 是因为做强化学习RLlib 当时是唯一能把实验从一台笔记本平滑搬到集群上的工具后来发现它其实可以管任意 Python 函数和类才意识到这东西的本质是“把 Python 的并发模型从线程/进程扩展到了整台集群”。今天这篇想把 Ray 从概念到落地完整拆一遍重点讲清楚它的统一 API 到底在解决什么问题以及在真实业务里怎么最快用起来、有哪些坑。1. Ray 的定位不是大数据框架是分布式 Python 运行时1.1 传统分布式为什么让人头疼在 Ray 出现之前搞分布式计算基本是三派天下。大数据派用 Spark核心抽象是 RDD/DataFrame擅长 SQL 分析和批处理但你想在里面写点复杂的状态逻辑或者做在线推理就感觉怎么都不顺手因为它本质上是一个“按 DAG 分阶段执行”的批处理引擎。另一派是高性能计算MPI 把通信露在外面写起来极其底层每个进程管理消息收发训练一个神经网络都要自己规划数据分发说实话不太适合普通业务开发。还有一派是“自己造轮子”用 Python 自带的 multiprocessing 写进程池或者用消息队列串任务单机还行一旦跨机器就要处理节点发现、失败重试、任务依赖整个工程复杂度会迅速失控。我当时踩过的典型坑是一个数据处理流程单机跑要十二个小时领导说“你上集群呀”结果我花了一周研究 Spark 接入写出来的逻辑既绕又丑——明明一个简单的 map 加 shuffle在 Spark 里要套一堆 transformation还得考虑到 JSON schema 的兼容性。后来想通了问题的核心不是“用什么框架”而是“分布式不该改变我表达业务的方式”。1.2 Ray 的“统一 API”到底统一了什么Ray 由 UC Berkeley RISELab 发起设计目标很明确做一个面向 AI 应用场景的通用分布式编排层让开发者用普通的 Python 函数、类、库就能完成分布式任务不需要重新学习一套 SQL 或者 MapReduce 抽象。Ray 的统一 API 表述上其实包含了四层任务级统一同一个函数在单机可以普通调用在集群就用ray.remote装饰器变成一个分布式任务不用改函数体内逻辑。状态级统一通过 Actor 在集群中维护有状态的对象不用自己搞分布式锁和共享存储。数据级统一Ray Data 提供了类似 DataFrame 的接口但它底层是分布式的对象集数据加载、预处理、训练集构建全都在同一套 API 下完成。服务级统一训练完的模型直接用 Ray Serve 拉起推理服务不需要再另起一套 Flask/FastAPI 架构直接在同一个 Ray 集群里实现弹性扩缩容。关键是所有组件都共用同一个底层调度器和对象存储所以“数据在哪里”、“计算在哪里”是可以互相感知的。这和传统大数据 Hadoop 家族用一堆子项目拼凑完全不同。用生活化类比解释Spark 像一条标准化的流水线你按工位送料它按流程产出Ray 更像一个调度中枢任何“干活的人”都能被实时派单还能给“带状态的监督员”安排常驻工位。所以Ray 最核心的价值不是某个单独的 API而是“一套 API 管理整个分布式生命周期”。从任务提交、数据分片、资源调度、故障恢复到可视化监控全都内聚在同一个运行时里。你只管写逻辑剩下的事交给 runtime。2. 核心组件与编程范式拆解2.1 Task把函数变成远程函数Ray 最基础的抽象是 Task它对标的就是“远程异步函数”。用法如下import ray ray.init() ray.remote def add(a, b): return a b # 远程调用立即返回一个 ObjectRef类似 Future ref add.remote(1, 2) # 获取结果 result ray.get(ref) print(result) # 3看到ray.remote之后最大的变化是什么add(...)变成了add.remote(...)前者是本地直接执行后者是“把任务对象连同参数一起提交给调度器由某个 worker 执行结果存到分布式对象存储里”。这个抽象最爽的一点是函数体内完全没有分布式代码不需要 socket、不需要消息序列化协议、不需要你自己处理重试。参数传递、结果回收、异常传播都由框架完成。如果你有一批相互独立的任务比如同时调用 100 个函数代码非常直观futures [add.remote(i, i) for i in range(100)] results ray.get(futures)这一句的等价物如果用 multiprocessing 写要管理进程池、队列、结果收集如果跨机器还要自己设计主从但在 Ray 里ray.get总会等着全部完成而且结果顺序与提交顺序一致。2.2 Actor有状态的工作进程Task 解决的是“无状态函数”的并行化但现实世界很多场景需要状态一个计数器、一个模型实例、一个模拟环境、一个数据库连接池。Ray 的 Actor 就是用来干这个的。ray.remote class Counter: def __init__(self): self.value 0 def increment(self): self.value 1 return self.value def get_value(self): return self.value # 创建一个 Actor 实例 counter Counter.remote() # 调用 Actor 方法也是异步的 final_value ray.get(counter.increment.remote()) print(final_value) # 1注意Actor 的方法调用默认按顺序执行同一个 Actor 实例内部是串行的但是不同 Actor 实例之间天然并行。你把Counter.remote()创建 10 次每个 Counter 有自己独立的状态这就相当于有了 10 个自治的“有状态工作进程”。实际业务里最常见的用法是“每个 Actor 加载一个模型副本然后接收一批推理请求”。因为模型加载是大开销不能每次递归加载而 Actor 可以做到“初始化一次、常驻内存、反复调方法”这在 GPU 场景中尤其重要。2.3 ObjectRef 与对象存储ray.get拿到的其实是一个ObjectRef它是分布式对象存储中的引用句柄。Ray 默认使用共享内存做对象存储进程间传递大对象不必走一遍序列化-反序列化而是通过零拷贝方式读取最大程度减少传输开销。你可以显式把一个大对象塞进对象存储big_data list(range(1_000_000)) ref ray.put(big_data) ray.remote def process(ref): data ray.get(ref) return sum(data) result ray.get(process.remote(ref))这样做的意义在于你不需要把 big_data 作为参数传给每个任务只会传一个指向共享内存的引用。多任务并发读同一批数据的时候这种设计非常省钱。还有个隐藏的技巧ray.get支持指定 timeout不会无限阻塞try: result ray.get(ref, timeout5) except ray.exceptions.GetTimeoutError: print(任务超时了先做别的)2.4 调度器与资源控制Ray 使用自研的分布式调度器原理上不是全局集中式调度而是“逻辑集中、物理分布式”。每个节点本地有个调度器全局通过 gossip 协议交换状态。这样既避免了单点瓶颈又能做局部性感知——把任务尽量调度到数据所在的节点减少跨节点通信。资源显式声明也很重要ray.remote(num_cpus2, num_gpus0.5) def heavy_task(): ...num_cpus表示这个任务占用多少 CPU 核num_gpus可以是浮点数比如 0.5 说明两个任务共享一张 GPU。资源算清楚后Ray 调度器就像一个“管家”如果集群总共 16 核四个num_cpus4的任务正好全部占满再来第五个任务就排队等待。我遇到比较多的死锁就是“任务自己占着资源然后又去等待其他任务释放资源”。比如一个 driver 本身默认占用 1 个 CPU 资源如果你在 driver 里提交大量任务同时又用ray.get等待结果而集群资源已经被任务占满driver 又没有多余 CPU 执行后续逻辑就会卡死。解决办法是在ray.init(num_cpus8)时给 driver 预留资源或者在任务内部不要嵌套海量子任务等待。3. 从零到一实操一个完整的 Ray 项目3.1 环境准备三步骤安装 Ray 非常简单用 pip 搞定pip install ray[default]这个[default]扩展包会附带 dashboard、cluster launcher 等常用工具建议直接装。如果只是最核心的 APIpip install ray也可以但没有可视化监控。装完后激活一个本地集群import ray # 不传参数时自动检测本机资源并启动一个本地模式集群 ray.init() # 也可以显式限制资源 ray.init(num_cpus4, object_store_memory2 * 1024 * 1024 * 1024)启动成功后可以在浏览器打开http://127.0.0.1:8265看 dashboard——节点状态、每个任务耗时、内存占用一应俱全。我建议第一次用的人先开 dashboard 观察一下因为你肉眼能看见 Tasks 在 worker 之间分发对理解“分布式”有很强的直观帮助。3.2 第一个分布式程序并行任务编排实战假设现在要批量计算一堆 URL 的可访问性单机串行可能要 5 分钟多线程又怕某些库不是线程安全的。用 Ray 的做法import ray import requests ray.init() ray.remote(max_retries3) def check_url(url): try: resp requests.head(url, timeout10) return url, resp.status_code except requests.RequestException as e: return url, str(e) urls [fhttps://example.com/{i} for i in range(100)] futures [check_url.remote(url) for url in urls] results ray.get(futures) for url, status in results[:10]: print(f{url}: {status})这个程序里可能涌现出几个问题为什么max_retries3因为网络请求经常有瞬时抖动重试可以兜底。如果requests库本身不是线程安全的在不同 worker 里跑完全没问题因为每个 task 在独立进程中执行。100 个任务会同时并发吗不一定。Ray 会根据集群核数决定并发度8 核机器默认最多 8 个任务并行其他任务排队。你可以通过num_cpus0.5强制提升并发度但过度并行也可能把带宽打满。3.3 进阶实战用 Ray 批量调用大模型 API最近几年大模型 API 应用突然爆发几乎每个项目都要调 OpenAI、DeepSeek 之类的接口。单条 API 调用看起来很快但当你有一万条文本需要批量处理时串行就非常痛苦。用 Ray 做并行调用大模型 API 有一个天然优势不掉進复杂的事件循环框架也不用自己维护线程池只要写普通函数就行。下面是一个示例使用 DeepSeek 兼容接口批量处理文本摘要import ray import time from openai import OpenAI ray.init(num_cpus8) client OpenAI( api_keyYOUR_API_KEY, base_urlhttps://api.deepseek.com, ) ray.remote(max_retries2) def summarize_with_llm(text): try: response client.chat.completions.create( modeldeepseek-chat, messages[ {role: system, content: 你是文本摘要助手输出简洁摘要。}, {role: user, content: text}, ], temperature0.3, ) return response.choices[0].message.content except Exception as e: # API 调用失败时抛出特定格式让 Ray 记录并重试 raise RuntimeError(fAPI call failed: {e}) from e texts [很长的一段文本1, 很长的一段文本2] * 100 # 200条 start time.time() futures [summarize_with_llm.remote(t) for t in texts] results ray.get(futures) elapsed time.time() - start print(f200 条文本并行调用耗时 {elapsed:.2f} 秒)使用体验有几点心得。一是限流容错。大模型 API 通常有速率限制你一口气发 200 个请求非常容易触发 429 错误。Ray 的max_retries只能解决“崩溃重试”不能解决“请求过多被拒”。我的做法是给远程函数外层包一个限速器或者细分批次10 个请求一组组内用信号量控制最大并发。二是异常要抛出来。max_retries依赖任务抛出异常才能触发重试。如果你在函数里把所有异常都吞掉、最终返回一个缺省值框架会认为任务成功真正的失败就被掩盖了。三是成本控制。并行度不是越高越好。API 是按 token 计费的你同时发 200 个请求中间任何一个模型返回超出上下文的错误整个批次都会受影响。所以我还是建议在真正线上跑之前先拿 20 条数据做并发度测试找到本地并行度和 API 限流之间的平衡点。3.4 模型服务部署Ray Serve 实验模型训练完了总要对外提供推理接口。Ray Serve 是 Ray 生态里的模型服务组件它和 Spark 的服务化完全不一样直接把ray.remote的 Actor 变成 HTTP 服务。from ray import serve import ray ray.init() serve.start() serve.deployment(ray_actor_options{num_gpus: 0}) class SummarizeService: def __init__(self): self.model mock-model async def __call__(self, request): text await request.json() return {summary: fprocessed: {text[text][:20]}} # 部署服务 serve.run(SummarizeService.bind())部署完成后可以通过 HTTP 调用curl -X POST http://127.0.0.1:8000/ -H Content-Type: application/json -d {text: This is a long article...}Ray Serve 最实用的地方是支持多副本自动扩缩容、灰度发布和请求级负载均衡。你不需要额外部署一整套 K8s 配置来承载模型服务用 Ray 内部机制就能完成。但我要提醒如果你们的服务是超大规模公网网关Ray Serve 不一定比专门的网关框架更合适。因为 Ray Serve 的定位是“模型服务与计算编排的融合”不是高并发 Web 网关。它更适合的业务是“推理逻辑复杂、需要和数据处理/训练管线无缝衔接”的场景。4. 常见问题与排查实录4.1 初始化与连接异常最常见的初始化问题有两个。一是RayContext已经存在重复调用ray.init()会报 “A Ray runtime has already been started”。解决方式是用ray.shutdown()先关闭再重启或者让程序里所有模块共享同一个 runtime不要到处初始化。二是集群资源不足。比如你在一个 4 核机器上跑ray.init(num_cpus16)会报 “At least 16 CPUs are needed, but only 4 are available”。这时要么修改num_cpus要么升级机器资源。这个报错信息看起来硬其实是个善意的保护防止任务调度时无限等待。4.2 资源占满与任务死锁我前面提过“driver 也占用资源”这个点再展开讲一个具体场景。假设集群 8 核你从 driver 提交了 8 个任务每个任务都num_cpus1。然后你又在下边调用ray.get(futures)等待结果这本身没问题因为 driver 已经被占了 1 个 CPU提交完 8 个任务后集群还有资源给 driver 运行等待逻辑。但如果这 8 个任务内部又各自提交了 4 个子任务这些子任务也占 CPU而父任务在等待子任务完成此时 CPU 已经被占满整个集群进入互相等待的循环。解决手段有三个在ray.init时把num_cpus设置为实际核数减一预留一部分给 driver或者在父任务里用ray.get加 timeout让它在等待时主动退出释放资源或使用ray.util.placement_group精细控制父子任务资源。4.3 对象存储内存不足大批量把数据塞进对象存储很容易触发 “ObjectStoreMemoryError”。默认对象存储是按主内存比例分配的通常 30%但不是所有场景都够用。比如你要加载一个 200GB 的数据文件预处理默认配置就很难受。解决办法是在ray.init()里显式调大object_store_memory但这只是缓兵之计因为物理内存是有限的。更合理的思路是利用 Ray Data 的流式处理避免所有数据一次性进内存import ray ds ray.data.read_csv(s3://bucket/path/*.csv) # 这时并没有把所有数据加载进内存而是逻辑图 ds ds.map(lambda row: {value: row[value] * 2}) # 按批次触发执行 for batch in ds.iter_batches(batch_size1000): process_local(batch)这种“惰性执行”模式和 Spark 很像但和ray.put()的“全部塞进内存”是两种风格需要根据场景灵活选择。4.4 序列化失败的疑难杂症Ray 在跨进程传递函数和数据时依赖 cloudpickle但某些类或库的对象无法被序列化。比如你在远程函数里引入了threading.Lock或者自定义了一个引用了文件句柄的对象这时候大概率报序列化错误。一个通用的排查思路是如果对象必须跨进程就改成可序列化形式比如 lock 可以改成进程级安全状态如果对象确实不能序列化就放在 Actor 中初始化不让它跨进程传递。比如数据库连接池、模型实例就不适合作为参数传到 Task 里而应该在Actor.__init__里创建保证它只存在于某个固定进程。4.5 常见问题速查表现象可能原因解决方案ray.init()重复启动报错运行时已经存在使用ray.shutdown()后重新初始化任务一直排队不执行资源不足或声明过大检查 dashboard调整num_cpus/num_gpusGetTimeoutError任务超时设置合理的 timeout 或优化任务逻辑对象存储内存溢满数据大量塞入增大object_store_memory或改用 Ray Data 流式远程函数内创建远程函数报错嵌套远程函数不支持直接调用把内部函数写成普通函数或使用 Actor 中转模型 API 返回 429请求并发过高添加限速逻辑分批提交请求Actor 频繁崩溃代码逻辑有状态破坏设置max_restarts利用ray.util.state.ActorState检测排查技巧就一个打开 dashboard看任务调度的时间线和资源占用图。80% 的问题在 dashboard 上都能一眼定位不用盲目猜。5. 选型对比什么时候用 Ray什么时候别用5.1 跟 Spark、MPI、Dask 的横向对比很多人问过我“Ray 和 Spark 到底选哪个”我的回答是看你要算的东西长什么样。维度RaySparkMPIDask核心抽象远程函数 / Actor / Ray DataRDD / DataFrame消息通信原语任务图对 Python 支持一等公民DataFrame 为主UDF 勉强差一等公民有状态服务强力支持Actor不太适合进程状态自行管理有限支持动态任务依赖天然支持 DAG 级动态DAG 静态构建手动设计DAG 动态支持AI 场景契合度很高训练 推理 数据中等数据清洗为主低中等运维复杂度中等可单机也可集群较高高较低社区活跃度很活跃成熟但偏冷稳定活跃核心结论Spark 是做“SQL 分析、ETL、数据仓库”场景的王者MPI 是 HPC 数值计算的老牌方案Dask 和 Ray 都在做 Python 任务调度但 Dask 更偏向 DataFrame 和科学计算而 Ray 的目标是 AI 应用的全栈编排。如果在项目里已经大量用 Pandas/NumPy/SklearnDask 的上手成本更低但如果你要写强化学习、要托管模型、要跨计算与数据两个层面统一调度Ray 的生态更齐全。还有一点Ray 的未来方向也越来越清晰服务端 Serverless 化。Ray 现在支持serve.run一键部署也支持跟 KubernetesKubeRay集成从开发到生产的环境过渡很顺滑。这个优势不是 Spark 能给的。5.2 大模型时代的真实开场大模型时代对计算框架的诉求变得非常具体训练要大集群推理要低延迟数据预处理要和模型输入直接联动还要处理多租户、混合部署。Ray 几乎就是照着这个需求设计的。比如最近比较常见的做法用 HuggingFace 的 trainer 做分布式训练底层是 Ray Train训练完后用 Ray Serve 管理 GPU 推理服务自动扩缩容数据处理阶段用 Ray Data 从 S3/HDFS 上批量加载语料用同一个集群完成“数据预处理 → 训练 → 服务发布”的全流程。这样一个开发者就可以同时做训练和推理不必维护两套完全不同的基础设施。说个更具体的如果你要基于开源模型微调一个行业大模型数据量可能从几万条涨到几百万条。传统做法是先把数据导到业务库里清洗一遍然后导成 CSV 发给训练脚本中间各种格式转换、数据版本管理非常痛苦。用 Ray 可以把这些步骤写进同一个 Pipeline读取、清洗、tokenize、分片、训练数据生成一气呵成。这种“数据管道和训练代码零切换”的能力对开发效率提升非常明显。当然Ray 不是银弹。如果你的业务只有 10 万条数据、单机 Pandas 也能跑完那没必要引入分布式复杂度。我默认的建议是“单机先跑通Ray 做加分项”而不要一开始就在纯单机场景里堆一个分布式框架。6. 我个人在实际操作中的体会最后说点虚的也最真实。Ray 用了几年下来我最感慨的其实不是它的性能而是它改变了我对“并行”的思考方式。以前写并行要老惦记“那几个 worker、什么队列、什么锁”写出的代码像是做手术到处是缝合痕迹。用 Ray 之后我可以先按单机逻辑把“函数”和“类”写清楚再想哪些可以并行拆分然后加个装饰器剩下的交给调度器。这种思维转变比任何性能调优都重要。一个小技巧收尾如果你的需求是“把现有 Python 脚本变成分布式”优先动装饰器不要动函数内部逻辑——先在函数头上加ray.remote然后把调用点从f(x)改成f.remote(x)最后拿结果时用ray.get三步走完。很多时候你不用重写任何算法就已经完成了最基本的分布式化。再提到 API 时代的一点。现在不管是大模型调用、第三方数据服务还是内部微服务大家都在和 API 打交道。Ray 提供的这层统一 API 其实和“外部 API 网关”不是一个层次的东西——它管的是你计算资源的编排而不是通信协议的转换两者可以完美共存。所以别担心 Ray 会吞掉你现有的 API 基建它只会让调用这些 API 的计算过程更灵活、更快、更不容易崩。我的建议是在下一个需要并行化处理的 Python 任务里试着把第一段代码装饰上ray.remote感受一下那种“单机到集群的丝滑过渡”你会用上瘾的。