ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

DORA tensor-pool 内存池传输示例实战:从 CPU↔CPU 到跨机 CUDA 的零拷贝张量搬运

DORA tensor-pool 内存池传输示例实战:从 CPU↔CPU 到跨机 CUDA 的零拷贝张量搬运 DORA tensor-pool 内存池传输示例实战从 CPU↔CPU 到跨机 CUDA 的零拷贝张量搬运【免费下载链接】doraDORA (Dataflow-Oriented Robotic Architecture) is middleware designed to streamline and simplify the creation of AI-based robotic applications. It offers low latency, composable, and distributed dataflow capabilities. Applications are modeled as directed graphs, also referred to as pipelines.项目地址: https://gitcode.com/GitHub_Trending/do/dora导读本文以 dora 仓库中libraries/extensions/tensor-pool/examples/目录的官方示例为核心系统讲解 DORA 固定内存池pinned memory-pool传输的使用方式包括正路径吞吐测试CPU↔CPU、CPU↔CUDA、CUDA↔CPU、同机多守护进程的零拷贝直读、跨机器镜像传输以及重复释放、释放后读写等生命周期负路径场景。读完本文你将掌握register_tensor_pool/write_tensor_pool/read_tensor_pool/free_tensor_pool四个操作的实际调用方式、各 YAML 的配置语义与运行命令并能基于仓库源码理解零拷贝、commit-acknowledged 写确认与帧序保证的底层原理。⚠️ 适用范围说明tensor-pool 是一个 opt-in 扩展不属于 dora 1.0 兼容性保证范围其 API、段布局与描述符格式可能在任意版本变化并存在已知未决缺陷。如果吞吐量对你比稳定性更重要请锁定 dora 版本后使用。详见 libraries/extensions/tensor-pool/README.md。1. 示例概览它在验证什么示例的核心目标是验证 DORA 的固定内存池传输sender 节点和 receiver 节点之间通过共享内存反复搬运张量而不是每帧都走普通消息路径。正路径场景cpu2cpu.yml、cpu2cuda.yml、cuda2cpu.yml等保持吞吐导向的行为生产者注册一次共享张量池之后每帧原地覆写接收者零拷贝读取并统计平均吞吐。负路径场景duplicate_free.yml、read_after_free.yml、write_after_free.yml、auto_cleanup.yml验证生命周期错误被以警告warning形式呈现而不会导致节点崩溃。跨守护进程/跨机器场景*_cross*.yml验证register_tensor_pool(machine...)在另一台 daemon同机或异机上镜像内存池同机多 daemon 时接收者直接读取发送者共享内存段零拷贝异机时走可靠的 zenoh/TCP 数据面。2. 安装与环境准备示例运行前需要 dora CLI 与 Python 依赖pip install dora-rs-cli # 如果尚未安装使用--uv运行时每个 YAML 的build:步骤会自动把 torchCPU 版来源download.pytorch.org/whl/cpu装进 per-node 托管环境无需预装。如果你不使用--uv请手动安装pip install torch numpy pyarrow tqdmCUDA 接收者场景cpu2cuda.yml、cuda2cpu.yml还需要确认 CUDA 可用python -c import torch; assert torch.cuda.is_available()3. 示例文件清单与职责划分文件说明sender.py发送侧注册并更新内存池读取cross_machine目标机器 id环境变量切换为跨机注册receiver.py接收侧从内存池读取、测量吞吐、按场景触发生命周期操作cpu2cpu.yml正路径吞吐测试CPU 发送 → CPU 接收无 GPU 的 CI 也可运行cpu2cuda.yml正路径吞吐测试CPU 发送 → CUDA 接收cuda2cpu.yml正路径吞吐测试CUDA 发送 → CPU 接收duplicate_free.yml接收者重复释放同一内存池CPU 接收者read_after_free.yml接收者先释放、再读取同一内存池CPU 接收者write_after_free.yml发送者先释放、再写入同一内存池CPU 接收者auto_cleanup.yml接收者不释放期望 daemon 在关闭时自动清理CPU 接收者cpu2cpu_cross.yml/cpu2cpu_cross_local.ymlCPU→CPU 跨机器 / 同机跨 daemon 吞吐测试cpu2cuda_cross.yml/cuda2cpu_cross.yml/cuda2cuda_cross.yml涉及 GPU 的跨机器吞吐测试GPU 侧会自动创建 CPU staging 池全部文件位于 libraries/extensions/tensor-pool/examples/发送与接收逻辑分别见 sender.py 与 receiver.py。3.1 核心环境变量所有 YAML 通过env:段注入环境变量来控制场景源码中由os.getenv读取默认值见 sender.py/receiver.py 顶部sender_device/receiver_device取值为cpu或cuda:idx决定张量所在的设备message_num帧数默认 100负路径场景降为 2 以快速完成生命周期校验memory_pool_scenariothroughput默认或duplicate_free/read_after_free/write_after_free/auto_cleanupcross_machine跨机场景下目标机器 id如B。4. 正路径吞吐场景单 daemon4.1 运行命令dora run libraries/extensions/tensor-pool/examples/cpu2cpu.yml dora run libraries/extensions/tensor-pool/examples/cpu2cuda.yml dora run libraries/extensions/tensor-pool/examples/cuda2cpu.yml4.2 预期行为dataflow 运行至完成sender 与 receiver 各打印一次预览张量Sender preview: .../Receiver preview: ...receiver 在结尾打印平均吞吐Average transfer throughput: xxx MB/s无崩溃、无明显内存错误。以 cpu2cpu.yml 为例其结构为两个节点sender_node通过inputs: next_require: receiver_node/next_require订阅接收者的回执receiver_node通过inputs: latency: sender_node/data订阅数据build:步骤负责安装 torchCPU 源与 numpy/tqdm。4.3 源码视角turn-based 握手如何保证帧序从 sender.py 可以看到数据流不是“发完就完”而是轮流制turn-based握手第 0 帧注册内存池后立即send_output(data, memory_pool_id, metadata)然后node.next()等待接收者回执后续帧write_tensor_pool原地覆写共享内存 →send_output发送空数组通知 →node.next()等待接收者的next_require确认。帧序保证来自接收者的 ack而不是发送节奏。接收者在 receiver.py 里每帧读完后send_output(next_require, pa.array([]))。发送者在写下一帧前必然已收到 ack因此接收者读取时发送者必然已完成写入——这就是零拷贝原地覆写receiver 的张量对象自动反映新数据能够成立的同步前提。单调计数校验发送者把迭代号写入random_data[0]第 31 行random_data[0] i接收者校验int(torch_tensor[0].item()) ireceiver.py 第 86、101 行。相比 “sum-of-8” 校验约 3% 冲突率单调计数可以确定性验证传播且无碰撞风险。帧序守卫与 WAN 重试跨机场景下镜像写入可能滞后于通知接收者对每一帧都做时限化300s而非计数化的重试time.monotonic() deadline循环内捕获read_tensor_pool的RuntimeError后 sleep 0.1s 重试直到读到期望帧。选择单调时钟是因为窗口内 NTP 回调会压缩/拉伸墙钟重试窗口同机直读只需毫秒级 ack第一遍重试即退出。5. 同机多 daemon自动检测的零拷贝直读5.1 拓扑与“防呆”设计cpu2cpu_cross_local.yml描述同一台机器上两个 daemonsender 挂在 daemon Areceiver 挂在 daemon B。无需任何额外配置内存池会自动检测两个 daemon 共享/dev/shm注册 ack 报告directtrue于是 receiver 直接读取 sender 的段——零拷贝完全绕过 daemon 中继。这是刻意的“防呆”设计同机多 daemon 部署永远不会静默回退到中继路径89–113 MB/s而是始终走直读约 5.8 GB/s。5.2 启动命令# coordinator 两个 daemon同机每个 daemon 需要各自的 --local-listen-port dora coordinator --port 6025 --store memory dora daemon --machine-id A --coordinator-addr 127.0.0.1 --coordinator-port 6025 \ --zenoh-peer tcp/127.0.0.1:5463 --local-listen-port 0 dora daemon --machine-id B --coordinator-addr 127.0.0.1 --coordinator-port 6025 \ --zenoh-peer tcp/127.0.0.1:5463 --local-listen-port 0 # 通过 coordinator 构建deploy 段要求走 build然后 attach 启动 dora build --coordinator-port 6025 libraries/extensions/tensor-pool/examples/cpu2cpu_cross_local.yml dora start --coordinator-port 6025 libraries/extensions/tensor-pool/examples/cpu2cpu_cross_local.yml --attach5.3 预期结果Average transfer throughput约为5800 MB/s同机直读同拓扑中继基线为 89–113 MB/s差距约50–65 倍。该场景由 torch 门控的冒烟测试smoke_local_memory_pool_cpu2cpu_cross_local覆盖并纳入memory-pool-smoke夜间 CI 任务见示例 README 的 Notes 部分。6. 跨机器场景镜像池 commit-acknowledged 写确认6.1 集群启动sender 在机器 Areceiver 在机器 B数据经两个 daemon 之间的 zenoh TCP 传输。每台机器一个 daemoncoordinator 可运行在任一台# 对外发布/作为 zenoh rendezvous 的机器 dora coordinator --interface 0.0.0.0 --port 6025 --store memory dora daemon --machine-id B --coordinator-addr 127.0.0.1 --coordinator-port 6025 \ --zenoh-peer tcp/0.0.0.0:5463 # 监听 5463供对端机器连接 # 另一台机器 dora daemon --machine-id A --coordinator-addr rendezvous-ip --coordinator-port 6025 \ --zenoh-peer tcp/rendezvous-ip:5463 # 拨号连接 rendezvous6.2 YAML 必需的三个要素所有*_cross*.yml都已内置这三项可对照 cpu2cpu_cross.ymlenv: cross_machine: B—— sender 以register_tensor_pool(machineB)注册缺少它池保持本地receiver 永远看不到镜像deploy: machine: A|B—— 每个节点由哪个 daemon 拉起deploy: working_dir: .相对 yml 所在目录—— 相对于 daemon 的 cwd仓库根目录遵循 multiple-daemons 约定绝对路径不可移植。6.3 真 WAN 的 zenoh 配置默认的多播发现无法跨越路由网络因此在拨号方 daemon 上通过ZENOH_CONFIG指定三点配置缺一不可// zenoh_wan.json5 { connect: { endpoints: [tcp/rendezvous-ip:5463] }, // 显式连接禁用多播 scouting: { multicast: { enabled: false } }, // 否则报 Scouting delay elapsed transport: { link: { tx: { queue: { congestion_control: { block: { wait_before_close: 60000000 } // 60s5s 默认值会在慢链突发时 } } } } }, // 杀掉会话 }ZENOH_CONFIG/path/to/zenoh_wan.json5 dora daemon --machine-id A --coordinator-addr ... --zenoh-peer tcp/ip:54636.4 运行无build:步骤的 YAML 可直接dora start——dora build会与快速构建竞争并可能报告 “no running build”这是无害的dora build --coordinator-addr rendezvous-ip --coordinator-port 6025 libraries/extensions/tensor-pool/examples/cpu2cpu_cross.yml dora start --coordinator-addr rendezvous-ip --coordinator-port 6025 libraries/extensions/tensor-pool/examples/cpu2cpu_cross.yml --attach6.5 跨机已知行为跨机池接受张量上限为1 GiB 注册上限。每帧写入只向 daemon 发送元数据共享内存引用daemon 读取发送者段并转发因此64 MiB 的 node→daemon 请求上限不适用写入是commit-acknowledgedwrite_tensor_pool只有在镜像 daemon 确认段写入后才返回因此后续的send_output通知永远不会超过数据本身receiver 不可能拿到过期帧。镜像写失败或 120s ack 超时会响亮地失败涉及 GPU 的跨机路径自动经过 CPU 池中转GPU_A → DtoH → CPU_A → zenoh TCP → CPU_B → HtoD → GPU_B原生 dora 的跨机中继只能承载小帧超过即挂起Drop express 在 16-batch TX 队列积压时会静默丢弃分片ROS 2 网络 DDS 在高延迟链路上受 RTT 节流。内存池是唯一能在 WAN 上搬运大帧的路径同机跨 daemon 读取完全绕过写入/ack 机制directtrue。6.6 跨机调试检查清单设置WALL_CLOCK: 1YAML env 中或WALL_CLOCK1跨机计时必须用墙钟——perf_counter的纪元是每台机器的启动时间跨机差值会被启动时间差主导主机已 NTP 同步墙钟差值才是真实传输时间receiver 曾用 perf_counter 测得约 0.00002 MB/s 的假象值见 sender.py 注释验证连接ss -tn | grep 5463显示拨号 daemon 到 rendezvous 的 ESTAB 连接coordinator WS6025必须被每个 daemon 可达拨号方处于 NAT 之后时镜像/对端 daemon 反复打印的 “Unable to connect to any locator of scouted peer” WARN 是无害的——数据路径由拨号方出站连接建立receiver preview sender preview字节一致张量是完整性检查过期镜像会表现为首元素不匹配测试脚本原生 dora 控制 harness、扫描脚本、会话日志位于/home/tcr/dora_test/可复用做重新测量外部目录仅作参考。7. 负路径场景生命周期错误不崩溃dora run libraries/extensions/tensor-pool/examples/duplicate_free.yml dora run libraries/extensions/tensor-pool/examples/read_after_free.yml dora run libraries/extensions/tensor-pool/examples/write_after_free.yml dora run libraries/extensions/tensor-pool/examples/auto_cleanup.yml7.1 预期警告/信息场景预期输出重复释放duplicate freeAttempt to release memory pool [memory_pool_id] failed - reason: pool does not exist. Operation aborted.释放后读read after freeAttempt to read memory pool [memory_pool_id] failed - reason: pool does not exist. Operation aborted.释放后写write after freeAttempt to write memory pool [memory_pool_id] failed - reason: pool does not exist. Operation aborted.自动清理auto cleanupDetected xx unreleased memory pool, releasing...与Successfully released xx unreleased memory pools!7.2 源码如何构造这些场景duplicate_free.ymlreceiver 在最后一帧连续调用两次free_tensor_pool见 receiver.py 第 113–115 行read_after_free.ymlreceiver 先释放、再尝试read_tensor_pool并用try/except Exception: pass吞掉预期异常第 116–121 行write_after_free.ymlsender 在第 1 帧写入前先free_tensor_poolsender.py 第 67–68 行此场景下 receiver 跳过单调计数断言SCENARIO ! write_after_free才校验receiver.py 第 99 行auto_cleanup.ymlreceiver 完全不释放由 daemon 在关闭时回收孤儿段——对应 daemon 侧的/dev/shm回收逻辑TensorPoolManager见 libraries/extensions/tensor-pool/src/lib.rs。四个负路径 YAML 使用message_num: 2减少消息数让生命周期校验短小聚焦它们与cpu2cpu.yml一样使用 CPU-only 接收者receiver_device: cpu对无 GPU 的 CI 运行器安全。8. 源码级原理四个操作与底层机制8.1 与 dora 的边界通用扩展通道tensor-pool只通过公开扩展通道与 dora 交互详见 docs/extensions.md。dora 本身对池、CUDA、段布局一无所知它只中介一个不透明描述符的生命周期本传输dora 侧对应发布池描述符extension_store(dora-tensor-pool, id, bytes)查找extension_load(dora-tensor-pool, id, remove…)撤销extension_drop(dora-tensor-pool, id)获知哪些 key 消失drain_dropped_extension_keys(dora-tensor-pool)调用 daemon 半部extension_request(dora-tensor-pool, bytes)触达对端 daemon 上的同名扩展InterDaemonEvent::ExtensionMessage由 daemon 半部发布描述符在 python/src/seam.rs 中是 JSON请求/对端消息在 src/protocol.rs 中是 postcarddora 两者都不解析。这个边界是刻意的它让扩展能自改元数据与自改协议而不触碰 dora 的线格式也让unsafe指针运算、seqlock、段布局与内嵌的libcudart绑定全部留在扩展这一侧。8.2 Python 侧的四个操作python/src/transport.rs 实现传输本身暴露四个操作使用方法见 tensor-pool README 的最小示例from dora import Node from dora_tensor_pool import get_tensor_info, tensor_from_info # helpers仅 torch node Node() pool_id node.register_tensor_pool(get_tensor_info(tensor), devicecuda:0) node.send_output(frame, pool_id) # 把 id 当普通数据传出去 ... node.write_tensor_pool(pool_id, get_tensor_info(next_tensor)) node.free_tensor_pool(pool_id)辅助函数get_tensor_info/tensor_from_info与 CUDA IPC 工具位于 python/dora_tensor_pool/。8.3 DORADMA 段头布局共享内存段以 DORADMA 魔数开头后面依次是元数据 JSON 长度、数据偏移、IPC 标志与句柄、seqlock 写世代号偏移大小字段08magic—bDORADMA\x0088json_len— u64 LE元数据 JSON 字节长度168data_off— u64 LE张量数据相对 shmem 基址的字节偏移248ipc_flag— u64 LEipc_handle有效时为 13264ipc_handle— CUDA IPC 内存句柄仅ipc_flag 1时有效968write_gen— u64 LEseqlock偶数写完奇数写入中104152reserved256Njson — 按 256 字节对齐填充的元数据 JSON256NMdata — 张量载荷8.4 传输路径分类与固定阈值从transport.rs源码看传输路径的选择是纯逻辑、可在无 GPU 的 CI 中穷举测试的should_pinCPU 源且大小 25 MiBDMA_PIN_THRESHOLD_BYTES才做cudaHostRegister固定——pinned-DMA 带宽在 25 MiB 处超越可换页拷贝 固定/解固定固定成本约 100µs这是 2026-06-27 消融实验测得的跨界点classify_transport2³ 8 种源 CUDA × 同设备 × P2P源非 CUDA 恒为SameDeviceDtoD跨设备且无 P2P如 RTX 5090 / Blackwell 系走HostStagingTransitCPU 页锁定中转 DtoH→HtoDclassify_write_pathipc_present × is_cuda × transit_ptr决定 5 条可到达写路径CpuToGpuPoolDma、GpuToGpuPoolTransit、GpuToGpuPoolDtoD、GpuToShmem、CpuToShmem。这些矩阵都有对应的#[cfg(test)]单元测试pin_tests、transport_tests中的 8-case 全矩阵测试。8.5 池 id 与随机种子#3015 修复共享内存段名格式为dora_pool_{machine_id}_{dataflow_id}_{node_id}_{counter}无 machine 时省略该前缀counter 恒为最后一段parse_pool_counter据此解析。为避免节点重启后复现上一化身共享内存名导致ShmemConf::create()O_EXCL 语义碰撞进程级计数器PINNED_COUNTER现在以随机u64种子初始化random_u64_seed()基于RandomState的 OS 级种子。trade-off 是崩溃循环且接收者永不释放时每次重启会泄漏一个注册表项与段直到注册表上限——完整的修复是 owner-death 池回收已跟踪为后续事项。9. 注意事项与限制汇总场景由各 YAML 的memory_pool_scenario环境变量控制cpu2cpu.yml与四个负路径 YAML 使用 CPU-only 接收者对无 GPU 的 CI 运行器安全cpu2cuda.yml、cuda2cpu.yml需要可用的 CUDA 运行时负路径场景使用减少的消息数保证生命周期校验短小聚焦--uv运行时各 YAML 的build:步骤自动从download.pytorch.org/whl/cpu安装 CPU 版 torch无需预装*_cross*.yml需要两个 daemon无法在标准 CI 上运行同机变体cpu2cpu_cross_local.yml由 torch 门控的memory-pool-smoke夜间任务覆盖已知未决缺陷详见 tensor-pool README跨进程释放清理可能被静默跳过#2935seqlock 溢出修复不完整两条内联 end-write 路径仍使用非环绕的old_gen 1#2890。池 id 跨重启碰撞#3015与节点崩溃时池未释放#2881已修复构建启用方式Python wheel 用maturin develop -m apis/python/node/Cargo.toml --features tensor-pooldaemon 半部用cargo build -p dora-daemon --features tensor-pool托管dora run/dora up进程内 daemon 的dora二进制用cargo build -p dora-cli --features tensor-pool。wheel 标志与 daemon 标志相互独立没有 daemon 半部时节点侧传输在单 daemon 内仍可用失去的是启动时孤儿段回收与整个跨机路径镜像段、直连 TCP 数据面、跨机注册。10. 总结从cpu2cpu.yml的单 daemon 吞吐测试到cpu2cpu_cross_local.yml的同机零拷贝直读约 5800 MB/s vs 中继 89–113 MB/s再到*_cross*.yml的跨机 commit-acknowledged 镜像写这套示例完整覆盖了 DORA 内存池传输的三大拓扑与全部生命周期负路径。配合源码阅读你可以确认零拷贝的正确性依赖 turn-based 握手与帧序守卫跨机可靠性依赖镜像确认与 120s 超时而“大帧走 WAN”目前只有这条路径能胜任。若要在自己的数据流中使用请务必先读 tensor-pool README 的兼容性声明并锁定 dora 版本。【免费下载链接】doraDORA (Dataflow-Oriented Robotic Architecture) is middleware designed to streamline and simplify the creation of AI-based robotic applications. It offers low latency, composable, and distributed dataflow capabilities. Applications are modeled as directed graphs, also referred to as pipelines.项目地址: https://gitcode.com/GitHub_Trending/do/dora创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表