如果你和我一样,日常要同时处理数据清洗、模型训练和在线推理这三摊子事,大概率会在某个时间点被同一件事逼疯:每个环节都有一套自己的分布式框架。数据预处理要用 Spark,分布式训练要自己拼多机同步逻辑,上线推理服务还得再学一套新的接口。工具之间互不了解,光是把训练好的模型传给线上服务就够折腾几个通宵。Ray 正是冲着这个痛点来的,它把分布式计算抽象成一套统一 API,底层帮你把调度、通信、容错全部扛下来。它不只是一个任务调度器,更像是一种重新组织分布式编程思路的范式:你在本地怎么写 Python,在集群上就怎么写,把函数变成分布式的唯一动作,是给它加一个装饰器。
这篇文章不是 Ray 的 API 手册,而是从我用它把数据处理、模型调参、推理上线一口气跑通的经验出发,把 Ray 为什么能统一这么多场景的底层逻辑、三个核心原语怎么用、集群落地时的工程细节,以及和 Spark、Dask 的选型边界一次讲透。适合那些已经厌倦在多个分布式框架之间来回切换、想用一套东西把 AI 工程链路串起来的人。
有人说 Ray 是"分布式计算的 Python 化",这个说法我大体认同,但只说对了一半。Python 只是它的门面,真正让它能重塑范式的,是它把批处理框架、流计算框架、任务调度器、参数服务器都看成了同一件事的不同形态。怎么做到的?往下看。
1. 传统分布式框架为什么和 AI 负载"八字不合"
1.1 批处理模型的隐藏假设:计算图在执行前就定型
Hadoop 和 Spark 的核心理念,是把任务预先画成一张静态 DAG,调度器拿着这张图分段执行。这个模型有一系列隐含假设:计算流程是确定的,每个阶段的输入输出是清晰的,某个节点失败了,整个 Stage 重跑是划算的。这些假设在传统数据分析场景里完全成立,但一碰到真实的 AI 训练管线就开始漏气。
举个例子,跑强化学习环境交互的时候,下一步要计算什么取决于当前这一步动作的结果,这是典型的动态控制流;做超参搜索的时候,配置空间不是写死的一层循环,而是要动态看指标、动态剪枝、动态增删试验点。这种"计算图边跑边变"的负载,硬塞进一个提前规划好执行计划、阶段边界必须清晰的框架里,写出来的代码你自己都认不出来。有人形容批处理框架像高铁运行图:发车时间、停靠站、编组全固定,适合大批量稳定运输。而 AI 负载更像早高峰的网约车:乘客随时上车、目的地随时改、路线随时绕。高铁再准时,也没法在这种不确定性里表现出色。
1.2 藏不住的"有状态":训练不是若干无状态函数的组合
传统分布式框架把计算模型化为"对一批数据进行变换",任务本身是无状态的,挂了就重跑。但一个模型训练任务最麻烦的地方恰恰在于全局状态:参数服务器的参数、优化器的动量项、样本的 epoch 索引,都是活的、巨大的、时刻在变的。你用无状态任务的视角去看参数同步和模型副本,就会发现根本映射不过来——你当然可以把梯度同步硬写成一个 MapReduce 步骤,但那种丑陋程度和性能损失,没人愿意在生产环境里碰。这也是为什么很多团队在 Spark 上做机器学习越做越痛苦:框架的正确用法和你实际要解决的任务形状不匹配。
1.3 异构资源感知:不是所有任务都吃 CPU
AI 负载里 GPU 是主力,数据清洗在 CPU 上跑,推理服务又要按请求量弹性伸缩。传统调度器把资源抽象成统一的 slot,对 GPU 的支持要么是后补的插件,要么干脆没有这概念。Ray 从底层就把 CPU、GPU、内存、对象存储当成一等公民来调度,一个任务标了 num_gpus=1,调度器就只把它放到有 GPU 的节点上。这种"专门的框架做专门的事"的错位,正是 Ray 出现的最基本面原因——它不是来抢 Spark 饭碗的,而是来补传统分布式范式够不着的那块 AI 负载。
2. Ray 的三块积木:Task、Actor、ObjectRef 是如何撑起统一 API 的
2.1 一句话理解三个原语
Ray 不搞新语言,不搞 DSL,直接建立在 Python 原生的函数和类之上。三个核心原语各有各的活,但组合起来能撑起几乎整个 AI 工程链路,这是它叫"统一 API"的根本底气。
ObjectRef 是分布式对象的"快递单号"。你把一个 Python 对象交给 Ray,它可能留在本机共享内存,也可能被搬到另外一台机器上,但你手里始终攥着单号,什么时候需要,拿着单号去取就行。拿单号的动作本身不阻塞,真正的等待发生在调用ray.get()的时候。Task 就是你扔出去给空闲 worker 跑的远程函数,@ray.remote装饰器一加,原来同步调用的函数变成异步提交:fn.remote(...)立刻返回一个 ObjectRef,后台的调度器会决定它在哪台机器跑、要不要走网络、内存够不够、要不要重试。Actor 则是常驻的、跨调用保状态的进程。你把一个类用@ray.remote包装,实例化之后,它就像一个坐拥固定工位的同事,工位上的档案和私人物品(也就是类属性)不会因为一次调用结束就被清走。多个 Task 可以排队向同一个 Actor 提交请求,Actor 在自己的进程里一个个消化。
2.2 为什么这三个原语就够覆盖整个分布式 AI 链路
因为你把真实的机器学习 pipeline 拆开看,需要的无非就四种形态。第一种是数据并行:一批样本做同样的变换,每个 Task 拿一部分数据跑,结果再聚回来。第二种是有状态的服务和训练器:模型参数、数据库连接、加载好的词表这些"重东西",只有 Actor 能挂在进程里反复复用,否则每个任务都重新反序列化一次,机器先炸给你看。第三种是异步控制和数据流转:训练循环里要不停看指标、决定下一步动作,ObjectRef 让你不用每次阻塞等结果,可以在等一个任务的同时提交另一个。第四种是资源隔离与弹性:给这个任务 2 个 CPU、那个 Actor 1 张 GPU,资源声明直接写在装饰器的参数里。
这三种原语加在一起,覆盖的是"分布式 Python 程序"这个完整的空间。反过来看,传统框架大多只覆盖其中一种:Spark 擅长数据并行,Celery 擅长任务分发,参数服务器得自己搭。Ray 的统一,是原语层面的统一,而不是把几个框架拼在一起的适配器。
| 原语 | 生活类比 | 核心特征 | 典型用途 |
|---|---|---|---|
| ObjectRef | 快递单号 | 异步获取、可传递引用、对象可复用 | 多任务数据流转、大对象共享 |
| Task | 跑腿任务 | 无状态、可重试、天然数据并行 | 预处理、批处理、超参搜索子任务 |
| Actor | 固定工位的同事 | 有状态、常驻、串行处理请求 | 模型服务、分布式训练 worker、参数持有者 |
2.3 所谓"统一 API",统一到了这里
明白这三块积木后,再看 Ray 的上层生态就不神秘了:Ray Data 是跑在 Task 上的数据处理,Ray Train 是跑在多个 Actor 上的分布式训练,Ray Serve 是常驻 Actor 组成的在线服务层,Ray Tune 是基于函数级调度的超参搜索。它们看起来是四个库,底层全复用同一套 Task/Actor/ObjectRef 机制。这也是我特别喜欢 Ray 的地方——你在 Data 里理解的内存模型,到 Train 里还是同一套直觉,不用每换一层框架就重塑认知。
3. 用同一套代码串起预处理、调参、推理:一个直接能跑的实战
3.1 场景设定
为了把"统一"落到地面上,我设计了一个很小的场景:有一批英文句子,先做分词和向量化,然后用不同的超参跑一个简单分类器,挑出准确率最高的配置,最后把它部署成一个能并发接收请求的 HTTP 服务。整个过程中,你只会用到 Ray 的同一套 API,数据在不同阶段之间的流动全靠 ObjectRef 完成。我的环境是 Python 3.9 + Ray 2.x,一行pip install ray[default]装完。
3.2 数据并行预处理:Task 的典型用法
数据集按行数切成 8 个分片,每个分片丢给一个远程函数处理。核心点在于,函数体里写的完全就是普通 Python,没有任何 DataFrame 或者 SQL 概念,这恰恰是 Ray 和传统大数据框架体感差异最大的地方。
import ray # 不传 address 时,本地单机自动初始化; # 在集群上,作业会连接当前环境里的已有集群 ray.init() from sklearn.feature_extraction.text import CountVectorizer # 假设这是加载好的已训练向量化器 vectorizer = CountVectorizer(vocabulary=loaded_vocab) vectorizer_ref = ray.put(vectorizer) # 放进对象存储,共享引用 @ray.remote(num_cpus=1) def preprocess_partition(lines, vectorizer_ref): vec = ray.get(vectorizer_ref) return [vec.transform([line]) for line in lines] # 把数据集切成 8 份 partitions = get_partitions(data_path, num_parts=8) refs = [preprocess_partition.remote(p, vectorizer_ref) for p in partitions] processed = ray.get(refs)这里有个非常重要的习惯:把大对象(比如这个向量化器)先用ray.put()放进对象存储,然后所有 Task 只传vectorizer_ref这个引用。如果不这么做,每次调用.remote()都会把整个对象从主进程序列化一遍传过去,数据一大就是性能灾难。记住这条,能帮你少走一半弯路。
3.3 并行超参搜索:Task 组合调优
预处理完就是调参环节。我把训练评估逻辑写成一个独立函数,每个配置是一个独立参数,100 个配置同时发出去。如果资源不够,Ray 会自动排队,不会直接失败。
@ray.remote(num_cpus=1) def train_and_evaluate(config, processed_ref): data = ray.get(processed_ref) accuracy = run_training(data, config) return config, accuracy # 生成了 100 组超参配置 configs = [generate_config(i) for i in range(100)] refs = [train_and_evaluate.remote(c, processed_ref) for c in configs] results = ray.get(refs) best_config = max(results, key=lambda item: item[1])[0] print("best config:", best_config)这段代码的价值在于:你完全不需要管理"哪台机器跑哪个任务",不需要自己写负载均衡,不需要考虑单点故障。装饰器一加,100 个 Task 会自动分散到集群里的空闲 CPU 上。如果你的训练逻辑需要 GPU,装饰器里改成@ray.remote(num_gpus=1),调度器就只会把任务放到有 GPU 的节点上;如果你不声明,哪怕集群有 50 张显卡,任务也只会安静地排在 CPU 上,然后向你报错。
3.4 Actor 里的模型副本:训练好的模型不该反复加载
调参完毕,要拿着最好的模型做批量预测。这里如果还用无状态 Task,每个请求或每个 batch 都重新加载一次模型文件,光是反序列化的时间就足以拖垮整个预测流程。正确做法是用 Actor 把模型挂在进程里,只加载一次。
@ray.remote class Predictor: def __init__(self, model_path, vectorizer_ref): # 构造阶段只加载一次,常驻内存 self.vectorizer = ray.get(vectorizer_ref) self.model = load_model(model_path) def predict(self, text): vec = self.vectorizer.transform([text]) return self.model.predict(vec)[0] predictor_actor = Predictor.remote("/data/model.pkl", vectorizer_ref) pred_refs = [predictor_actor.predict.remote(text) for text in sample_texts] labels = ray.get(pred_refs)看到没有?从无状态的 Task 切换到有状态的 Actor,唯一的区别是装饰的对象从函数变成了类。这个思维切换是理解 Ray 的关键一步,也是"范式重塑"里最巧妙的地方——不是发明新概念,而是把 Python 世界里再熟悉不过的东西原样保留,只是多给了它们并行和分布的能力。
3.5 模型服务上线:Ray Serve 把上面的 Actor 直接包成 HTTP 接口
模型选完、验证完,下一步是上线。Ray Serve 做的就是把上面这个 Actor 直接暴露成 HTTP 服务。以下写法基于 Ray 2.x 较新的接口,不同小版本 API 略有差异,但思路是一致的:
from ray import serve serve.start() @serve.deployment class ModelServer: def __init__(self, model_path, vectorizer_ref): self.predictor = Predictor.remote(model_path, vectorizer_ref) async def __call__(self, request): body = await request.json() text = body["text"] label = await self.predictor.predict.remote(text) return {"label": ray.get(label)} serve.run(ModelServer.bind(model_path, vectorizer_ref))注意看,这一层服务和上面的 Actor 跑的还是在同一个 Ray 集群上,用的是同一套对象存储。模型从训练到推理的传递,本质上就是把一个 ObjectRef 从一段逻辑传到另一段逻辑,中间没有数据落盘,没有文件传输,没有薛定谔式的环境不匹配。这就是"统一 API"在生产上真正值钱的地方——整套链路共用一个运行时,少了一整类集成问题。
4. Ray、Spark、Dask:技术选型对照,以及为什么不建议无脑替换
4.1 三个框架各自的"主场"
很多人一听到分布式计算就想到 Spark,一听到并行 Python 就想到 Dask,看到 Ray 冒出来,第一反应是"又一个框架"。我的建议是,别把它们当成同类产品比,你要先看清楚它们各自的假设是什么。
Spark 的假设是你的负载是一个预先定义好的、大规模稳定的批处理过程,它在 SQL 和 DataFrame 层面的优化器非常强,生态也非常成熟。Dask 的假设是你的代码本身是以 NumPy/Pandas 为中心的,你只需要用图式计算的方式扩大规模。而 Ray 的假设是:你的程序是动态的、有状态的、异步的,并且需要一套基础设施来处理从任务调度到分布式对象管理的完整问题。
4.2 用真实场景做选型决策
我整理了一张对比表,但这张表的作用不是帮你站队,而是帮你在做技术选型时快速对应自己的负载类型。
| 维度 | Spark | Dask | Ray |
|---|---|---|---|
| 调度粒度 | 粗粒度 Stage/批式 | 图级/惰性计算 | Task 级/细粒度函数调用 |
| 状态管理 | 无状态为主 | 无状态计算图 | Actor 常驻有状态 |
| GPU 感知 | 较弱,依赖外部插件 | 有限 | 原生 CPU/GPU/内存资源声明 |
| 编程接口 | DataFrame/SQL/RDD | Pandas/NumPy 风格 | 原生 Python 函数 + 装饰器 |
| 上层生态 | SQL/流处理/MLlib | 数组/DataFrame/分布式 ML | Data/Train/Tune/Serve |
| 最适合负载 | 离线 ETL、批分析 | 科学计算、中型 DataFrame | AI Pipeline、强化学习、在线服务 |
几个快速决策场景:每晚定时对几亿行日志做聚合分析,选 Spark,因为它稳定、成熟、优化器强,这类工作负载几十年没变过。你手里的代码就是一个大 Pandas 循环,只是数据大得单机放不下了,选 Dask,因为它和 Pandas/NumPy 的亲和度最高,改造成本最小。但如果你要跑训练循环,训练循环里有超参搜索、有动态剪枝、还要把训练好的模型直接部署成服务,那 Ray 是唯一能让你用一套心智模型从头撑到尾的选项。Uber、微软、OpenAI 这些公司早年都在这条链路上踩过坑,最后不约而同地靠向 Ray,不是没有原因的。
4.3 为什么我不建议"无脑替换"
技术选型的本质,是看哪个框架模型的隐含假设和你的实际负载最吻合。有些团队把 Spark 换成 Ray 之后反而更痛苦,因为他们的负载是规规矩矩的 SQL 聚合,Spark 优化器已经帮他们把执行计划安排得明明白白,而 Ray 把控制权全部交给你,也就意味着把优化责任也交给你了。如果你的负载根本不需要细粒度控制和有状态计算,换到 Ray 只会徒增自己调度的负担。这个判断比"谁更新谁先进"重要得多。
5. 从单机 demo 到集群落地:这些工程细节决定能不能跑起来
5.1 集群怎么起,节点怎么加
本地开发时ray.init()什么都不传就能跑,但正式集群不是这么玩的。Ray 集群至少有一个 head 节点负责调度和全局状态管理,其他节点都叫 worker,启动方式很简单:
# 在 head 节点上执行 ray start --head --port=6379 # 在每台 worker 上执行,<head_ip> 换成实际 IP ray start --address=<head_ip>:6379跑完ray status能看到当前集群里的节点数、资源总量和正在运行的任务。生产上更多用ray up拉起云上集群,或者直接用 K8s Operator 把 Ray 集群编排进 Kubernetes。这里提醒一句:客户端机器上的 Ray 版本必须和集群节点上的版本保持一致,这个小问题能浪费你一下午。
5.2 资源声明是要命的一环:你不说,调度器就瞎猜
我在群里见过最多的线上事故,就是有人装饰器上什么都不写,只写了一个@ray.remote。这样做的后果是:第一,任务默认只占 1 个 CPU,明明有 GPU 的代码因为没写num_gpus=1,调度器把它派到没有 GPU 的节点,然后一路运行到中途才报错;第二,嵌套并行时资源会被迅速占满,各个子任务互相争抢,整个集群吞吐暴跌;第三,Actor 不声明资源时同样按默认规则分配,你预期 10 个并发,结果调度器只放了一个,后面全部排队。
正确的姿势是显式声明:
@ray.remote(num_cpus=2) def heavy_task(data): ... @ray.remote(num_cpus=0.5) def io_task(data): # 适合 I/O 型任务,让多个任务共享一个 CPU ... @ray.remote(num_gpus=1) class GPUWorker: ...尤其注意num_cpus=0.5这种小数写法,它不是玩笑,而是让 I/O 密集任务合理共享 CPU 的常用手段。资源声明是 Ray 里最便宜也最有效的调优手段,一行参数能避免一晚上排查。
5.3 内存和对象存储:最容易被遗忘的容量规划
Ray 的每个节点都有一块共享内存对象存储,默认占用物理内存的 30%,你可以通过object_store_memory参数调整。很多人往ray.put()里放超大数据集,又不及时释放引用,一个任务跑完,集群直接 OOM 给脸色看。我的经验是记住三点:一是大对象能不进对象存储就不进,尽量让 worker 直接从 S3 或本地磁盘读;二是非核心的中间 ObjectRef 用完立刻del;三是时刻看一眼 Dashboard 的内存占用曲线,别等节点挂掉才追悔莫及。Dashboard 默认跑在 head 节点的 8265 端口,能清楚看到每个 Actor 的内存、每个 Task 的事件日志,这是排查问题时的第一入口。
5.4 容错设定:哪些能自动恢复,哪些不能
Ray 的容错不是一刀切的。无状态的 Task 挂了可以自动重试,默认会重试很多次,你可以用max_retries控制;但 Actor 默认不会自动重启,需要显式设置max_restarts。这就造成一个很有趣的工程权衡:如果你的 Actor 崩溃是因为代码 bug,调高max_restarts只会让问题在每次重启后再次暴露,徒增日志噪音;如果是偶发的 OOM 或者网络抖动,合理设置重试却能救回整个作业。我的建议是,对无状态任务明确设max_retries=3,对有状态任务设max_restarts=2的同时在旁边配套 checkpoint,让它重启后能恢复状态,而不是从零再来。这比默认配置可靠得多。
5.5 自动伸缩:别指望它即时生效
Ray 支持通过配置文件设置min_workers、max_workers,当任务队列增长时自动加机器。但自动伸缩是有冷启动时间的,少则一两分钟,多则更久。在线推理服务如果全指望弹性扩展,高峰期第一波请求一定被冷启动打垮。我的做法是保留一个合理的最小 worker 数,让弹性只用于消化突发流量,而不是承担所有伸缩期望。
6. 踩坑实录:几个我花了一整晚才弄明白的问题
6.1 ObjectRef 被垃圾回收,报 ObjectLostError
现象:Task 跑到一半,突然报ObjectLostError,而且那个对象明明刚创建没多久,按理说应该在内存里。
排查链路是这样的:我先怀疑是节点挂了,检查了一圈,节点都很健康;又怀疑是对象存储被清空,排查了容量配置,也正常。最后定位到,问题出在我把 ObjectRef 放在了一个循环作用域里,创建完引用之后既没有存到外部列表,也没有立刻 get,Python 的引用计数随即归零,Ray 底层回收了对象。等你后面再拿这个 ref 去取对象,自然是"物件已被处理掉"。
解决方式很朴素但必须执行到位:重要对象的引用一定放在主进程的长生命周期变量里,比如一个固定的列表或者字典;需要长时间保留的大对象,创建后立刻ray.put()一次,让自己手里永远握着单号。这一类问题最恶心的地方在于它不报错,只会在某个深夜让你的任务诡异失败。
6.2 ray.init 连不上集群,报版本不匹配
现象:照理说作业应该连到本地集群跑,结果报了一堆连接被拒、地址不对,甚至版本不匹配。
排查步骤我走了一遍:先看ray status,集群节点确实都在;再看脚本里的初始化代码,问题就出在这里——本地环境里跑过ray.init()的旧脚本,它会默认起一个全新的本地集群,而不是连接你已经启动的那个。正确做法是显式指定ray.init(address="auto"),让客户端自动发现当前环境里的已有集群。另一个常被忽略的是客户端和集群的 Ray 版本必须一致,当年我用pip install -U ray升级了本地包,但集群节点还是老版本,两边握手直接失败。解决并不难:所有节点统一重装同版本,然后逐个重启。这个问题没有技术难度,但极其消磨耐心。
6.3 Actor 反复崩溃重建,整个作业卡死
现象:一个 Actor 在生产环境里反反复复崩溃,日志里全是 restart 记录,作业永远完不成。
我的排查链路:先去看 Dashboard 里这个 Actor 的状态,发现它总是重建后几分钟又挂;查日志尾部,施工单位是构造函数里加载一个大模型,内存直接超过容器限额被系统杀掉。我一开始想通过调大max_restarts来硬扛,结果只是延长了痛苦,因为问题每次都出在同一个地方。
最后把模型从构造函数挪到了第一次调用时的懒加载,同时在启动参数里预留足够的内存配额,问题才真正解决。这里的关键教训是:Actor 的初始化阶段尽量轻量,凡是重活都挪到首次使用做懒加载。不要把网络请求、超大模型加载塞进__init__,那些东西一旦失败,重试成本高得离谱。
6.4 lambda 函数不能直接当远程函数用
这是新人最常见、也最隐蔽的问题。Ray 内部用 CloudPickle 做序列化,虽然比标准 pickle 强不少,但 lambda 闭包会携带大量隐式环境变量,有时候能序列化过去,有时候序列化过去之后在另一台机器上反序列化失败,报错信息还特别难追。更麻烦的是,不同环境的全局变量和导入状态不一致,闭包里的引用经常对不上。我的结论简单粗暴:所有远程函数都写成模块顶层的具名函数,参数全部通过 ObjectRef 传下去。牺牲一点点代码华丽度,换来的是跨环境稳定,这笔账怎么算都划算。
这些坑有一个共同点,它们都不是 Ray 的 bug,而是对"分布式运行时"这套心智模型理解不够到位时产生的预期偏差。你一旦接受了"对象引用是有生命周期的""远程函数是一次独立部署的执行单元"这些设定,再回头看这些问题,会发现全部理所当然。
最后聊点个人体会。Ray 不是银弹,它没有帮你解决算法问题,也没有让分布式变得和单机完全一样简单——它只是把分布式编程的心智负担压缩到 Task、Actor、ObjectRef 这三个原语里,其余全交给调度器和对象存储。我最近的项目把预处理、调参、推理三块逻辑放在同一套 Ray 集群上跑,中间数据几乎不需要落盘,不同阶段之间靠引用直接喂过去,这对效率的提升非常可观。如果你正在搭建 AI 工程链路,我建议先用手头的小数据把这三块积木跑通,再考虑搬上集群。一旦你在实战里完成了从"无状态 Task"到"有状态 Actor"的思维切换,再回头看传统那套"多个框架拼凑一条流水线"的搞法,真的会觉得绕了太多弯路。