Ray 新手性能四戒:延迟 ray.get、避免微任务、复用对象引用与流水线处理

Ray 新手性能四戒:延迟 ray.get、避免微任务、复用对象引用与流水线处理 Ray 新手性能四戒延迟 ray.get、避免微任务、复用对象引用与流水线处理【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray本篇文章围绕 Ray 官方文档《Tips for first-time users》即仓库中的 tips-for-first-time.rst展开面向第一次接触 Ray 的开发者系统讲解四个最容易犯、却会显著拖慢程序性能的常见错误过早调用ray.get()、切分过小的远程任务、重复向远程任务传同一大对象、以及等全部结果就绪再处理的同步式数据处理方式。读完本文你将掌握ray.init()、ray.remote、.remote()、ray.put()、ray.get()、ray.wait()六个核心 API 的正确使用节奏并能用一套可量化的方法评估自己的任务粒度是否合适。Ray 提供了高度灵活、却又极简易用的 API但恰恰是这种把函数加个装饰器就能并行的低门槛容易让新手写出看似并行、实则串行的代码。文中的四个 Tip 正是针对这些问题给出的实战纪律每个 Tip 都配有可直接运行的完整示例、真实耗时数据以及源码层面的解释。核心 API 速查表本文所有示例都只依赖以下六个核心 API先建立整体认知API说明ray.init()初始化 Ray 运行时上下文ray.remote函数/类装饰器将函数标记为在独立进程中执行的任务或将类标记为 Actor.remote()远程函数调用、远程类实例化、远程方法调用的统一后缀远程操作是异步的ray.put()将对象存入对象存储object store并返回其 ID同步操作ray.get()根据对象 ID或 ID 列表取回对象**同步阻塞**操作ray.wait()输入对象 ID 列表返回 (1) 已就绪的 ID 列表和 (2) 未就绪的 ID 列表默认每次返回一个就绪 ID从源码看ray.get、ray.put、ray.wait三个核心函数的 Python 层实现集中在 python/ray/_private/worker.py 中get位于 L2854–L3032put位于 L3037–L3084wait位于 L3093 之后。其中ray.get的文档字符串明确标注该方法会一直阻塞到对象在本地对象存储中可用为止若传入的是对象 ID 列表则保持输入顺序返回对应对象列表。这一点对理解 Tip 1 和 Tip 4 至关重要。文中所有性能数据均来自原文档在 13 英寸 MacBook Pro2.7 GHz Core i7、16GB 内存上的实测结果。为避免机器差异带来的波动示例统一使用ray.init(num_cpus4)指定 4 个 CPU由于每个任务默认申请 1 个 CPU这个设置允许最多 4 个任务并行执行系统构成为 1 个 driver驱动程序加最多 4 个 worker。Tip 1延迟调用 ray.get()把阻塞留在最后Ray 中所有远程操作任务、Actor 方法的调用都是异步的调用会立即返回一个 future本质上是一个结果 ID / ObjectRef这正是并行得以实现的关键——driver 可以在不等待结果的情况下批量发起多个操作。要拿到真正的结果必须对结果 ID 调用ray.get()而该调用会阻塞直到结果可用。副作用是阻塞期间 driver 也无法发起新的远程调用从而拖累并行度。遗憾的是新手顺手就写ray.get()太自然了。先看一个串行基线下面的代码调用 4 次do_some_work()每次耗时约 1 秒import ray import time def do_some_work(x): time.sleep(1) # 替换为你需要做的实际工作 return x start time.time() results [do_some_work(x) for x in range(4)] print(duration , time.time() - start) print(results , results)运行结果符合预期总耗时约 4 秒duration 4.0149290561676025 results [0, 1, 2, 3]常见错误一只加装饰器忘取结果很多新手并行化的第一步就是给函数加ray.remote并把调用改成.remote(...)import time import ray ray.init(num_cpus4) # 指定系统有 4 个 CPU ray.remote def do_some_work(x): time.sleep(1) return x start time.time() results [do_some_work.remote(x) for x in range(4)] print(duration , time.time() - start) print(results , results)执行结果让新手一头雾水duration 0.0003619194030761719 results [ObjectRef(df5a1a828c9685d3ffffffff0100000001000000), ObjectRef(cb230a572350ff44ffffffff0100000001000000), ...]两个关键信息其一程序瞬间完成不到 1ms——因为这里测的只是发起任务调用的时间不是任务运行时间其二拿到的是 4 个 ObjectRef 而不是[0, 1, 2, 3]——因为远程操作异步返回 future 而非结果本身。这也印证了源码中ray.get的语义ObjectRef 只是结果的标识必须显式取回。常见错误二每个任务立刻 get阻塞扼杀并行那就在调用后取结果呗于是改成results [ray.get(do_some_work.remote(x)) for x in range(4)]结果正确了但耗时依然 4 秒零加速duration 4.018050909042358 results [0, 1, 2, 3]原因正如前文所说ray.get()是阻塞调用在每次远程调用后立刻调用它等于等上一个任务跑完才发起下一个——本质上还是逐个执行并行度为零。正确姿势先全部提交再一次取回要让 4 个任务真正并行应当先批量发起所有远程调用再统一取结果results ray.get([do_some_work.remote(x) for x in range(4)])此时耗时降到约 1 秒说明 4 个do_some_work()确实在并行执行duration 1.0064549446105957 results [0, 1, 2, 3]小结ray.get()是阻塞操作过早调用会破坏并行度应尽量把ray.get()推迟到程序最后、把所有远程调用提交完毕之后再调用。关于循环内反复 get、不必要的 get等更细的对抗模式仓库中的 ray-get-loop.rst、unnecessary-ray-get.rst 和 ray-get-submission-order.rst 有更深入的讨论源码中ray.get的 docstring 也直接挂接了这些模式文档。Tip 2避免过小的任务用大任务摊销调度开销新手并行化的自然冲动是把每个函数、每个类都变成 remote但任务过小反而可能让 Ray 程序比串行 Python 更慢。再次以do_some_work为例这次把单次任务缩短到 0.1ms并把调用次数放大到 100,000 次import time def tiny_work(x): time.sleep(0.0001) # 替换为你需要做的实际工作 return x start time.time() results [tiny_work(x) for x in range(100000)] print(duration , time.time() - start)串行基线约 13.4 秒duration 13.36544418334961这与理论下限吻合10 万个 0.1ms 任务的下限是 10 秒再加上函数调用等开销13 秒在预期之内。现在用 Ray 并行化让每个tiny_work()调用都变成远程任务import time import ray ray.remote def tiny_work(x): time.sleep(0.0001) return x start time.time() result_ids [tiny_work.remote(x) for x in range(100000)] results ray.get(result_ids) print(duration , time.time() - start)结果出人意料——不仅没有加速反而更慢了duration 27.46447515487671原因在于每一次任务调用都有不可忽略的开销任务调度、进程间通信、系统状态更新等当任务本身只有 0.1ms 时调度开销完全支配了执行时间。解法聚合小任务摊销单次开销一种有效的提速方案是把远程任务做大让一次性调用摊销掉启动开销。下面用mega_work把每 1000 次tiny_work()聚合进一个更大的远程函数import time import ray def tiny_work(x): time.sleep(0.0001) return x ray.remote def mega_work(start, end): return [tiny_work(x) for x in range(start, end)] start time.time() result_ids [] [result_ids.append(mega_work.remote(x * 1000, (x 1) * 1000)) for x in range(100)] results ray.get(result_ids) print(duration , time.time() - start)运行耗时约 3.25 秒duration 3.2539820671081543大约是串行执行的 1/4与 4 个 CPU 并行执行的预期完全吻合。如何估算多大的任务才算够大自然的问题是任务到底多大才足以摊销远程调用开销一个实用方法是直接测量单任务调用开销。运行下面这个空任务基准ray.remote def no_work(x): return x start time.time() num_calls 1000 [ray.get(no_work.remote(x)) for x in range(num_calls)] print(per task overhead (ms) , (time.time() - start) * 1000 / num_calls)在 2018 款 MacBook Pro 上的实测结果为per task overhead (ms) 0.4739549160003662即执行一个空任务也要近 0.5ms。这意味着任务至少应耗时数毫秒才能摊销调用开销。当然单任务开销因机器而异、因本机任务与跨机远程任务而异但让任务至少运行几毫秒是开发 Ray 程序时一个非常实用的经验法则。这一主题的更多讨论可参考仓库中的 too-fine-grained-tasks.rst。Tip 3避免向远程任务重复传递同一大对象当把一个大对象作为参数传给远程函数时Ray 会在底层自动调用ray.put()把它存入本地对象存储。任务在本地执行时这会显著提升性能因为所有本地任务共享同一个对象存储。但某些场景下这种自动ray.put()反而成为性能瓶颈——典型例子就是重复传递同一个大对象import time import numpy as np import ray ray.remote def no_work(a): return start time.time() a np.zeros((5000, 5000)) result_ids [no_work.remote(a) for x in range(10)] results ray.get(result_ids) print(duration , time.time() - start)只调用 10 个什么都不做的远程任务却耗时约 1.08 秒duration 1.0837509632110596原因每次调用no_work(a)Ray 都会自动执行ray.put(a)把数组a复制进对象存储。a有 250 万个元素5000×5000复制开销不容小觑10 次调用就是 10 次全量复制。解法显式 ray.put 一次传递对象 ID避免重复复制的方法很简单显式调用一次ray.put(a)然后把a的 ID 传给no_work()import time import numpy as np import ray ray.init(num_cpus4) ray.remote def no_work(a): return start time.time() a_id ray.put(np.zeros((5000, 5000))) result_ids [no_work.remote(a_id) for x in range(10)] results ray.get(result_ids) print(duration , time.time() - start)耗时骤降至约 0.13 秒duration 0.132796049118042比原程序快了约 7 倍——因为复制数组a的操作从 10 次降为 1 次。从源码看ray.put在 python/ray/_private/worker.py 中通过worker.put_object()落盘到对象存储并返回 ObjectRefL3074–L3083且对象在仍有引用期间不会被逐出put的 docstring 明确说明The object may not be evicted while a reference to the returned ID exists。另一个隐藏收益防止对象存储过早填满相比提速避免同一对象的多次拷贝还有一个更重要的好处防止对象存储过早被填满、触发对象逐出eviction。频繁的逐出与重建会带来额外的传输和计算开销。仓库中的 pass-large-arg-by-value.rst 与 return-ray-put.rst 分别从大参数按值传递与返回值 put 后再传两个角度对对象存储与引用的正确用法做了更细致的剖析。Tip 4流水线处理数据谁先完成谁先处理如果对多个任务的结果统一调用ray.get()就必须等最慢的那个任务跑完才能开始处理。当各任务耗时差异很大时这会造成明显的等待浪费。考虑这样一个场景4 个do_some_work()并行执行每个任务耗时在 0~4 秒之间均匀随机分布随后由process_results()处理这些结果每个结果处理 1 秒。预期总耗时 最慢任务的耗时4 秒处理时间。import time import random import ray ray.remote def do_some_work(x): time.sleep(random.uniform(0, 4)) # 替换为你需要做的实际工作 return x def process_results(results): sum 0 for x in results: time.sleep(1) # 替换为实际处理代码 sum x return sum start time.time() data_list ray.get([do_some_work.remote(x) for x in range(4)]) sum process_results(data_list) print(duration , time.time() - start, \nresult , sum)实测接近 8 秒duration 7.82636022567749 result 6等最慢任务的同时其余任务早就完成却只能干等白白拉长了总时长。更优的做法是数据一就绪就立即处理——这正是ray.wait()的用武之地。不指定额外参数时ray.wait()会在参数列表中任意一个对象就绪时立即返回返回值为两个列表(1) 已就绪对象的 ID(2) 尚未就绪对象的 ID。把process_results()替换为每次只处理一个结果的process_incremental()import time import random import ray ray.remote def do_some_work(x): time.sleep(random.uniform(0, 4)) # 替换为你需要做的实际工作 return x def process_incremental(sum, result): time.sleep(1) # 替换为实际处理代码 return sum result start time.time() result_ids [do_some_work.remote(x) for x in range(4)] sum 0 while len(result_ids): done_id, result_ids ray.wait(result_ids) sum process_incremental(sum, ray.get(done_id[0])) print(duration , time.time() - start, \nresult , sum)总耗时降到约 4.85 秒提升显著duration 4.852453231811523 result 6两种执行方式的差异如下图所示图中 (a) 展示使用ray.get()等所有do_some_work()任务完成后才调用process_results()的时间线最慢任务约在时刻 4 结束随后 4 秒顺序处理总计约 8 秒(b) 展示使用ray.wait()的流水线时间线每个任务完成即触发process_incremental()处理与剩余任务的执行重叠总耗时降至约 4.8 秒。wait() 的关键参数num_returns 与 timeoutray.wait()的能力远不止每轮返回一个就绪 ID。从 python/ray/_private/worker.py 中wait的签名L3093–L3098可以看到两个实用参数num_returns默认 1指定本轮返回的就绪对象数量。例如ray.wait(result_ids, num_returns4)会等到 4 个结果全部就绪才返回等价于一次性的全量等待num_returns2则可每两个一批进行流水线处理。timeout默认 None最多等待的秒数。设置后函数在就绪数量达标或超时两者中先到者触发返回timeoutNone表示无限等待直到满足num_returns。注意timeout必须为非负数源码中对此有显式校验L3191–L3194。fetch_local默认 True为 True 时等待对象下载到本地节点后才算就绪为 False 时对象在集群任意位置可用即返回不触发向本地节点的拉取。此外wait的输入必须是 ObjectRef或 ObjectRefGenerator的列表传入单个引用会直接抛出TypeErrorL3175–L3182返回的两个列表均保持输入顺序L3126–L3129。循环配合ray.wait实现生产-消费式流水线是 Ray 中提升吞吐的通用范式仓库中的 pipelining.rst 专门讨论了如何通过提前请求下一项、再处理当前项来让计算与 RPC 传输重叠从而压满 CPU该文档也指出 Ray Data 等官方库重度依赖这种流水线技术。总结四条性能纪律延迟ray.get()先批量提交所有远程调用再统一取结果阻塞调用越晚并行窗口越大。避免微任务任务至少运行几毫秒才能摊销约 0.5ms 级的单次调用开销过小的任务应聚合为更大的远程函数。复用对象引用同一大对象反复传参时先ray.put()一次再传 ID既省 7 倍时间也避免对象存储被反复拷贝填满。流水线处理任务耗时参差时用ray.wait()让处理紧跟就绪结果用num_returns/timeout控制批量与超时语义让计算与传输/等待重叠。如果你希望进一步深入本仓库的 ray-core 模式库 提供了二十余篇针对具体反模式的专题文章如不必要的 get、get 循环、过细任务、大参数传值、流水线等是这四个 Tip 的进阶版本核心 API 的权威语义则以 python/ray/_private/worker.py 中的get、put、wait实现与 docstring 为准。【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考