
如果你和我一样日常要同时处理数据清洗、模型训练和在线推理这三摊子事大概率会在某个时间点被同一件事逼疯每个环节都有一套自己的分布式框架。数据预处理要用 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 异构资源感知不是所有任务都吃 CPUAI 负载里 GPU 是主力数据清洗在 CPU 上跑推理服务又要按请求量弹性伸缩。传统调度器把资源抽象成统一的 slot对 GPU 的支持要么是后补的插件要么干脆没有这概念。Ray 从底层就把 CPU、GPU、内存、对象存储当成一等公民来调度一个任务标了 num_gpus1调度器就只把它放到有 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(vocabularyloaded_vocab) vectorizer_ref ray.put(vectorizer) # 放进对象存储共享引用 ray.remote(num_cpus1) 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_parts8) 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_cpus1) 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, keylambda item: item[1])[0] print(best config:, best_config)这段代码的价值在于你完全不需要管理哪台机器跑哪个任务不需要自己写负载均衡不需要考虑单点故障。装饰器一加100 个 Task 会自动分散到集群里的空闲 CPU 上。如果你的训练逻辑需要 GPU装饰器里改成ray.remote(num_gpus1)调度器就只会把任务放到有 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 用真实场景做选型决策我整理了一张对比表但这张表的作用不是帮你站队而是帮你在做技术选型时快速对应自己的负载类型。维度SparkDaskRay调度粒度粗粒度 Stage/批式图级/惰性计算Task 级/细粒度函数调用状态管理无状态为主无状态计算图Actor 常驻有状态GPU 感知较弱依赖外部插件有限原生 CPU/GPU/内存资源声明编程接口DataFrame/SQL/RDDPandas/NumPy 风格原生 Python 函数 装饰器上层生态SQL/流处理/MLlib数组/DataFrame/分布式 MLData/Train/Tune/Serve最适合负载离线 ETL、批分析科学计算、中型 DataFrameAI 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 --port6379 # 在每台 worker 上执行head_ip 换成实际 IP ray start --addresshead_ip:6379跑完ray status能看到当前集群里的节点数、资源总量和正在运行的任务。生产上更多用ray up拉起云上集群或者直接用 K8s Operator 把 Ray 集群编排进 Kubernetes。这里提醒一句客户端机器上的 Ray 版本必须和集群节点上的版本保持一致这个小问题能浪费你一下午。5.2 资源声明是要命的一环你不说调度器就瞎猜我在群里见过最多的线上事故就是有人装饰器上什么都不写只写了一个ray.remote。这样做的后果是第一任务默认只占 1 个 CPU明明有 GPU 的代码因为没写num_gpus1调度器把它派到没有 GPU 的节点然后一路运行到中途才报错第二嵌套并行时资源会被迅速占满各个子任务互相争抢整个集群吞吐暴跌第三Actor 不声明资源时同样按默认规则分配你预期 10 个并发结果调度器只放了一个后面全部排队。正确的姿势是显式声明ray.remote(num_cpus2) def heavy_task(data): ... ray.remote(num_cpus0.5) def io_task(data): # 适合 I/O 型任务让多个任务共享一个 CPU ... ray.remote(num_gpus1) class GPUWorker: ...尤其注意num_cpus0.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_retries3对有状态任务设max_restarts2的同时在旁边配套 checkpoint让它重启后能恢复状态而不是从零再来。这比默认配置可靠得多。5.5 自动伸缩别指望它即时生效Ray 支持通过配置文件设置min_workers、max_workers当任务队列增长时自动加机器。但自动伸缩是有冷启动时间的少则一两分钟多则更久。在线推理服务如果全指望弹性扩展高峰期第一波请求一定被冷启动打垮。我的做法是保留一个合理的最小 worker 数让弹性只用于消化突发流量而不是承担所有伸缩期望。6. 踩坑实录几个我花了一整晚才弄明白的问题6.1 ObjectRef 被垃圾回收报 ObjectLostError现象Task 跑到一半突然报ObjectLostError而且那个对象明明刚创建没多久按理说应该在内存里。排查链路是这样的我先怀疑是节点挂了检查了一圈节点都很健康又怀疑是对象存储被清空排查了容量配置也正常。最后定位到问题出在我把 ObjectRef 放在了一个循环作用域里创建完引用之后既没有存到外部列表也没有立刻 getPython 的引用计数随即归零Ray 底层回收了对象。等你后面再拿这个 ref 去取对象自然是物件已被处理掉。解决方式很朴素但必须执行到位重要对象的引用一定放在主进程的长生命周期变量里比如一个固定的列表或者字典需要长时间保留的大对象创建后立刻ray.put()一次让自己手里永远握着单号。这一类问题最恶心的地方在于它不报错只会在某个深夜让你的任务诡异失败。6.2 ray.init 连不上集群报版本不匹配现象照理说作业应该连到本地集群跑结果报了一堆连接被拒、地址不对甚至版本不匹配。排查步骤我走了一遍先看ray status集群节点确实都在再看脚本里的初始化代码问题就出在这里——本地环境里跑过ray.init()的旧脚本它会默认起一个全新的本地集群而不是连接你已经启动的那个。正确做法是显式指定ray.init(addressauto)让客户端自动发现当前环境里的已有集群。另一个常被忽略的是客户端和集群的 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的思维切换再回头看传统那套多个框架拼凑一条流水线的搞法真的会觉得绕了太多弯路。