尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Ray 嵌套任务(Nested Tasks)指南:用远程任务实现递归式嵌套并行

Ray 嵌套任务(Nested Tasks)指南:用远程任务实现递归式嵌套并行 人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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 中远程任务remote task不仅可以被 Driver 直接调用还可以在任务内部动态地提交其他远程任务甚至提交自身这种任务调任务的递归结构被称为嵌套任务nested tasks是表达分治divide-and-conquer类嵌套并行度的核心模式。本文将基于 Ray 官方设计模式文档从快速排序quick sort的分布式实现出发完整讲解嵌套任务的使用方法、ray.get()的正确调用时机、嵌套带来的调度开销与成本控制并结合本仓库源码给出可复现的完整代码与基准测试。读完本文你将掌握如何用嵌套任务写出可扩展的分治算法并理解何时应该避免过度细粒度的并行化。嵌套任务让远程函数动态派生新的远程函数本模式出自 Ray Core 官方设计模式目录设计模式与反模式目录原文见 nested-tasks.rst。Ray 的核心抽象是ray.remote装饰的远程函数与 Actor。一个远程函数在被调用时Ray 会将其封装为任务task提交给调度器由集群中的某个 worker 进程执行。而嵌套任务模式的核心思想是远程函数的执行体内部可以继续调用其他远程函数包括调用它自己从而在运行期动态地展开一棵任务树实现嵌套并行。这种模式的价值在于子任务之间彼此独立、可以并行执行特别适合分治类算法——父任务把问题切分成子问题后并不需要自己串行地完成全部计算而是把子问题再次提交为远程任务让 Ray 调度到其他空闲 worker 上并行处理。从实现层面看嵌套任务之所以天然成立是因为 Ray 的任务提交链路并不区分谁在提交无论是 Driver 还是 worker 进程都会通过 CoreWorker 的提交接口向 GCS 与调度器发送任务说明TaskSpec然后由调度器决定在哪个节点、哪个 worker 上执行。本仓库中任务提交与执行的底层逻辑集中在 src/ray/core_worker/core_worker.cc如SubmitTask/ExecuteTask相关路径而 Python 侧的ray.remote装饰器与任务封装实现在 python/ray/remote_function.py 与 python/ray/_private/worker.py 中。因此worker 内部再次调用fn.remote(...)与 Driver 调用完全走同一条调度链路只是调度的发起者变成了 worker。典型应用场景分布式快速排序官方文档给出的场景非常直观你需要对一个很大的数字列表做快速排序。快排天然是分治的——选定 pivot 后小于和大于 pivot 的两个子数组可以独立排序。如果整个排序都在一个进程里串行完成耗时与数据规模线性相关而利用嵌套任务父任务把两个子数组分别派发成新的远程任务Ray 就可以让它们在不同 worker 上并行执行从而在分布式环境下显著缩短排序时间。下图为嵌套任务执行时形成的任务树tree of tasks结构根任务不断派生子任务每个节点对应一次远程调用叶子节点完成实际排序后沿树向上归并结果。完整代码示例递归快速排序的分布式版本官方文档对应的可运行示例位于 pattern_nested_tasks.py以下为该示例的完整实现与文档收录内容一致import ray import time from numpy import random def partition(collection): # Use the last element as the pivot pivot collection.pop() greater, lesser [], [] for element in collection: if element pivot: greater.append(element) else: lesser.append(element) return lesser, pivot, greater def quick_sort(collection): if len(collection) 200000: # magic number return sorted(collection) else: lesser, pivot, greater partition(collection) lesser quick_sort(lesser) greater quick_sort(greater) return lesser [pivot] greater ray.remote def quick_sort_distributed(collection): # Tiny tasks are an antipattern. # Thus, in our example we have a magic number to # toggle when distributed recursion should be used vs # when the sorting should be done in place. The rule # of thumb is that the duration of an individual task # should be at least 1 second. if len(collection) 200000: # magic number return sorted(collection) else: lesser, pivot, greater partition(collection) lesser quick_sort_distributed.remote(lesser) greater quick_sort_distributed.remote(greater) return ray.get(lesser) [pivot] ray.get(greater) for size in [200000, 4000000, 8000000]: print(fArray size: {size}) unsorted random.randint(1000000, size(size)).tolist() s time.time() quick_sort(unsorted) print(fSequential execution: {(time.time() - s):.3f}) s time.time() ray.get(quick_sort_distributed.remote(unsorted)) print(fDistributed execution: {(time.time() - s):.3f}) print(-- * 10)代码要点拆解串行版本quick_sort标准的递归快排当列表长度不超过阈值20 万时直接用内置sorted收尾避免过深的递归。分布式版本quick_sort_distributed用ray.remote装饰。当列表长度大于阈值时先partition切分然后同时发起两个远程调用lesser quick_sort_distributed.remote(lesser)greater quick_sort_distributed.remote(greater)这两行是嵌套任务的核心父任务并不阻塞地等待第一个子任务完成后再发起第二个而是先把两个子任务都提交出去让它们并行运行。之后才用ray.get(lesser)与ray.get(greater)分别取回结果并拼接。魔术数字magic number阈值200000决定了何时停止继续派发子任务、改为本地排序。这个阈值正是为了避免任务过小带来的反模式详见下文成本与反模式一节。官方代码注释给出的经验法则是单个任务执行时长至少约 1 秒才值得走分布式递归。返回值拼接ray.get(lesser) [pivot] ray.get(greater)沿递归路径向上归并最终返回完整有序列表。嵌套任务可以像普通函数调用一样直接组合返回值这也是该模式易用性的体现。官方基准测试结果示例代码末尾自带的基准测试在示例运行环境的输出清晰地展示了嵌套任务的适用边界数组大小串行执行秒分布式执行秒结论200,0000.0400.152任务太小分布式反而更慢4,000,0006.1615.779分布式略占优势8,000,00015.45911.282分布式明显更快可以看到当任务很小时非分布式版本更快随着列表规模增大、单个任务执行时间变长分布式版本的优势才体现出来。这一现象与嵌套任务的开销模型直接相关下文详细展开。ray.get() 的调用时机先全部提交再统一取结果官方文档特别强调了一个关键点我们在两次quick_sort_distributed远程调用都发起之后才调用ray.get()lesser quick_sort_distributed.remote(lesser) greater quick_sort_distributed.remote(greater) return ray.get(lesser) [pivot] ray.get(greater)这样安排的目的是最大化工作负载中的并行度。与之对应的反模式是在循环中逐个调用ray.get()详见 ray-get-loop.rstray.get()是阻塞调用会一直等到对应的 ObjectRef 结果就绪才返回。如果在一个循环里提交一个任务 → 立即ray.get()等它完成 → 再提交下一个任务那么所有远程任务实际上被串行化了并行度为零。正确做法是先把所有远程调用都提交出去它们会在后台并行执行再统一获取结果。ray.get()也支持传入 ObjectRef 列表一次性等待全部任务完成。在本例中两个ray.get()紧挨着出现在返回语句里虽然它们是先后执行的但此时两个子任务早已在后台并行运行因此ray.get(lesser)的等待时间不会拖累greater任务的执行。如果改成先ray.get(lesser)再提交greater就会损失掉一半的并行度。另外与嵌套任务配套的还有一条反模式提醒尽可能把 ObjectRef 直接作为任务参数传递避免在任务内部对参数调用ray.get()详见 nested-ray-get.rst。Ray 会自动解析任务参数中的 ObjectRef 依赖等待其就绪后再执行任务而如果在任务内部先对 ObjectRef 调ray.get()会让 worker 进程保持占用并阻塞在 CPU 资源全部被占用时甚至可能造成死锁。本文快排示例中的ray.get()是合法用例——它需要按顺序拼接两个子结果属于嵌套场景下允许的必要获取。嵌套任务的三类固有开销官方文档明确提醒嵌套任务并非免费它至少带来三类额外成本额外的 worker 进程每个被派发的远程任务都需要占用一个 worker 进程来执行。任务树展开得越深、分叉越多同时存活的 worker 进程数就越多。如果子任务本身执行得极快进程创建/复用的开销可能远超计算本身。调度开销每个嵌套子任务都要经历一次完整的提交 → GCS 记录 → 调度器分配 → 目标节点执行流程底层对应 src/ray/core_worker/core_worker.cc 中的任务提交与执行链路这与普通函数调用的开销不在一个量级。簿记bookkeeping开销每个任务都涉及 ObjectRef 的创建、引用计数、结果回传与内存管理任务数量越多这类元数据开销越大。因此要用嵌套并行获得真实加速必须保证每个嵌套任务都做了足够有意义的工作。这与官方反模式文档 too-fine-grained-tasks.rst 的结论完全一致并行化/分布化任务的固定开销通常高于普通函数调用。如果并行化一个执行极快的函数开销本身可能比函数执行还久快排示例中的魔术数字阈值正是该原则的工程化体现只有当列表还足够大单个排序任务预计执行时间够长时才继续递归派发否则就地用sorted()完成避免为纳秒级任务付出毫秒级调度的倒挂。何时使用嵌套任务适用性判断综合官方文档与示例数据可以总结出嵌套任务模式的适用条件子任务天然可并行算法本身具有可分治结构快排、归并、矩阵分块、递归搜索等切分出的子问题互不依赖。单任务粒度足够大每个嵌套任务的实际工作量足以覆盖调度与簿记开销。官方经验阈值是单任务至少约 1 秒量级示例中的200000元素阈值在不同硬件上可以按此原则调整。深度可控递归深度过大会导致任务树节点数爆炸式增长。示例中通过阈值提前收尾到本地排序将任务总数限制在合理范围。并行收益可观测在数据规模较小时分布式版本反而更慢见上文 200,000 元素的基准0.040s vs 0.152s应在规模足够大时才启用嵌套并行。如果发现任务被切分得过于细碎官方建议使用batching批处理技术把多个小工作合并成一次远程调用让每个任务承载更多实际计算从而摊薄amortize固定开销详见 too-fine-grained-tasks.rst 中对应的反模式与批处理示例。总结嵌套任务是 Ray 表达分治型嵌套并行度的标准模式远程函数内部再次调用远程函数包括递归调用自身动态展开一棵可并行的任务树。使用时需要把握三个要点先提交、后获取把子任务的remote()调用全部发出后再统一ray.get()避免阻塞式串行化参考 ray-get-loop.rst。控制任务粒度用阈值或批处理保证单任务工作量足够大摊薄 worker 进程、调度与簿记三类固定开销参考 too-fine-grained-tasks.rst。优先传 ObjectRef 而非在任务内 get让 Ray 自动解析参数依赖避免 worker 阻塞甚至死锁参考 nested-ray-get.rst。完整的可运行示例见 pattern_nested_tasks.py运行前请确保已通过pip install ray安装 Ray 并import ray; ray.init()初始化运行时。赞分享人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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点击查看免费下载相关推荐mold 项目内嵌 oneTBB 指南使用嵌套流图Nested Flow Graphs构建分层并行计算mold 项目内嵌 oneTBB 指南使用嵌套流图Nested Flow Graphs构建分层并行计算 导读 本指南围绕 mold 仓库中集成的 oneT开发工具构建工具系统编程Web3-ui与Chakra UI深度集成扩展Web3功能的完整解决方案Web3 ui与Chakra UI深度集成扩展Web3功能的完整解决方案 Web3 ui是一个专为Web3开发打造的React UI库它与Chakra UIWouter高级路由模式递归路由与无限嵌套实现Wouter高级路由模式递归路由与无限嵌套实现 你是否曾在开发单页面应用SPA时遇到过需要构建多层级嵌套路由的场景比如电商平台的分类页面 /categ前端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表