资讯动态

DataFlow分布式加速实战:RayOrch让数据流水线在多GPU上并行飞起来

发布时间:2026/9/2 13:56:44 来源:尧图企业网站定制
DataFlow分布式加速实战RayOrch让数据流水线在多GPU上并行飞起来【免费下载链接】DataFlow基于大模型算子和工作流的高效文本大模型训练数据合成框架项目地址: https://gitcode.com/OpenDCAI/DataFlowDataFlow是一个基于大模型算子和工作流的文本训练数据合成框架而它的RayOrch模块正是把数据流水线在多 GPU 上并行加速的发动机。你只需把已有的算子用一行代码包装起来推理就能自动分发到多张 GPU——无需修改算子本身也无需改动 pipeline 的存储逻辑。本文带你从零上手 DataFlow 分布式加速看懂原理、跑通示例并附上真实基准数据。为什么需要多 GPU 并行 DataFlow 的算子多为行独立map-style的推理任务给一行数据打分、生成、过滤彼此互不依赖。这类任务天然适合数据并行——把 4096 行数据拆给 8 张卡每张卡处理一份理论耗时就是串行的 1/8。问题在于过去要实现这一点你得手动切分数据、管理多进程、合并结果代码量大且易错。RayOrch 的RayAcceleratedOperator把这套脏活全干了让你只写业务不写调度。核心思路透明数据并行。从 pipeline 的视角看它就是一个普通算子读DataFlowStorage、写结果内部才把 DataFrame 扇出到 N 个 Ray Actor每个 Actor 各持有一份独立算子含模型。核心原理一个包装器搞定数据并行整个加速能力由三个轻量文件组成逻辑非常清晰dataflow/rayorch/accelerated_op.pyRayAcceleratedOperator包装器负责把数据分片、分发到 Actor、回收结果。dataflow/rayorch/memory_storage.pyActor 内部用的内存存储InMemoryStorage避免文件系统 I/O。dataflow/rayorch/init.py统一对外出口from dataflow.rayorch import RayAcceleratedOperator即可使用。它的工作流可以概括为三步懒加载Actor 在首次run()时才创建pipeline 的compile()阶段不会触发模型加载避免编译期空耗资源。分片分发读取存储中的 DataFrame → 转成记录列表 → 按连续切分策略发给 N 个 Actor。透明回收每个 Actor 内部用InMemoryStorage包裹数据分片、调用原始算子的run()最后把结果拼回 DataFrame 写回存储。这种外层FileStorage 内层InMemoryStorage的分层设计让调用方完全无感知——你对存储的使用方式和原来一模一样。快速开始三步跑通多 GPU 加速 ⚡1. 安装依赖RayOrch 是可选依赖两种安装方式任选其一# 方式一在 DataFlow 目录下整体安装已含 rayorch 依赖 pip install -e . # 方式二仅单独安装 RayOrch pip install rayorch0.0.12. 包装任意算子以 Superfiltering 质量打分器为例包装逻辑极简from dataflow.rayorch import RayAcceleratedOperator from dataflow.operators.text_sft.eval.superfiltering_sample_evaluator import ( SuperfilteringSampleEvaluator, ) from dataflow.utils.storage import FileStorage # 包装算子4 个并行副本每个副本 1 张 GPU scorer RayAcceleratedOperator( SuperfilteringSampleEvaluator, replicas4, # 4 个并行 Actor num_gpus_per_replica1.0, # 每个 Actor 占 1 张 GPU ).op_cls_init(devicecuda, max_length512) # 原始 __init__ 参数两个关键参数决定了并行规模replicas并行副本数即同时处理数据的 Actor 数量。num_gpus_per_replica每个副本的 GPU 配额支持小数。例如0.25表示 4 个副本共享一张卡非常适合小模型。3. 像用原始算子一样运行storage FileStorage( first_entry_file_namedata/input.jsonl, cache_path./cache, file_name_prefixstep, cache_typejsonl, ) scorer.run( storagestorage.step(), input_instruction_keyinstruction, input_output_keyoutput, ) scorer.shutdown() # 释放 GPU 资源见下文说明注意run()的调用方式和原算子完全一致IDE 还能自动补全参数——因为包装器通过ParamSpec把内部算子的__init__/run签名透传了出来。与 Pipeline.compile 无缝集成如果你用的是 DataFlow 的标准 pipeline 工作流继承PipelineABC并实现forward()RayOrch 同样能无感接入。这里有一个贴心的细节调用PipelineABC.compile()时编译后的_compiled_forward会在每个 stage 结束后自动调用shutdown()自动释放该阶段占用的 Actor 与 GPU 资源无需你手动管理。这段自动释放逻辑实现在 dataflow/pipeline/Pipeline.py 中if hasattr(op_node.op_obj, shutdown): self.logger.debug(fAuto-shutting down {op_node.op_name} to release actor resources.) op_node.op_obj.shutdown()什么时候必须手动 shutdown()场景是否必须调用单算子 / 脚本运行到结束可选进程退出时 Ray 自动清理多个RayAcceleratedOperator串联且都占 GPU必须否则后续 stage 会因 GPU 资源被预留而永久阻塞num_gpus_per_replica0纯 CPU可选CPU 充裕不阻塞但模型仍占内存多 GPU 并行实测加速比有多香 作者在 8× NVIDIA A800-SXM4-80GB 上用真实算子跑了串行 vs 并行的对比基准。以 Superfiltering 打分器基于 GPT-2124M处理tatsu-lab/alpaca数据集 4096 行为例配置Warm 耗时加速比Cold Start正确性串行1 GPU66.41s1.0×—基线2 GPU36.48s1.8×12.1s✓ 一致4 GPU18.45s3.6×13.0s✓ 一致8 GPU10.80s6.1×18.4s✓ 一致几点观察结果完全正确并行输出与串行逐行对比rtol1e-3全部一致加速是免费的不牺牲精度。加速比接近线性8 卡拿到 6.1× 加速。GPT-2 单行推理极快约 60 it/sRay 调度/序列化开销占比相对偏高所以没到理论 8×。换更重的模型如 7B 的 Deita时加速比会更接近线性。冷启动一次性的首次run()会创建 Actor 并加载模型约 12–20s后续复用 warm Actor耗时可忽略。你也能用仓库自带的 benchmark 脚本自己验证命令支持 argparse 传参# Superfiltering4096 行2/4/8 卡并行对比 python test/rayorch/test_real_operators.py --op superfiltering --rows 4096 --replicas 2 4 8 # Deita 7B256 行2/4 卡并行 python test/rayorch/test_real_operators.py --op deita --rows 256 --replicas 2 4 # 自定义 GPU 分配每副本 0.5 卡 python test/rayorch/test_real_operators.py --op superfiltering --rows 2048 --replicas 4 8 16 --gpus-per-replica 0.5脚本会自动包含1 GPU 串行基线并输出各副本数的耗时、加速比与正确性校验结果可保存为 JSON--save-json便于分析。完整参数说明见 test/rayorch/README-zh.md。新手避坑指南 只适合行独立算子。需要跨行全局状态的算子比如基于完整相似度矩阵的语义去重不要用这个包装器否则会破坏正确性。数据量越大加速越明显。数据量小时 Ray 通信开销占比高加速比打折推荐 ≥1024 行做基准测试。善用小数 GPU 配额。小模型如 124M没必要每副本独占一张卡num_gpus_per_replica0.25可以让 4 个副本共享一张卡硬件利用率更高。多 stage 串联务必释放。多个占 GPU 的加速算子串联时前一个 stage 的 Actor 若不清理会卡住后续 stage——手动场景记得shutdown()compile()场景则已自动处理。想先在 CPU 上试逻辑。仓库提供了 CPU-only 的 Dummy 算子测试 test/rayorch/test_accelerated_op.py约 30 秒即可跑通串行 vs 并行的完整对比无需 GPU。总结DataFlow 的 RayOrch 把多 GPU 数据并行这件复杂的事压缩成了一次包装零侵入不改算子、不改存储RayAcceleratedOperator一行接入。真加速8 卡实测最高 6.1× 加速且结果与串行完全一致。好集成与PipelineABC.compile()无缝配合自动管理 Actor 生命周期。弹性强replicasnum_gpus_per_replica两个参数灵活匹配任意硬件与模型规模。只要你的数据流水线是行独立的推理任务RayOrch 就是让它在多 GPU 上并行飞起来的最省力方案。从单卡到多卡你只需改一个数字。【免费下载链接】DataFlow基于大模型算子和工作流的高效文本大模型训练数据合成框架项目地址: https://gitcode.com/OpenDCAI/DataFlow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价