 在 Ray 中构建可组合的分布式任务图)
Ray DAG 惰性计算图 API 全解用 .bind() 在 Ray 中构建可组合的分布式任务图【免费下载链接】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/rayRay 的ray.remote装饰器提供了运行时远程执行的灵活性而 Ray DAG API 则在此基础上更进一步通过.bind()将远程函数、类Actor与方法绑定为静态计算图的中间表示IR节点从而构建惰性求值的分布式计算图。本文将围绕 doc/source/ray-core/ray-dag.rst 的完整内容结合 python/ray/dag 目录下的源码实现系统讲解 Ray DAG 的节点类型、执行语义、输入输出抽象与 Actor 生命周期管理并给出可直接复制运行的完整示例。读完本文你将掌握用 Ray DAG API 编排多函数链、Actor 方法链、自定义输入与多输出 DAG 的完整能力并理解其与 Ray Compiled Graph 等上层高性能 API 的关系。Ray DAG 是什么面向开发者与库作者的图构建层从 ray.remote 到静态计算图ray.remote的优势在于运行时灵活性任务与 Actor 在调用.remote()时才真正被调度执行。而 Ray DAG 提供的是另一种编程模型——构建时静态声明。对于ray.remote装饰的类或函数你可以调用其.bind()方法生成一个中间表示IR节点这些节点充当 DAG 的骨架与构建块静态地把整个计算图钉在一起在执行阶段每个 IR 节点按拓扑顺序被解析为具体值。IR 节点还可以赋值给变量并作为其他节点的参数传入这使得开发者可以用普通 Python 变量组合出任意复杂的图结构a_ref func.bind(1, inc2) # a_ref 是 IR 节点 b_ref func.bind(a_ref, inc3) # 节点作为参数传给其他节点设计定位与推荐使用场景需要特别强调的是Ray DAG 是一个面向开发者developer facing的 API官方在文档中明确了两个推荐使用场景本地迭代与测试由更上层库编写的应用在 Ray DAG API 之上构建新的库。也就是说普通用户通常接触的是基于它构建的库例如 Ray Serve 的模型组合而 DAG API 本身是给库作者与需要精细控制图结构的开发者使用的底层工具。源码中DAGNode、InputNode、MultiOutputNode等类均以DeveloperAPI标注也印证了这一定位见 python/ray/dag/input_node.py。此外Ray 还提供了构建在 Ray DAG API 之上的实验性高性能 API——Ray Compiled Graph尤其适合多 GPU 应用场景可参考 ray-compiled-graph.rst。用函数构建 DAG任务链与任意根节点执行对ray.remote装饰的函数调用.bind()生成的 IR 节点在 DAG 执行时会被当作一个Ray Task运行并解析为任务的输出。下面是 doc/source/ray-core/doc_code/ray-dag.py 中的完整示例import ray ray.init() ray.remote def func(src, inc1): return src inc a_ref func.bind(1, inc2) assert ray.get(a_ref.execute()) 3 # 1 2 3 b_ref func.bind(a_ref, inc3) assert ray.get(b_ref.execute()) 6 # (1 2) 3 6 c_ref func.bind(b_ref, inca_ref) assert ray.get(c_ref.execute()) 9 # ((1 2) 3) (1 2) 9这段代码展示了三个关键特性链式组合b_ref以a_ref为位置参数c_ref以b_ref为位置参数、a_ref为关键字参数形成一条依赖链任意节点可作根任何 IR 节点都可以直接调用dag_node.execute()作为 DAG 的根节点从根节点不可达的其他节点会被忽略。例如执行b_ref.execute()时c_ref不会运行普通值与节点混用func.bind(1, inc2)中1是普通 Python 值a_ref是 IR 节点两者可以自由混合作为参数。底层实现FunctionNode 如何执行从源码看函数节点对应FunctionNode类python/ray/dag/function_node.py。它的_execute_impl在运行时把绑定的函数体、参数与 options 重新组装成一个普通的 Ray 远程调用return ( ray.remote(self._body) .options(**self._bound_options) .remote(*self._bound_args, **self._bound_kwargs) )这意味着 DAG 执行阶段本质上仍是标准 Ray Task 调度节点之间通过参数依赖自然形成任务间的数据依赖。测试 python/ray/dag/tests/test_function_dag.py 中构建了一个d(d2_ref, d_ref)的多层嵌套 DAG最终得到 28验证了节点复用时同一个节点作为多个节点的输入的执行正确性。为节点配置资源与行为选项.bind()之前还可以通过.options()为每个节点单独指定 Ray 资源与行为选项例如任务名、num_cpus、num_returns、max_retries等。测试 test_function_dag.py 中展示了完整用法b_ref b.options(nameb, num_returns1).bind(a_ref) c_ref c.options(namec, max_retries3).bind(a_ref) dag d.options(named, num_cpus2).bind(b_ref, c_ref)这些选项会被存入节点的_bound_options并在执行时透传给ray.remote(...).options(...)对应 function_node.py 的实现。若传入非法选项如负的num_cpus会在执行期抛出预期的ValueError。用类与类方法构建 DAGActor 编排对ray.remote装饰的类调用.bind()生成的 IR 节点在 DAG 执行时会被当作一个Ray Actor实例化每次执行 DAG 时都会重新实例化该 Actor。类方法调用则形成与该父 Actor 实例绑定的调用链。函数、类与类方法产生的 IR 节点可以自由混用组合成一个 DAG。以下示例同样来自 ray-dag.pyimport ray ray.init() ray.remote class Actor: def __init__(self, init_value): self.i init_value def inc(self, x): self.i x def get(self): return self.i a1 Actor.bind(10) # Instantiate Actor with init_value 10. val a1.get.bind() # ClassMethod that returns value from get() from # the actor created. assert ray.get(val.execute()) 10 ray.remote def combine(x, y): return x y a2 Actor.bind(10) # Instantiate another Actor with init_value 10. a1.inc.bind(2) # Call inc() on the actor created with increment of 2. a1.inc.bind(4) # Call inc() on the actor created with increment of 4. a2.inc.bind(6) # Call inc() on the actor created with increment of 6. # Combine outputs from a1.get() and a2.get() dag combine.bind(a1.get.bind(), a2.get.bind()) # a1 a2 inc(2) inc(4) inc(6) # 10 (10 ( 2 4 6)) 32 assert ray.get(dag.execute()) 32需要留意的执行细节a1.inc.bind(2)与a1.inc.bind(4)是在声明图结构不会立即执行而是作为同一 Actor 实例上按绑定顺序排列的方法调用链DAG 执行时inc调用按声明顺序2、4、6串行作用于各自 Actor 的状态最终根节点是combine它同时以a1.get与a2.get两个类方法节点为参数实现了函数节点 类方法节点的混合 DAG计算过程为a1.get返回10 2 4 16a2.get返回10 6 16combine求和得 32。底层实现ClassNode 与 ClassMethodNode类与方法分别对应ClassNode与ClassMethodNodepython/ray/dag/class_node.pyClassNode._execute_impl与函数节点类似把ray.remote(cls).options(...).remote(*args, **kwargs)作为执行动作即Actor 创建对ClassNode访问方法名如a1.get会返回一个未绑定的类方法节点_UnboundClassMethodNode其.bind()方法生成ClassMethodNode并记录父类节点PARENT_CLASS_NODE_KEY与上一次方法调用PREV_CLASS_METHOD_CALL_KEY后者用于保证同一 Actor 上方法调用的确定性提交与执行顺序ClassMethodNode._execute_impl在执行时通过getattr(self._parent_class_node, self._method_name)取得 Actor 句柄上的远程方法并调用从而把方法调用路由到正确的 Actor 实例。另外注意ClassNode的构造函数会拒绝把InputNode作为其参数传入class_node.py因为动态用户输入在类构造/绑定阶段尚不存在。用 InputNode 定义 DAG 的动态输入InputNode是一个 DAG 的单例节点代表运行时的用户输入值。它的使用规范是在无参数的上下文管理器with InputNode()中使用并在dag_node.execute()调用时作为参数传入。示例ray-dag.pyimport ray ray.init() from ray.dag.input_node import InputNode ray.remote def a(user_input): return user_input * 2 ray.remote def b(user_input): return user_input 1 ray.remote def c(x, y): return x y with InputNode() as dag_input: a_ref a.bind(dag_input) b_ref b.bind(dag_input) dag c.bind(a_ref, b_ref) # a(2) b(2) c # (2 * 2) (2 1) assert ray.get(dag.execute(2)) 7 # a(3) b(3) c # (3 * 2) (3 1) assert ray.get(dag.execute(3)) 10这里的dag_input是单例输入节点同一 DAG 中所有节点共享它同一个输入会被广播给所有引用它的节点a_ref与b_ref都引用了它。同一个 DAG 可以携带不同输入反复执行而不必重新构建图——这正是惰性的核心价值图只构建一次输入每次执行时注入。输入访问的三种方式从 input_node.py 的类文档可知InputNode支持三种取值方式且可通过__getitem__索引/字典键与__getattr__对象属性访问输入内部的字段访问方式写法对应执行参数位置索引dag_input[0]、dag_input[1]dag.execute([1, 2])对象属性dag_input.x、dag_input.m1dag.execute(obj)obj含对应属性字典键dag_input[m1]dag.execute({m1: 1})底层实现上InputNode的__getitem__/__getattr__会按需创建InputAttributeNode子节点input_node.py若execute()收到多个位置参数与关键字参数它们会被包装进DAGInputData由InputAttributeNode在执行时按__getitem__/__getattr__语义解析出具体值input_node.py。此外InputNode构造函数不接受任何 args/kwargs传入会直接抛出ValueErrorinput_node.py。用 MultiOutputNode 处理多输出 DAG当一个 DAG 需要产生多个输出时使用MultiOutputNode作为 DAG 的输出节点。此时dag_node.execute()返回一个Ray ObjectRef 列表列表长度与传给MultiOutputNode的节点数一致。示例ray-dag.pyimport ray from ray.dag.input_node import InputNode from ray.dag.output_node import MultiOutputNode ray.remote def f(input): return input 1 with InputNode() as input_data: dag MultiOutputNode([f.bind(input_data[x]), f.bind(input_data[y])]) refs dag.execute({x: 1, y: 2}) assert ray.get(refs) [2, 3]这个示例同时演示了两个要点多输出MultiOutputNode接收一个节点列表list或tupleexecute()返回每个分支输出的 ObjectRef 列表输入字段访问input_data[x]与input_data[y]分别从同一份用户输入字典中取不同字段作为两个并行分支各自的输入。从源码看MultiOutputNode的_execute_impl直接返回其绑定的参数即子节点解析后的结果列表并在构造时校验输入必须为list或tuplepython/ray/dag/output_node.py。复用 Actorbind 与 remote 的生命周期差异默认情况下通过Actor.bind()创建的 Actor 属于 DAG 定义的一部分当 DAG 执行结束后Ray 会销毁这些 Actor。如果希望在多次 DAG 执行之间复用 Actor例如承载有状态的服务应改用Actor.remote()创建 Actor再在 DAG 中引用它的方法。示例ray-dag.pyimport ray from ray.dag.input_node import InputNode from ray.dag.output_node import MultiOutputNode ray.remote class Worker: def __init__(self): self.forwarded 0 def forward(self, input_data: int): self.forwarded 1 return input_data 1 def num_forwarded(self): return self.forwarded # Create an actor via remote API not bind API to avoid # killing actors when a DAG is finished. worker Worker.remote() with InputNode() as input_data: dag MultiOutputNode([worker.forward.bind(input_data)]) # Actors are reused. The DAG definition doesnt include # actor creation. assert ray.get(dag.execute(1)) [2] assert ray.get(dag.execute(2)) [3] assert ray.get(dag.execute(3)) [4] # You can still use other actor methods via remote API. assert ray.get(worker.num_forwarded.remote()) 3这里的关键区别worker Worker.remote()创建的 Actor不属于DAG 定义DAG 执行不会销毁它DAG 中只绑定该 Actor 的方法worker.forward.bind(...)因此多次execute()会复用一个 Actor 实例其内部状态self.forwarded得以跨执行累积执行 3 次后为 3该 Actor 仍可像普通 Ray Actor 一样通过.remote()调用其他方法worker.num_forwarded.remote()。这一模式在源码层面对应ClassMethodNode支持父节点为ActorHandle的情况——_get_actor_handle()会检查_parent_class_node是否为ray.actor.ActorHandleclass_node.py从而把方法调用直接路由到既有的 Actor 句柄上。深入原理DAGNode 的执行机制与扩展方向节点的统一抽象所有 IR 节点函数节点、类节点、类方法节点、InputNode、MultiOutputNode都继承自抽象基类DAGNodepython/ray/dag/dag_node.py。每个节点持有四类数据_bound_args/_bound_kwargs绑定的位置/关键字参数其中可以是普通值也可以是其他 IR 节点_bound_optionsray.remote风格的资源与行为选项_other_args_to_resolve子类特定的、需要序列化解析的附加信息如类方法节点的父类节点_upstream_nodes/_downstream_nodes由参数依赖自动推导出的上下游关系构成图结构。execute() 的拓扑执行DAGNode.execute()dag_node.py通过apply_recursive对 DAG 进行自底向上的递归遍历先解析子节点再用已解析的结果替换父节点参数最终在根节点返回结果。apply_recursive带缓存机制保证同一节点在图内被多处引用时只执行一次。需要说明的是源码中DAGNode.execute()已标注RayDeprecationWarning弃用警告未来版本可能被移除建议关注其替代 API如 Compiled Graph 的execute系列。与 Ray Compiled Graph 的关系DAGNode上还提供了experimental_compile()dag_node.py可将 DAG 编译为加速执行路径支持_submit_timeout、_buffer_size_bytes、enable_asyncio、GPU 通信重叠等参数并返回CompiledDAG。这就是官方文档所述构建在 Ray DAG API 之上的实验性高性能 API关于编译图的概念、用法与可视化profiling/visualization详见 doc/source/ray-core/compiled-graph 文档目录。测试与验证资源Ray DAG 在仓库中有完善的测试覆盖可作为学习与验证的参考函数 DAGpython/ray/dag/tests/test_function_dag.py含链式任务、options 透传、非法选项报错等类 DAGpython/ray/dag/tests/test_class_dag.py输入节点python/ray/dag/tests/test_input_node.py输出节点python/ray/dag/tests/test_output_node.py编译图与 GPU 相关python/ray/dag/tests/experimental/下的test_compiled_graphs.py、test_torch_tensor_dag.py、test_multi_node_dag.py等。更多学习资源Ray DAG 是 Ray 生态中多个上层组件的基础机制其中最典型的是Ray Serve 的模型组合Model Composition——Serve 的serve.deployment内部即基于 DAG 节点构建服务图可参阅仓库中的 doc/source/serve 文档目录。此外Ray Compiled Graph 的可视化profiling 与 DAG 图形化展示可参考 doc/source/ray-core/compiled-graph 相关文档。若想直接从代码入手建议通读 python/ray/dag 目录下的dag_node.py、function_node.py、class_node.py、input_node.py、output_node.py五个核心文件并运行 doc/source/ray-core/doc_code/ray-dag.py 中的各段示例注意示例需要先pip install ray并在本地或集群中ray.init()后运行。总结Ray DAG API 通过.bind()把ray.remote函数与类提升为可组合的图节点提供了函数任务链、Actor 方法链、动态输入InputNode、多输出MultiOutputNode与 Actor 生命周期控制等完整能力并作为 Ray Compiled Graph、Ray Serve 等上层框架的构建基石。它的核心设计哲学是构建时静态声明、执行时拓扑求解、输入可重复注入图结构定义一次用户输入每次执行时提供非常适合库作者封装复杂流水线以及开发者对多 Actor/多函数协作做细粒度的编排与本地迭代调试。【免费下载链接】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),仅供参考