先说结论:Ray不是要替代谁,而是把“分布式”这件事下沉成了一个Python原生的运行时。它既没有像Spark那样把一切都抽象成RDD/DataFrame的批计算模型,也没有像Celery那样把任务做成一条投递队列就完事。它站在两者之间的空白地带——让普通的Python函数、类方法、状态对象,透明地跑在一个分布式集群上,同时保留了任务的动态依赖和对象引用的自动传递。
这篇文章我想用一次对比式拆解来聊透它。如果你已经用过Spark或者Celery,你会更容易理解Ray的设计动机;如果你只接触过其中一方,那这篇文章能帮你把“分布式计算到底在解决什么问题”这块拼图补完整。
1. Spark与Celery的能力边界:谁都能做分布式,但对“分布式”的理解不同
在进入Ray的细节之前,有必要先对齐一个基础认知:Spark和Celery都是分布式系统,但它们各自解决的是完全不同的“分布式问题”。很多人在做技术选型时踩坑,根源往往不是框架本身,而是把这两个问题搞混了。
1.1 Spark的分布式:对“数据”的分布式,而非对“应用”的分布式
Spark的核心抽象是RDD和DataFrame。它的分布式体现在:一份大规模数据集切分成多个partition,分布在多个executor进程里,每个executor只处理自己手里的那份数据。任务图(DAG)在Driver端生成,然后按照Stage提交到集群上执行。
这意味着Spark的应用形态被约束在了“批处理”这个范式里。你可以写Spark SQL、DataFrame算子、MLlib,但本质上你是在“描述一次大规模数据转换”,而不是在“运行一个业务应用程序”。它的调度模型是粗粒度的,一个Stage内的所有task执行同样的逻辑,遇到shuffle就同步阻断,等所有分区数据落地后再进入下一Stage。
这种设计的优势是容错模型简洁:通过RDD血缘可以重算丢失的分区。但代价是:Spark不擅长处理任务之间有复杂依赖、执行时间长短不一、需要动态改变执行路径的场景。比如你想在某个中间结果上根据运行时的数据分布动态决定是继续聚合还是提前终止,Spark的静态DAG模型表达起来非常拧巴,得借助累加器、外部存储状态等hack手段。
1.2 Celery的分布式:对“任务”的分布式,而非对“状态”的分布式
Celery的逻辑简单得多:把函数调用序列化成消息,投递到消息中间件(RabbitMQ/Redis),worker进程从队列里拉消息,执行完后把结果写进backend。它的分布式体现在任务可以被任意一个空闲worker消费,天然支持水平扩容和异步化。
但它有三个先天边界:
- 无状态优先:每个task默认是独立无状态的,任务之间如果要共享数据,得通过结果backend或者外部存储(Redis、数据库)中转。这是消息模型的产物——你不应该在消息里传递大块数据。
- 调度能力有限:Celery支持定时、路由、优先级,但任务间的依赖关系需要靠chain/group/chord等原语手动编排,依赖一多,代码就变成了嵌套回调地狱。
- 实时性不足:即使使用prefork池,任务执行结果从进程返回给调用方也要经过backend读写,很难做毫秒级的高频交互式计算。
Celery适合“把一条业务链路上的小任务分发出去”,比如发邮件、生成缩略图、触发第三方API回调。但它不是为“多个任务共享一个大对象,且要在这个对象上做多轮迭代计算”这种场景设计的。
1.3 中间地带的空缺:Ray看到的信号
Spark之外,有一类工作负载无法被批计算覆盖:需要细粒度任务调度、对象在任务间传递、多个有状态worker协同完成一个更大的业务逻辑——比如强化学习里的环境并行采样、超参数搜索里的模型并发训练、在线推荐系统的特征计算服务。
Celery之外,也有一个缺口:任务之间需要传递的不只是小对象,还有Azure Blob、视频帧、超大内存DataFrame这些重物件;任务之间需要动态生成新任务而不是提前铺好一张DAG;有状态actor需要从一个worker迁移到另一个worker,而不仅仅是一则消息流转。
Ray就是在这一层发力:用一个全局控制平面(GCS)管理集群状态,用raylet管理节点上的共享内存和worker进程,把Python对象直接放进分布式对象存储引用传递。这样在体验上接近写单机Python,在能力上覆盖了细粒度任务、动态依赖、有状态计算,同时又保留了横向扩展的天花板。
2. Ray架构核心拆解:GCS、对象存储与调度器背后的协同逻辑
如果只用一句话概括Ray的架构,我会说:它把“调度”和“数据”彻底解耦,用分布式共享内存把对象传递变成了幽灵引用,用中心化控制平面维护全局视图,但把真正的执行决策权下放给每个节点。
2.1 全局控制面(GCS):Ray为什么能容忍动态任务?答案在这里
Ray集群的基础节点包括一个Head节点和若干Worker节点。Head节点上跑着一个全局控制存储(GCS,Global Control Store),它本质上是一个Redis后端,保存整个集群的元数据:节点列表、对象目录、任务队列跟踪、actor的位置。
当你调用ray.get()阻塞等待一个ObjectRef时,本地raylet会向GCS查询该对象的位置,然后建立一条从拥有者节点到当前节点的数据通路。这个“查询+拉取”是lazy的——不是任务提交时就传输数据,而是真正需要这个值时才传输。这一条设计几乎决定了Ray能干所有Spark干不了的事情。
为什么动态任务在Spark里没法做?因为Spark的DAG是静态构建的,一个stage执行完,下一个stage的物理执行计划已经固定了。而Ray中,一个任务可以随时返回新的ObjectRef,调用方拿到这个ref后再提交新的函数调用,整个执行图是边运行边增长的。GCS负责维护这张动态图的一致性,raylet负责具体的数据落地和执行。
2.2 分布式对象存储:raylet里的Plasma与共享内存机制
每个raylet进程启动时会预分配一块大的共享内存区域(默认是可用内存的30%),这块区域就是该节点上的分布式对象存储。所有被放进Ray对象存储的对象,序列化后写入这块共享内存,本机其他进程可以通过零拷贝直接读取,跨节点则通过网络传输。
这个设计的收益是:同一节点上的多个worker处理同一个大对象时,不需要反复反序列化。比如你分发了一个500MB的机器学习数据集作为任务的输入参数,它序列化一次进本地共享内存,后续任何来自同一节点的tail任务都直接从共享内存零拷贝读取,IO成本几乎与线程间共享一致。
但这里有个很容易被忽略的坑:Ray ObjectRef的生命周期是引用计数的。当你调用ray.put()后拿到ref,这个对象会一直留在分布式对象存储里,直到所有引用都被清除。如果循环里反复put大对象不释放,内存会被吃干净。之前我见过一个跑数据预处理脚本的案例,循环2万次向对象存储写了中间DataFrame,内存直接爆掉,而单机Python根本不会暴露这个问题。
2.3 调度器与worker池:raylet的双重角色
一个raylet内部跑着两套核心组件:
- 调度器组件:接收来自driver或者下游actor的任务,检查本节点是否有空闲worker、所需对象是否本地可用,决定立即执行、等待传输还是转发给其他节点。
- 对象管理组件:维护共享内存存储、处理对象的引用计数和驱逐策略。
当一个Python函数被ray.remote装饰后调用,它的执行过程是这样的:
- 调用方本地序列化参数,生成TaskSpec并提交给本地raylet。
- 本地raylet查询GCS,看依赖对象在哪些节点上。
- 决策调度:优先放在输入对象所在的节点,避免网络传输。这就是Ray的locality-aware调度。
- 如果目标raylet上没有空闲worker,任务进入该raylet的待执行队列,等worker释放。
- worker执行完函数后,把返回值放进节点上的对象存储。
这套机制比Spark更细的地方在于:调度单位到单个Python函数的粒度,而且不做bulk同步。每个任务独立推进,数据局部性由调度器自动维护,用户不需要关心对象的物理分布。
3. 任务间的数据流:从ObjectRef到共享内存,再到跨节点传输
现在很多人一看到“分布式对象存储”这个词就下意识觉得重,但Ray的实现其实很轻。它的对象存储并不复制一份完整数据,而是序列化后写入共享内存段,并返回一个6字节的ObjectID作为引用。真正重的地方藏在序列化策略里。
3.1 序列化:为什么默认情况下你的类能直接跨节点传递
Ray自带的序列化器基于Apache Arrow和Cloud Pickle的结合。Arrow负责处理numpy数组、DataFrame等科学计算的数据格式,做到零拷贝反序列化;Cloud Pickle负责处理任意的Python对象。
如果你传给任务的参数是numpy array,Ray会走Arrow的零拷贝路径:直接把buffer指针传给worker,不需要反序列化。这比Spark的Java序列化高效得多——Spark需要把JVM对象转成字节流再通过网络传输,而Ray可以做到本机共享内存上零拷贝、跨节点只传Arrow的物理布局描述。
如果你自定义的Python类里嵌套了句柄对象、生成器、lambda闭包,Cloud Pickle会把它序列化成字节流。这里要提醒一句:尽量保证自定义类是“数据型”的,不要在构造函数里持有锁、文件句柄、socket连接。Ray确实能把它们序列化过去,但序列化后的对象在目标进程里其实已经断开了底层资源。
3.2 ray.put与ray.get的隐性代价
很多人把ray.put当作免费的午餐,稍微大点的对象就放进共享内存,然后到处传ref。问题在于,ray.put是同步的——它会阻塞当前线程直到对象写入完成。如果频繁put大数据,反而会成为性能瓶颈。
ray.get更是有隐藏的同步成本。当你wait或者get一个远端结果时,当前线程会阻塞,而且这个阻塞不是简单的睡眠,而是让出CPU给同进程的其他worker。如果你在主driver里同时get几千个ref,会导致驱动线程阻塞频繁切换。更合理的方式是:用ray.wait分批处理,或者把大任务切碎成阶段流水线,让任务在worker之间流动,而不是在driver上做集中收集。
一个真实的案例:我做过一个数据清洗项目,需要把1400万行日志按用户ID分组聚合,再做特征工程。初版代码直接for循环依次ray.get每个分片结果,跑一次要4分钟。改成.map批量提交,等全部结束后统一取ref,运行时间降到40秒。差别就是get的等待策略导致的。
3.3 零拷贝数据传输的“副作用”:嵌套对象引用的隐式依赖
Ray的ObjectRef还允许被嵌套进另一个对象的内部。比如你让任务A返回一个包含ObjectRef的列表,这个列表本身存储在对象存储里,但里面的ref仍然有效,调用方拿到列表后可以继续按需取子对象。
这个能力非常强大,可以构造“分布式惰性计算图”。但副作用是:引用链会让对象存储里的对象无法被立刻释放,必须等最外层引用结束。调优时要重点观察ray memory命令输出的引用树,防止出现父对象持有子对象导致整棵子树无法清理的情况。这个在我们线上环境踩过一次,一个吃内存的reference tree把节点内存撑到95%,排查后发现是一个任务把全部子任务的结果ref塞进了一个列表然后put到对象存储里保存,导致所有子任务的结果都无法被回收。
4. Actor编程模型:为什么有状态任务能一直留在岗位上
Ray的任务模型是无状态的,函数执行完就退出,所有局部变量随着进程销毁。但有大量场景需要状态常驻:模拟器环境、模型推理服务、数据库连接池、外部API客户端的会话保活。为此Ray提供了一套Actor原语。
4.1 从无状态函数到有状态对象:用类来锚定状态
把一个类用@ray.remote修饰后,任何一次类实例化都会在一个worker进程上创建一个长驻对象。该actor上的方法调用会自动序列化并发送到目标进程执行,本地只保存一个ActorHandle引用。
上面这个降级模式在真实项目中很常见:一个服务需要调用多个模型做融合打分,每个模型都放在单独actor里。第一次启动时会冷加载权重文件,几百MB的模型可能耗时十几秒,但后续每个请求都只做inference,延迟降到了毫秒级。如果用纯无状态任务,每次调用都要重新加载模型,根本没法用。
4.2 actor的并发模型:串行执行如何影响吞吐
默认情况下,一个actor内的所有方法调用是串行执行的——每个方法执行期间,其他方法调用会排队。这是为了让你无脑安全地操作实例变量,不用加锁。
但如果你的actor内部操作IO密集,串行执行会浪费吞吐。这时可以给类方法加@ray.method(num_returns=2)之类的注解,或者指定max_concurrency参数,让Ray创建多个线程在同一个actor实例上执行方法。注意:一旦开启并发,你就得自己负责线程安全,Ray不保证同一时刻只有一个线程修改实例变量。
实际项目中我倾向于这样设计:把无状态的计算逻辑拆成几个独立actor,每个actor只负责一种资源(比如一个模型、一个数据库连接),完全串行;需要并发时,客户端侧用线程池并发生成多个actor句柄,每个句柄一个独立实例。这样规避了所有并发安全问题的验证,性能也能打满。
4.3 actor的生命周期管理与容错
actor没有自动回收机制。当你不再使用某个actor且没有做显式的ray.kill(actor_handle),它会一直驻留在worker进程里直到集群结束。这在长周期服务里积累多了会成为内存黑洞。
Ray的容错策略是:如果actor所在节点宕机,actor状态彻底丢失,无法像任务那样从血缘重建。所以actor适合保存“可以快速重建的状态”,不适合保存“唯一的事实数据”。生产环境下建议为actor设计一个状态恢复钩子,比如从外部存储重新加载模型检查点、重放最近一段日志等。
5. 高级模式:嵌套任务、DAG、动态异步与Shuffle的真实用例
前面所有内容都在铺架构基础,这一节直接上实战。我不会泛泛而谈“Ray很强大”,而是把我在项目中真正用过的四个高级模式拆开讲。
5.1 工作流模式:嵌套任务与动态生成图
Ray最“反Spark直觉”的特点之一就是嵌套任务。一个函数内部可以继续提交新的Ray任务,然后ray.get它们的结果,再继续计算。这在逻辑上让分布式系统可以表达递归算法、树形搜索、自适应分而治之等模式。
比如经典的并行merge sort,在Ray里可以这样嵌套:
@ray.remote def parallel_sort(arr): if len(arr) <= 1: return arr mid = len(arr) // 2 left_future = parallel_sort.remote(arr[:mid]) right_future = parallel_sort.remote(arr[mid:]) left, right = ray.get([left_future, right_future]) return merge(left, right)这段代码在Spark里几乎写不出来——Spark的任务是预设DAG,不允许在executor内部再动态提交子任务。而Ray映射出来的是一个动态深度的任务树,每个内部节点产生两个子任务,执行结束再merge回本节点。
代价也是存在的:嵌套调度会在每个节点触碰一次raylet,层级一多调度延迟会累积。合理设计是人为限定嵌套深度,比如超过某一层后改成普通函数直接递归,避免纯ray.remote表达。
5.2 细粒度依赖:可以阻塞整个DAG,也可以只等待部分对象
ray.wait是Ray给依赖管理提供的另一个有力工具。它允许你给定一组ObjectRef,等待其中任意num_returns个完成即返回,未完成的部分可以后续继续wait。
这个API非常适合实现“袋鼠策略”:批量启动一堆任务,谁先回来谁先处理,新任务继续补充进池子,维持固定的并行度。相比Spark的barrier同步等待,这种模式在耗时方差很大的场景下能显著压缩整体时间。
举一个偏工程化的例子:你有一万个请求需要调用第三方接口,第三方接口延迟波动很大(有的10ms,有的5秒)。如果用延迟最低的策略,一次性提交所有请求,最后整体延迟被最慢的请求绑架;用ray.wait后每完成一个就立即处理结果,总耗时取决于请求的平均延迟而不是最大延迟。
5.3 Ray DAG与函数式图:Pipeline编程模型
Ray 2.x引入了一个轻量的DAG编程接口,通过with InputNode()构建任务流水线,然后用dag.experimental_compile()编译成可执行的计算图。它的作用与Spark的DAG类似,但表达的是函数级别的数据流。
它的优势在于:Ray的DAG不需要全部数据先落地再进入下一阶段,可以做到流式pipeline。比如一个包含预处理、模型推理、后处理的流程,可以在对象尚未完全生成时就启动后续阶段,对首个样本的响应时间优化明显。
5.4 数据并行:从map到shuffle的完整链路
Ray在数据并行上提供了ray.data.Dataset,它的算子类似Spark的RDD,但是底层内存模型是Arrow列式存储。map_batches、groupby、sort这些操作能跑在分布式shared-memory对象存储之上,避免Spark那样的大规模shuffle写磁盘。
但请务必记住:ray.data这层抽象比Spark DataFrame年轻得多,复杂SQL优化器、物化视图这些能力都还在补齐中。如果你的核心场景是复杂ETL,建议还是用Spark,把Ray放在模型训练、推理、实时计算通道这一层。两者不是互斥的,而是可以在同一个pipeline里共享Arrow格式的数据。
6. 横向对比:Ray vs Spark vs Celery,什么时候别用Ray
以下是基于我实际项目中的选型观察和线上运行数据做的对比:
- 无状态任务:三者都可以,Celery最轻量;Ray的
ray.remote也行,但需要一个Ray集群环境;Spark最重,只适合批处理。 - 有状态长驻服务:Ray最适合,Celery不行,Spark更难。
- 大对象跨worker传递:Ray有对象存储零拷贝优势,Celery需要外部存储中转,Spark依赖磁盘IO。
- 高容错离线批处理:Spark的checkpoint机制能精确到stage,Ray只有任务级的自动重试,复杂场景下恢复精度不如Spark。
- 毫秒级交互:Celery消息队列延迟高,Ray的进程内延迟在几十毫秒量级,可以用。
- 动态依赖/递归任务:Ray天生支持,Spark的静态DAG匹配困难,Celery需要硬编码chain。
从集群资源角度看,Ray的最小部署是一个Head节点加一个Worker节点,内存至少2GB,和Spark比已经算轻量。但和Celery的单worker内存占用(几十MB)相比,Ray也算重武器。做一个判断:如果只是异步邮件通知、简单后台任务队列,就别碰Ray;如果任务间开始传递复杂对象、有动态并行度、需要状态复用,Ray的性价比就显现出来了。
7. 生产环境落地与踩坑记录:三个值得警惕的细节
最后分享几个我们线上实际踩过后才明白的坑,这些在官方文档里不会写得很显眼。
7.1 内存别交给默认值,必须显式配置
Ray的共享内存默认占节点物理内存的30%,堆内存默认分配也很宽裕。如果你同时跑模型推理和数据处理,容易碰壁。生产环境我建议在集群初始化时就用ray.init(object_store_memory=..., _memory=...)把上限卡住,或者通过环境变量设置,宁可在调度时OOM清晰爆出来,也不要默默使用swap导致整个节点hang住。
7.2 task返回大对象时,小心重复复制
当一个task返回一个巨大的DataFrame时,Ray机制上是先写入worker进程堆内存,再复制进对象存储。这个复制过程中,同一份数据在节点上会有两份,内存峰值接近两倍。解析方法是在函数内部就用Arrow格式构造返回对象,减少从Python字节码到Arrow的转换损耗。
7.3 任务失败重试默认无限制,必须限制
@ray.remote(max_retries=...)如果不指定,默认会无限重试。对于幂等任务问题不大,但如果任务里有外部写入,比如写数据库、发消息、扣库存,无限重试会造成重复副作用。建议默认设1,关键业务设0,宁可失败让整个workflow显式fail,也不要做隐蔽的重复执行。
7.4 绘图一个大快乐:GCS不是万能的
GCS是单点。如果Head节点宕机,整个集群的任务调度会全部停止,driver与worker的连接会超时。Ray 2.x在GCS容错上做了不少工作,但远远没有达到K8s控制平面那套高可用标准。如果你跑的是长周期工作流,建议给Head节点挂一个自动重启机制,并在代码里设计checkpoint和恢复逻辑。没有这个兜底,一个36小时的训练作业在Head重启后只能从头跑。
写在最后
走完这一轮架构拆解和模式实战,我的感受是:Ray不是万能的银弹,也不是Spark的简单替代品。它把分布式计算从“数据平台”下沉到了“应用运行时”,让每个Python开发者都能在不改变思维模式的前提下写出分布式的代码。如果你正好处在Spark太重、Celery太弱的场景里,Ray大概率是目前最合适的选择。我自己在超大模型推理、强化学习环境并行、特征服务这三个项目上的经验都验证了这一点。希望这篇文章能帮你在选型和实际使用上少踩几个暗坑,多做几个正确的架构决策。