☰
Ray分布式Python运行时:一套API搞定单机到集群并行
2026/9/26 16:54:42 网站建设 项目流程

先说结论:如果你正在写 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, timeout=5) except ray.exceptions.GetTimeoutError: print("任务超时了,先做别的")

2.4 调度器与资源控制

Ray 使用自研的分布式调度器,原理上不是全局集中式调度,而是“逻辑集中、物理分布式”。每个节点本地有个调度器,全局通过 gossip 协议交换状态。这样既避免了单点瓶颈,又能做局部性感知——把任务尽量调度到数据所在的节点,减少跨节点通信。

资源显式声明也很重要:

@ray.remote(num_cpus=2, num_gpus=0.5) def heavy_task(): ...

num_cpus表示这个任务占用多少 CPU 核,num_gpus可以是浮点数,比如 0.5 说明两个任务共享一张 GPU。资源算清楚后,Ray 调度器就像一个“管家”:如果集群总共 16 核,四个num_cpus=4的任务正好全部占满,再来第五个任务就排队等待。

我遇到比较多的死锁就是“任务自己占着资源,然后又去等待其他任务释放资源”。比如一个 driver 本身默认占用 1 个 CPU 资源,如果你在 driver 里提交大量任务,同时又用ray.get等待结果,而集群资源已经被任务占满,driver 又没有多余 CPU 执行后续逻辑,就会卡死。解决办法是在ray.init(num_cpus=8)时给 driver 预留资源,或者在任务内部不要嵌套海量子任务等待。

3. 从零到一:实操一个完整的 Ray 项目

3.1 环境准备三步骤

安装 Ray 非常简单,用 pip 搞定:

pip install "ray[default]"

这个[default]扩展包会附带 dashboard、cluster launcher 等常用工具,建议直接装。如果只是最核心的 API,pip install ray也可以,但没有可视化监控。

装完后激活一个本地集群:

import ray # 不传参数时,自动检测本机资源并启动一个本地模式集群 ray.init() # 也可以显式限制资源 ray.init(num_cpus=4, object_store_memory=2 * 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_retries=3) def check_url(url): try: resp = requests.head(url, timeout=10) return url, resp.status_code except requests.RequestException as e: return url, str(e) urls = [f"https://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_retries=3?因为网络请求经常有瞬时抖动,重试可以兜底。
  • 如果requests库本身不是线程安全的,在不同 worker 里跑完全没问题,因为每个 task 在独立进程中执行。
  • 100 个任务会同时并发吗?不一定。Ray 会根据集群核数决定并发度,8 核机器默认最多 8 个任务并行,其他任务排队。你可以通过num_cpus=0.5强制提升并发度,但过度并行也可能把带宽打满。

3.3 进阶实战:用 Ray 批量调用大模型 API

最近几年大模型 API 应用突然爆发,几乎每个项目都要调 OpenAI、DeepSeek 之类的接口。单条 API 调用看起来很快,但当你有一万条文本需要批量处理时,串行就非常痛苦。用 Ray 做并行调用大模型 API 有一个天然优势:不掉進复杂的事件循环框架,也不用自己维护线程池,只要写普通函数就行。

下面是一个示例,使用 DeepSeek 兼容接口批量处理文本摘要:

import ray import time from openai import OpenAI ray.init(num_cpus=8) client = OpenAI( api_key="YOUR_API_KEY", base_url="https://api.deepseek.com", ) @ray.remote(max_retries=2) def summarize_with_llm(text): try: response = client.chat.completions.create( model="deepseek-chat", messages=[ {"role": "system", "content": "你是文本摘要助手,输出简洁摘要。"}, {"role": "user", "content": text}, ], temperature=0.3, ) return response.choices[0].message.content except Exception as e: # API 调用失败时,抛出特定格式让 Ray 记录并重试 raise RuntimeError(f"API 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(f"200 条文本并行调用耗时 {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": f"processed: {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_cpus=16)会报 “At least 16 CPUs are needed, but only 4 are available”。这时要么修改num_cpus,要么升级机器资源。这个报错信息看起来硬,其实是个善意的保护,防止任务调度时无限等待。

4.2 资源占满与任务死锁

我前面提过“driver 也占用资源”这个点,再展开讲一个具体场景。假设集群 8 核,你从 driver 提交了 8 个任务,每个任务都num_cpus=1。然后你又在下边调用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_size=1000): 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_gpus
GetTimeoutError任务超时设置合理的 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/Sklearn,Dask 的上手成本更低;但如果你要写强化学习、要托管模型、要跨计算与数据两个层面统一调度,Ray 的生态更齐全。

还有一点,Ray 的未来方向也越来越清晰:服务端 Serverless 化。Ray 现在支持serve.run一键部署,也支持跟 Kubernetes(KubeRay)集成,从开发到生产的环境过渡很顺滑。这个优势不是 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,感受一下那种“单机到集群的丝滑过渡”,你会用上瘾的。

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

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

立即咨询