Skip to content

【Task.026】TransferQueue RDMA - RFC #217

Description

@gongshaotian

1 摘要

Task tracker: #86
Task list:https://github.com/redai-infra/community/blob/main/contributor-program/2026-cohort-1/official-task.md
Auther: @gongshaotian @overloadedHenry

1.1 问题

任务书描述的问题是"图片 pixel_values 体积大,现有 TransferQueue 传输可能成为多机
多模态流水线瓶颈",给出的手段是"增加 RDMA transport,并在不可用时安全回退"。

在固定基线上逐行核对代码后,有两个事实决定了方案设计。

事实一:RDMA 实现已经存在。 固定基线的 TransferQueue 内置 MooncakeStore,具备 tcp | rdma protocol、GDR stagingoversize chunkingretrymasterbootstrapconfig.yamlbackend.storage_backend 就是开关,而 Relax 从未给这个字段赋值。缺口不在传输实现,在 Relax 侧的配置入口、能力探测与一致回退

事实二:只把 RDMA 打开,它跑不出该有的速度。 Relax 数据面里唯一以百 MB 计的字段 multimodal_train_inputs 是一个 list[dict],必然落在 MooncakeStore唯一的慢路径上:多付一次全量 memcpy、每次 put/get 都重新注册并注销 MR(get 侧还额外每次全量新分配目标缓冲),并且结构上永远进不了 GDR(GDR 分支只接收 torch.Tensor 类型的 key)。同时这个字段在 GRPO 组内存在 n_samples_per_prompt 份逐字节相同的副本——而且这 N 份是被 HF processor 重复算了 N 次才产生的,尽管组内样本本来就按引用共享同一张原图;消费侧默认还是"单 rank 拉全量 + NCCL 扇出",而换成按 rank 直取后,读侧的这些固定开销会被同一 DP 组内的 TP×PP×CP 个 rank 各付一遍。

1.2 六层方案

待解决问题 手段 与 transport 的关系
L0 无法回答"多少字节、花在哪一段" 数据面字节会计 + 分段计时 前置
L1 没有 RDMA 入口,且朴素回退会挂死 backend/protocol/device 配置 + 能力探测 + job 级一致回退 本体
L2 最大字段走慢路、拿不到 GDR payload 正规化为 spec + jagged Tensor 正交,且是 GDR 的前提
L3 组内 N 份逐字节相同副本 构造侧按引用去重 + put 侧去重 正交,收益与 transport 无关
L4 单 rank 拉全量再 NCCL 扇出 消费侧按 rank 直取覆盖多模态 正交,但读侧开销被乘 TP×PP×CP,以 L5 为前置
L5 CPU 路径 MR 零复用;get 每次全量新分配 宿主常驻注册池 + get 目标缓冲复用(对齐上游 GDR 路径已有做法) 正交,是 L4 成立的条件

L1 让 RDMA 可用,L2–L5 让 RDMA 有用。 每层独立开关、独立度量、独立 PR,默认值等价当前
main 行为。生效顺序按依赖而非编号:L5 编号在后但必须先于 L4 打开,因为 L4 把读侧的字节量
与固定开销一起乘以 TP×PP×CP,而写侧每 DP 组只付 1 份(§4.5、§4.6)。


2 目标

2.1 交付目标

  1. Relax 能以显式配置启用 TransferQueue 的 RDMA 传输,覆盖 protocol、RDMA
    device、GDR 三个维度。
  2. 无 RDMA 环境自动、一致地回退:job 级唯一决策,分级降级,且不会挂死。
  3. 数据面 payload 在任何 backend/protocol 组合下逐字节一致,且一致性责任分段可测。
  4. 连接建立、背压、超时、断连、重试、清理六项行为各有对应测试。
  5. 同拓扑多档 payload 各跑 ≥3 次,有效带宽 ≥ +20% 或 p95 延迟 ≤ −20%。
  6. 超出任务书的一条:给出 RDMA 的净收益,与 payload 形状、去重、扇出三项
    收益分开报(理由见 §2.3)。
  7. 开关、前置条件核对、排障对照、回滚路径进 docs/

2.2 非目标

非目标 理由
在 TransferQueue 里新写 RDMA transport 目标是让 Relax 具备可用的 RDMA 传输路径,实现方式是接出既有 MooncakeStore(protocol=rdma)而非重写传输层。
改 TransferQueue 的 put/get 分流逻辑 本方案使用其现成的快路。唯一可选的上游小 PR 是 metrics 注册判定(§4.2.5)
权重同步链路 权重走 DCS(NCCL broadcast),与本任务无关
运行中动态切 backend in-flight 数据可能已在远端落盘而仅丢 ack,改写会造成重复数据;只在首次 tq.init 之前回退(§4.8)
Colocate(sync)模式 该模式数据面不过 TQ,配置被忽略
跨 batch 去重 与 TQ 的 clear/GC 语义耦合,收益不值这个风险(§4.4)

2.3 为什么需要归因

验收门槛是"有效带宽 ≥ +20% 或 p95 ≤ −20%"。但当前基线含 8× 组内冗余、且最大
字段多付一次全量 memcpy——在这个基线上几乎任何改动都能越过 20%,预估仅把
backend 从 SimpleStorage 换成 Mooncake/TCP 就能达到20%(待测试)。

因此"全开 vs main 提升 40%"这样的数字无法回答评审最想问的问题:这 40% 里有多少
是 RDMA 带来的?
如果答案是 3%,那这个方案就没有完成任务书交给它的事。所以
"基线是否诚实"是方案的一部分,而不是 benchmark 的执行细节(§5.4)。


3 现状

以下每条都在固定基线 commit 与 Relax main 上逐行核对过,位置索引见附录。

3.1 数据面在哪

fully-async / hybrid 模式下 Rollout 与 Actor 在不同 GPU 集群,两者之间唯一
数据通道是 TransferQueue。

Image

红色的 SU 与橙色的 rank0 fetcher 是两个作用点:**一个是字节怎么过网,一个是
字节到了之后怎么扇出。

3.2 payload 有多大

sglang_rollout.py:297-325generate() 是 per-sample 的,_run_image_processor
对每个 sample 各跑一次;utils.py 把结果作为 list[dict] 放进要 put 的
TensorDict,并显式不做张量化:

# relax/utils/utils.py:153-154
if args.multimodal_keys is not None:
    train_data["multimodal_train_inputs"] = [s.multimodal_train_inputs for s in samples]

# relax/utils/utils.py:276-277
if key == "multimodal_train_inputs":
    tensor = value      # 保留原始 list,进 TensorDict 后成为 NonTensorStack

量级上,Qwen3-VL 每个视觉 token 展开为 3 × temporal_patch(2) × patch(16)²
fp32 元素,即 bytes ≈ 96 × patch_size² × visual_tokens。一张约 400 token 的图
约 7.4 MB,--image-max-token-num 上限下单图可达数百 MB。一个 global batch 的
多模态字段是 GB 量级,而 input_ids / logprobs / advantages 等文本字段合计
通常在 MB 量级。数据面体积几乎完全由这一个字段决定。

3.3 十条核对结论

# 事实 位置 后果
F1 MooncakeStore 已存在,storage_backend 就是开关;另有第二条现成 RDMA 路径 Yuanrong.enable_rdma(UCX H2H) config.yaml:20-103;Relax controller.py:145-192 只构造 SimpleStorage 字段 全部流量走 ZMQ/SimpleStorage L1
F2 put 按值的运行时类型分流;慢路多一次全量 memcpy;GDR 分支只收 tensor_keys mooncake_client.py:162-200:263-293serial_utils.py:420-429 非张量字段付出两趟全量 memcpy(serial_utils.py:85 + :427);含非原生类型叶子时第三趟 pickle(:119/:124/:127)。快路 CUDA 张量付一次 D2H(mooncake_client.py:676-677) L2
F3 multimodal_train_inputsNonTensorStack,逐行取出是 Python dict utils.py:153-154,276-277managers/base.py:474-494 全 batch 最大的字段必然命中慢路 L2
F4 KV key = {global_index}@{field},无内容语义;组内样本按引用共享 multimodal_inputs_shallow_copy_sample),但 processor 仍 per-sample 跑。产出的 pixel_values 推断为逐字节相同(同一 processor 对同一输入的输出),但各份是独立分配的张量(sglang_rollout.py:250-254torch.from_numpy),字节相同性取决于 processor 确定性,需在实现前实测验证。 managers/base.py:456-471data_source.py:20-36,218-228sglang_rollout.py:322-325(对比已有的 base64 组级去重 :649-658 processor CPU 与数据面字节各 N 份,shipped 配置下 8× L3a, L3b
F5 MR 在每次 put / 每次 get 注册后即注销(四条路径同构,均 register + finally unregister),merge_contiguous_memory 只合并本次调用内的相邻区间;get 侧还每次 torch.empty 全量新分配目标缓冲。而上游已在 GDR 路径实现进程级常驻注册GdrStaging,docstring 自述 "registered once for the process lifetime … avoids repeated cudaMalloc/register overhead"),CPU 路径没有对应设施,:266 留着 TODO put :257-261:287-291;get :411-415:502-506tensor_utils.py:84,133-165mooncake_utils.py:90-101,143BATCH_SIZE_LIMIT=400 两侧都切批(put :178,183/get :377,388 RDMA 的 pin-page 固定开销按调用次数与 key 粒度重复支付,且无法跨 step 摊薄 L2, L5
F6 多模态默认由 cohort 内单 rank(TP=0, PP=0, CP=0)拉全量再 CP→TP→PP 三跳 NCCL 扇出;per_rank_fetch=True跳过全部三跳广播,改由同一 DP 组内每个 TP/PP/CP rank 各自 get_meta+get_data 拉一份逐字节相同的数据 stream_dataloader.py:566,597-605,633-649,722,750-782,819-820arguments.py:213-226 换 transport 不改扇出串行度;而打开按 rank 直取后,读侧全部字节与固定开销被乘 TP×PP×CP(写侧只有 1 个 rank put,actor.py:2239-2242 L4, L5
F7 spec/tensor 拆分代码已在仓库里,只用在 NCCL 那一跳 stream_dataloader.py:489-517 现成能力没交给 TQ L2
F8 MooncakeStore.auto_init 默认 true,会 pkill -f "[m]ooncake_master" mooncake_bootstrap.py:55-65 共享 Ray 集群上新 job 杀掉他人 master L1
F9 metrics 注册硬判定 storage_backend == "SimpleStorage" interface.py:198-215 换 backend 静默丢失存储侧指标 L0+L1
F10 命名 actor 先建、store_config 后存,且 storage 创建失败不抛异常 interface.py:83-87,109-118,152,181-195 朴素回退变成挂死 L1

其中三条决定了方案形状,展开说明。

F2 + F3:最大的字段一定走慢路

mooncake_client.py:162-200 按值的运行时类型分流:

Image
  • 快路 :252-261:直接对 tensor 自身内存 register_bufferbatch_upsert_from
    unregister_buffer。零序列化;CUDA 张量付一次 D2H(mooncake_client.py:676-677),CPU 张量零额外拷贝。
  • 慢路 :263-293:msgpack encode(张量叶子本身已是零拷贝)→
    serial_utils.pack_into:420-429 把所有 item 完整 memcpy 进一块新分配的
    uint8 region → 注册 → upsert → 注销。这次拷贝无法回避:Mooncake 要求一个 key
    对应一段连续可注册内存,上游自己也留了 TODO(:266)。
  • GDR 分支只接收 tensor_keys:170-172),慢路的值永远不进 GDR。get 侧用
    dtype is not None 判定,同理(:334-341)。

managers/base.py:474-494_generate_valuesNonTensorStack 是逐行取出
Python dict,只对 nested tensor 才调 unbind()

for field in sorted(data.keys()):
    field_data = data[field]
    if isinstance(field_data, Tensor) and field_data.is_nested:
        results.extend(field_data.unbind())     # → 逐样本 Tensor,命中快路
    else:
        results.extend(field_data)              # → 逐样本 Python dict,命中慢路

结论:全 batch 中体积最大的字段,是唯一必然付出"encode + 全量 pack memcpy +
每次 put 重注册 MR"的字段,且拿不到 GDR。

F4:同一张图在组内被编码、拷贝、注册、传输、存储各 N 次

managers/base.py:456-471 生成 key:

return [pfx + sfx for sfx, pfx in itertools.product(keys_suffixes, keys_prefixes)]
# → '0@input_ids', '1@input_ids', ..., '0@pixel_values', ...

key 里只有 global_index 和 field 名。而 GRPO 一个 prompt 的 n_samples_per_prompt 个 sample 按引用共享同一张图(data_source.py:218-228_shallow_copy_sample 展开组,:20-36 明确注明共享 multimodal_inputs),但 generate() 是 per-sample 的,各自跑一次 _run_image_processorsglang_rollout.py:322-325),产出的 pixel_values 内容相同但 global_index 不同 → N 个不同 key → 全链路各做 N 次。shipped 配置 --n-samples-per-prompt 8 即 8×。

需要区分的是:sglang_rollout.py:649-658 已经对base64 编码_encode_multimodal_inputs:264-294,喂给 SGLang HTTP 请求体的那一步)做了组级去重,判定方式是 multimodal_inputs is 同一对象。但真正产出百 MB 级 pixel_values_run_image_processor:207-261不在这个去重范围内。这条现成先例同时说明:组内引用共享是可以被直接利用的既有事实,不必绕道内容哈希——也因此不依赖上表 F4 里"processor 输出是否逐字节确定"这个待实测的前提(§4.4)。

F10:朴素回退不会报错,会挂死

interface.pyinit() 顺序是:创建命名 actor TransferQueueController
:181-184)→ _maybe_create_tq_storage:193)→ 最后
store_config:195)。而 _maybe_create_tq_storage 在 provider 返回 None 时只
logger.error不抛异常:83-87)。

于是 Mooncake bootstrap 失败后,命名 actor 已存在但 config 从未存入;此时直接重试
tq.init(simple) 会走 _init_from_existing():152):

_TQ_CONTROLLER = ray.get_actor("TransferQueueController", namespace="transfer_queue")
conf = None
while conf is None:                       # interface.py:109-118
    conf = ray.get(_TQ_CONTROLLER.get_config.remote())
    time.sleep(1)

→ 1 秒轮询的死循环,表现为训练卡住且无任何报错。 回退实现因此有两条硬约束:
能力探测发生在 tq.init 之前;回退前必须 ray.kill 这个命名 actor。

3.4 把 F1–F10 放到一次 put 的时间轴上

Image

非连续→连续 44.601 msCPU encode + pack 33.660 ms,而它要优化的 Mooncake put 只有 22.331 ms待加速的那一段比它
前面的串行开销还小。

这不是否定 L1——L1 必须先有,否则没有 RDMA 可谈。这是说 L1 之外还有比它更大的一块,
而那一块正好也是 GDR 能否生效的前提。


4 设计方案

4.0 改造前后

Image

改动位置:

位置 改造前 改造后
processor 调用 组内 N 个 sample 各跑一次(引用其实是共享的) 组内共享,跑一次 L3a
payload 形状 list[dict](非张量,慢路) spec + jagged Tensor(张量,快路) L2
组内副本 N 份实体 1 份实体 + N−1 份引用 L3b
transport ZMQ / SimpleStorage MooncakeStore RDMA(+GDR 可选) L1
get 侧缓冲 每次 torch.empty 新分配 + 注册 + 注销 常驻注册池,跨 step 复用 L5
扇出 cohort 内单 rank 拉全量 + CP→TP→PP 三跳 NCCL 串行扇出 跳过三跳广播,cohort 内每 rank 各自直取 L4

上图只画了数据面(过 TQ 的部分);L3a 发生在数据面之前的 rollout 侧(generate_and_rm_group 内),不在图上。注意"扇出"这一行改造后读侧总字节量是上升的(同一份数据被 cohort 内每个 rank 各读一次),换来的是去掉 pickle + 三跳 NCCL 的串行路径——因此它以 L5 为前置,见 §4.5、§4.6。

4.1 L0:数据面字节会计与分段计时

问题:换 backend 会丢掉 TQ 存储侧指标(F9),而 §3.4 的时间轴目前全靠推算。
改动之前需先让"多少字节、分几段、每段多久"可被回答。这一层不改行为,只加观测,
本身就是交付物——任务书要"有效带宽",而有效带宽 = 有效字节 / 端到端时间;
没有 L0,"有效字节"是个说不清的量,±20% 的结论无法复核。

实现点全在 Relax 侧,不改 TQ:

  1. relax/utils/utils.py 组装 put payload 处,按 field 统计
    nbytes / dtype / 是否张量 / 样本数,聚合成一条 per-step 记录。
  2. stream_dataloader.py 消费侧同样统计,并对 rank0 扇出路径单独计时。
  3. 分段计时点与 §3.4 时间轴一一对应:image_processor / 非连续→连续 / TQ put /
    TQ get_meta / TQ get / NCCL 扇出 / 重建 TensorDict。
  4. 新增 --dataplane-accounting {off,summary,verbose},默认 offsummary 只在
    step 边界打一行聚合,verbose 落 per-field JSONL。

核心指标(Prometheus,沿用 Relax 现有 exporter):

指标 类型 标签 用途
dataplane_bytes_total Counter field, path=fast|slow, direction 证明 F3,量化 L2 目标体积
dataplane_bytes_unique_total Counter field 量化 F4 冗余率 = L3b 收益上限
dataplane_stage_seconds Histogram stage 证明 §3.4 时间轴,定位真实瓶颈
dataplane_fanout_seconds Histogram mode=rank0|per_rank L4 收益
dataplane_backend_info Gauge=1 backend,protocol,device,gdr 证明跑的是声称的路径
dataplane_fallback_total Counter from, to, reason 暴露静默降级

另外,启动时无条件打印一行生效配置摘要(不受 --dataplane-accounting 控制):

[dataplane] backend=MooncakeStore protocol=rdma device=mlx5_0 gdr=on
            normalize=on dedup=on per_rank_fetch=on
            fallback=none (probe: 8/8 nodes ok)

这一行是排障的第一现场,也是 benchmark 报告必须附带的证据——它直接否掉"以为在跑
RDMA、其实在跑 TCP"这类无效结论。

4.2 L1:RDMA transport 入口(本体)

4.2.1 配置面

新增 CLI(relax/utils/arguments.py::add_transfer_queue_arguments,属于项目
"Ask First" 文件,需评审确认):

--tq-storage-backend {simple,mooncake}     # 默认 simple,语义等价现状
--tq-transport {tcp,rdma}                  # 仅 mooncake 生效
--tq-rdma-mode {off,auto,required}         # 默认 off
--tq-rdma-device NAME                      # 空 = Mooncake 自选
--tq-use-gdr / --no-tq-use-gdr             # 默认关
--tq-gdr-staging-mb N                      # 默认 1024
--tq-mooncake-metadata-server ADDR         # 默认 P2PHANDSHAKE
--tq-mooncake-master-address ADDR          # mooncake backend 下必填
--tq-mooncake-auto-init {never,job-owned}  # 默认 never,见 4.2.4
--tq-fallback-backend {simple,none}        # 默认 simple
--tq-init-timeout-seconds N                # 默认 30

现有 TQ 参数块(arguments.py:206-245)不使用前缀;新参数引入 --tq- 前缀以区分新增项与原有项,此命名约定待维护者确认。三态 --tq-rdma-mode 使用 choices,与 --checkpoint-engine-backend--distributed-backend(均无 choices)的现状不同,因为三态语义需要启动期校验拒绝非法组合。

三态 --tq-rdma-mode 的语义:

mode 探测失败 探测通过
off 不启用(即使硬件可用)
auto 记录原因 → 回退 --tq-fallback-backend 启用
required fail fast,非零退出 启用

required 存在的唯一理由是 benchmark 与生产:静默降级会让一次"RDMA 提升 20%"的
测量实际跑在 TCP 上。

4.2.2 启动时序:探测必须早于 tq.init

F10 决定了顺序不能随意:

Image

三个不可省的细节:

  1. 探测在 tq.init 之前。 事后探测意味着命名 actor 已创建,回退路径就落进
    F10 的死循环。
  2. 回退前必须 ray.kill 命名 actor(name TransferQueueController
    namespace transfer_queue),否则第二次 tq.init_init_from_existing()
    while conf is None 轮询 → 挂死。
  3. 不能依赖 TQ 抛错。 _maybe_create_tq_storage 失败不抛异常(F10),所以
    Relax 侧必须在 tq.init 返回后主动校验 storage 真实存在:校验返回 conf 与预期
    backend 一致,加一次 dry-run put/get 往返。

4.2.3 能力探测与分级降级

探测在每个将参与数据面的节点各跑一次,返回结构化报告:

检查 手段 失败含义
mooncake 可导入且版本一致 import mooncake + version 集群镜像不一致,必须降级
存在 RDMA 设备 /sys/class/infiniband/* 无 IB/RoCE 硬件
HCA 端口 ACTIVE .../ports/*/state 网卡未起链路
GID 可用 .../gids/* 非全零 RoCE 未配
memlock 无限制 resource.getrlimit(RLIMIT_MEMLOCK) MR 注册会失败(容器最常见)
节点间可达 两节点 handshake 跨机不通;单机 loopback 会假通过
GDR 前置 nvidia_peermem / gdrdrv 模块 只关 GDR,不关 RDMA

降级是分级的,不是一刀切

GDR 不满足        → use_gdr=false,保留 rdma
RDMA 不满足       → protocol=tcp,保留 MooncakeStore
mooncake 不满足   → fallback_backend(simple)

这比"能否 RDMA"的二元判断重要:GDR 的前置条件(peermem 模块 + 足够 staging 显存)
在很多集群上不满足,但 RDMA 本身完全可用。

4.2.4 多租户安全(F8)

MooncakeStore.auto_init 上游默认 true,会 pkill -f "[m]ooncake_master"
Relax 的 scripts/entrypoint/ray-job.sh 面向共享的、预先存在的 Ray 集群
直接透传这个默认值等于允许新提交的 job 杀掉同集群其他 job 的 master。处理方式:

  • Relax 侧默认 --tq-mooncake-auto-init never,并把 --tq-mooncake-master-address
    mooncake backend 下设为必填。master 由运维/启动脚本显式提供,TQ 不自己抢。
  • job-owned 模式下:端口从 job 专属区间分配、master 进程记录 job id、清理时只按
    job id 精确终止
    ,绝不使用宽匹配 pkill / killall
  • metadata_server 默认 P2PHANDSHAKE,避免额外引入一个需要独立运维的
    etcd / http metadata 服务。

4.2.5 可观测性回退补偿(F9)

interface.py:198-215 的硬判定意味着切到 MooncakeStore 后 TQ 存储侧 Prometheus
指标消失。两条应对:

  1. 短期(本任务范围内):由 L0 的 Relax 侧会计覆盖字节量与分段耗时,保证换
    backend 后仍能回答带宽/延迟问题;文档明示"存储侧 TQ 原生指标在 mooncake
    backend 下不可用"。
  2. 长期:向 TransferQueue 提一个独立小 PR,把该分支改为按"storage unit 是否
    实现 metrics 接口"判定而非按 backend 名。这是上游缺陷,不阻塞本任务。

4.3 L2:payload 正规化

问题:F2+F3——全 batch 最大的字段结构性地走在慢路上,且永远进不了 GDR。

形状变换

复用仓库里已有的 _encode_multimodal_inputs(F7,stream_dataloader.py:489-517),
把它从 NCCL 那一跳前移到 TQ 这一跳:

Image

TQ 侧不需要任何改动_generate_values:490-491 对 nested tensor 调 unbind()
得到逐样本 Tensor → 命中快路;_merge_tensors_to_tensordict:581-587
as_nested_tensor(layout=jagged) 还原。这条路径是 TQ 原本就设计好的,Relax 只是
没把这个形状交给它。

收益来源(三项互相独立)

来源 机制
消除 pack_into 全量 memcpy 快路直接 register_buffer(tensor.data_ptr())
打开 GDR 的可能性 GDR 分支只接收 tensor_keys(F2),非张量永远进不去
减少 MR 注册次数 同 field 的连续张量可被 merge_contiguous_memory 合并(F5)

第三项要求张量在内存中连续,因此正规化时对同一 field 的样本张量做一次连续化打包

兼容与开关

  • 新增 --tq-normalize-multimodal,默认关。
  • 打开后 put 侧写 mm_spec + mm_<key> 系列字段,不再写
    multimodal_train_inputs;消费侧按是否存在 mm_spec 自动选择解码路径。
  • 断点续训 / staleness 混批场景下 TQ 中可能同时存在两种形状,消费侧按字段存在性
    逐 batch 判定,不做全局假设。
  • stream_dataloader 原有的 rank0/NCCL 路径继续可用:spec+tensor 形状本来就是为它
    设计的。

4.4 L3:组内冗余消除

问题:F4——同一张图在 GRPO 组内被完整走完全链路 N 次。

冗余的起点是明确的:data_source.py:218-228_shallow_copy_sample:20-36)把一个 prompt 展开成 n_samples_per_prompt 个 sample,注释写明 multimodal_inputs 按引用共享。组内 N 个 sample 从此持有同一个对象。

链路上因此有两处独立的重复,各需一次去重:

位置 重复什么 代码位置
L3a 构造侧 HF processor 跑 N 次,产出 N 份内容相同的 multimodal_train_inputs sglang_rollout.py:322-325,在 per-sample generate()
L3b put 侧 utils.py:154 把 N 个引用摊成 list → N 个 {idx}@multimodal_train_inputs key → N 次 encode、N 次传输、N 份存储 managers/base.py:456-494

两处必须都做,不是二选一。 只做 L3a 省下的是 processor CPU:utils.py:154 仍然产出 N 个列表元素,key 按下标生成,TQ 侧依旧 N 次 encode、N 份字节。只做 L3b 省下的是带宽:processor 那 N−1 次调用已经花掉了。

L3a 构造侧去重

仓库里已有同型先例,且是按对象身份而非内容哈希判定的——generate_and_rm_group 对 base64 编码做的组级去重(sglang_rollout.py:649-658):

first_mm = getattr(group[0], "multimodal_inputs", None)
if first_mm is not None and all(getattr(s, "multimodal_inputs", None) is first_mm for s in group[1:]):
    encoded_mm, t_enc = await _encode_multimodal_inputs(first_mm)
    for sample in group:
        sample._pre_encoded_mm = encoded_mm

这个先例只覆盖 _encode_multimodal_inputs:264-294,把媒体转成 SGLang HTTP 请求体里的 base64)。真正产出百 MB 级 pixel_values_run_image_processor:207-261)仍在 per-sample generate() 内被调用(:322-325),没有组级去重。

L3a 把同一个 is 判定延伸到它:在 generate_and_rm_group 内、创建 per-sample task(:660-668)之前,若组内 multimodal_inputs 为同一对象且 prompt 相同,只跑一次 _run_image_processor,把 (processor_prompt_ids, multimodal_train_inputs) 挂到组内每个 sample 上;generate() 命中缓存则跳过 :322-325。这与 _pre_encoded_mm 的挂载/消费方式(:361-369)一致。

收益 / 性质 说明
省 processor CPU N=8 时省 7 次 HF processor 调用;该开销已在度量(perf_detail/rollout/image_processor_time:621-626
不依赖 processor 确定性 判定的是输入引用是否同一,与 F4 里"输出是否逐字节相同"这个待实测前提无关
无碰撞风险 is 判定不是内容哈希,R4 不适用
无哈希成本 不需要额外读一遍 N×MB 字节
让 L3b 退化为身份判定 组内 N 个 sample 从此持有同一个 multimodal_train_inputs 对象,put 侧 id() 即可识别重复

前置审计(合并前必须完成,因为共享引用把"只读"从惯例变成了硬约束):

  • 下游必须只读。已核对主训练路径 megatron/data.py:517-544:只做 append + torch.cat / pad_and_flatten,不 in-place 改源 dict;data_source.py:34-35 的注释也断言该字段是"被 set 而非被 mutate"。其余消费点(stream_dataloader.py:713-741utils/multimodal/stats.pytrain_dump_utils.py:224-225)需逐一确认,并补一条"共享引用后训练数值不变"的端到端比对。
  • prompt 必须一致才能共享 processor_prompt_ids。partial rollout、abort 重入等场景下若组内 prompt 已分叉,is + prompt 双条件判定失败,自动退回逐 sample。
  • agentic 路径自动 no-opagentic/pipeline/runtime.py:3474 逐轮生成新图,session/service.py:1185copy.deepcopy,引用天然不共享。

L3b put 侧去重

for each prompt group g:
    k = id(tensor)  if L3a 生效           # 身份判定,零成本
        blake2b(bytes(tensor))  否则       # 兜底,仅对 L2 产出的连续张量
    if k 首次出现:
        写 mm_pixel_values[idx] = tensor  # 实体
        ref[idx] = None
    else:
        写 mm_pixel_values[idx] = 空张量   # 占位
        ref[idx] = 首份的 global_index
写 mm_ref                                  # KB 级,随 spec 一起走慢路

消费侧按 mm_ref 把引用样本指向首份张量。同 batch 内必然可见,因为去重只在同一 partition / 同一 put 批次内生效。

边界与生命周期:

问题 决定
身份 vs 哈希 L3a 打开时用 id();L3a 未打开或引用已分叉时才走 blake2b。哈希是兜底,不是主路径
跨 batch 去重 不做。跨 batch 引用会与 TQ 的 clear/GC 语义耦合,收益不值这个风险
哈希碰撞(仅兜底路径) blake2b-128,同时校验 shape + dtype + nbytes;碰撞概率远低于硬件误码率
哈希成本(仅兜底路径) 与一次 memcpy 同量级,但只跑一次,且替代了 N−1 次 encode + memcpy
消费侧是否共享内存 默认 clone() 成独立张量,避免下游 in-place 改动串写;提供 --tq-dedup-share-view 供只读场景
图片本就各不相同(非 GRPO / n=1) 判定全部 miss,退化为一次多余扫描;由 --tq-dedup-multimodal 默认关 + L0 冗余率指标指导开启

这一层的收益与 transport 完全无关:即使继续用 SimpleStorage/ZMQ,8× → 1× 也是
8 倍的数据面体积削减。这正是它必须与 L1 分开度量的原因(§5.4)。

4.5 L4:消费侧按 rank 直取

问题:F6 — per_rank_fetch 默认 False(stream_dataloader.py:566),多模态默认由 cohort 内单 rank(TP=0, PP=0, CP=0,见 :640-649)拉全量,再经 CP→TP→PP 三跳 broadcast_object_list 扇出(:757-782),其中 rank0 侧要先 pickle 整个 payload。per_rank_fetch=True 时(:633-639三跳广播全部跳过,改由同一 DP 组内每个 TP/PP/CP rank 各自 get_meta + get_data;跨 rank 一致性靠 TQ sampler 在 (partition_id, task_name, dp_rank, batch_index) 上的缓存保证——即各 rank 拿到的是逐字节相同的同一份数据,不是各自的分片。多模态已被这条路径覆盖(:722 的条件是 has_multimodal and not per_rank_fetch,打开后该字段留在 TensorDict 内由各 rank 自取;:819-820 重建)。做法是在多模态场景下将其设为默认开启,并给出实测依据。
Image

L4 换的是什么,必须说准。 上游 docstring 写得很直白(:601-603):Trades a single rank-0 pickle + one NCCL bcast for N parallel ZMQ deserialises — wins when pickle dominates tgd_bcast_tp_time。所以 L4 不是省字节的优化,而是用读侧字节放大换掉 pickle + 三跳 NCCL 的串行路径。放大倍数是被跳过的广播 cohort 大小,即 TP×PP×CP:627 注释举的实际拓扑是 TP2/PP2/CP8。而写侧不放大:relax/backends/megatron/actor.py:2239-2242 把 put 限定在 TP==0 and is_pipeline_last_stage()CP!=0 直接 return,每个 DP 组恰好 1 个 rank put。

四个注意点:

  • 仅在 backend 为 KV(mooncake)时有意义。SimpleStorage 下 SU 数量有限,per-rank
    并发拉取可能反而打满少数 SU,因此 L4 的默认值随 --tq-storage-backend 联动。上游
    --per-rank-fetch 的 help 也要求配套 --num-data-storage-units >= TP world size
    arguments.py:223-224),这本身就是"读侧压力被放大"的自述。
  • 放大轴是 TP×PP×CP,不是 DP。 不同 DP rank 本来就读不同数据,与本开关无关;
    L4 恰恰是把 TP/PP/CP 三跳广播换成 cohort 内各 rank 各读一份。放大的也不只是固定
    开销,而是整份字节量
  • --per-rank-fetchrollout_routed_experts 互斥actor.py:2160 会在该字段
    存在时自动关闭(jagged NestedTensor 的广播路径不兼容)。多模态场景要确认 data_fields
    里没有该字段,否则开关静默失效、benchmark 结论会挂在错误配置上。
  • L4 的收益以 L5 为前置(§4.6):读侧每次 get 都要新分配全量目标缓冲、注册、注销
    (F5),这些开销乘 TP×PP×CP 之后可能吃掉省下的 pickle 时间。因此 L4 与 L5 必须成对
    评估,B 组里也要单独隔离(§5.4)。

与 L3b 的耦合点,成因需归正。 L3b 的"占位样本取不到首份"约束来自 DP 分片本身:consumer 只读自己 DP 分片的 global_index 集合,而 put 侧的去重域是整个 partition,首份可能落在别的 DP rank 的分片里。这与 L4 开关无关——per_rank_fetch 不改变一个 DP 组读到哪些样本,只改变有几个 rank 去读。所以约束应写成:只要 consumer 按 DP 分片读,L3b 的去重域就必须收缩到 DP 分片域,或按分片分别去重(风险 R5)。L3a 发生在 rollout 侧,与分片域无关,不受此约束。

4.6 L5:读侧固定开销削减(L4 的前置)

问题:F5 — MooncakeStore 的 CPU 路径对 MR 零复用。四条 worker 结构完全同构,都是"注册 → 传输 → finally 注销":

路径 位置 每次调用做的事
put 张量 :257-261 _preprocess_tensors_for_put(CUDA 张量付一次 D2H :675-676,非连续付一次 .contiguous() :677)→ merge_contiguous_memory → register → upsert → unregister
put 非张量 :287-291 allocate_empty_tensors 新开 uint8 region → batch_encode_into 全量 memcpy → register → upsert → unregister
get 张量 :411-415 allocate_empty_tensors 新开全量目标缓冲 → register → get_intounregister
get 非张量 :502-506 同上 + batch_decode_from(零拷贝视图)

三点需要说清楚,其中一点是对常见误解的纠正:

  1. 注册生命周期两侧是对称的,不存在"put 就地注册常驻、只有 get 才注销"。快路的"就地"只是指注册调用方已有的缓冲、不额外拷贝,注册本身照样在 finally 里注销。
  2. 注册区间数量上 put 可能比 get 更多。get 侧 allocate_empty_tensorstensor_utils.py:28-95)按 dtype 分组,每个 dtype 只开一块大连续内存,注册数 = dtype 种类数;put 快路注册的是调用方那批张量,彼此一般不相邻,merge_contiguous_memorytensor_utils.py:133-165,按地址排序合并相邻区间)合不掉,注册数最坏 = 张量数。因此本层的立论不是"读比写慢"——单次成本孰高孰低取决于负载形状,本方案不预设。
  3. 读侧真正独有的额外成本是每次全量新分配tensor_utils.py:84torch.empty(total_elements)。put 快路复用调用方已经存在的张量,不付这一笔;get 每次都要开一块 payload 尺寸的新内存,并对冷页做 pin。

为什么这一层必须存在,用两条不依赖单次快慢的论证:

  • 不对称放大:写侧每 DP 组 1 个 rank put(actor.py:2239-2242),读侧在 L4 打开后每 DP 组 TP×PP×CP 个 rank 各 get 一次。读侧的任何单次固定开销都被乘一遍,写侧不乘。这使读侧固定开销在 L4 之后必然成为主导项,与"单次谁更慢"无关。
  • 上游已证明该模式可行,只是没铺到 CPU 路径mooncake_utils.py:90-101GdrStaging 就是进程级常驻注册池,docstring 自述 One cudaMalloc buffer, registered once for the process lifetimeavoids repeated cudaMalloc/register overhead,注册发生在 :143lazy_init、注销只在 :151close。但它是 CUDA-only、且只对 dtype is not None 的张量 key 生效(mooncake_client.py:334-348),非张量 key 结构上进不去。L5 不是新机制,是把这个已有模式补到 CPU 路径。另有上游自留 TODO 佐证方向一致::266 switch to a pre-registered buffer from MooncakeStore once such an API is available

做法(三项独立,可分别开关):

内容 依赖
L5a 常驻注册池 进程级 HostStaging:按 dtype/尺寸档位预分配并一次性注册若干宿主缓冲,put/get 从池中取、用完归还不注销;池未命中时回落到当前的 per-call 注册路径
L5b get 目标缓冲复用 get 的目标缓冲从池中取,消除 torch.empty 全量新分配与冷页 pin;跨 step 复用同一批已 pin 页 L5a
L5c 批次合并 BATCH_SIZE_LIMIT=400 在两侧都按 key 数切批(put :178,183/get :377,388),每批一轮注册。L2 把百 MB 字段摊成 per-sample 张量后 key 数显著上升、注册轮数随之上升;按字节预算而非 key 数切批,使注册轮数与 payload 体积解耦 L2

L5a/L5b 需在 TQ 侧改(mooncake_client.py),属上游改动,按 §4.9 的 PR 计划走 upstream PR;若上游未合,Relax 侧无法单独实现,此时 L4 默认保持关闭,这是本层与 L4 联动的硬约束。L5c 只依赖已有 API,可先落。

度量:新增 dataplane_mr_register_total{direction,path}(Counter)与 dataplane_mr_register_seconds(Histogram),以及 dataplane_get_alloc_seconds。判定标准是注册次数随 step 数增长的斜率趋于 0(池命中)而非注册单次变快——后者会被负载形状干扰,前者不会。

兜底:池化引入的是"缓冲被复用"这一新前提,一旦消费侧持有归还后的缓冲引用就会读到脏数据。因此归还必须在数据被拷出或被上层接管之后,并配一条哨兵测试(§5.1)。--tq-host-staging 默认关闭。

4.7 配置全表与兼容矩阵

所有新开关默认值都等价于当前 main 的行为,即默认零行为变更

开关 默认 打开后影响面
--dataplane-accounting off 仅日志/指标 L0
--tq-storage-backend simple TQ backend 选择 L1
--tq-transport tcp Mooncake protocol L1
--tq-rdma-mode off 探测与降级策略 L1
--tq-use-gdr false GDR staging L1
--tq-mooncake-auto-init never 是否自建 master L1
--tq-normalize-multimodal false put/get payload 形状 L2
--rollout-group-share-processor false 组内共享 processor 结果 L3a
--tq-dedup-multimodal false put 侧组内去重 L3b
--tq-host-staging false 宿主常驻注册池 + get 目标缓冲复用 L5a, L5b
--tq-batch-bytes-budget 0(=沿用 BATCH_SIZE_LIMIT 按 key 数切批) 切批依据改为字节预算 L5c
--per-rank-fetch 覆盖多模态 跟随 backend 消费侧扇出 L4

组合合法性由启动时校验强制:

组合 结果
use_gdr=true + normalize=false 警告:最大字段是非张量,GDR 对它无效(F2)
dedup=true + normalize=false 拒绝启动:put 侧去重依赖 L2 产出的连续张量
group_share_processor=true + normalize 任意 合法。L3a 在 rollout 侧生效,不依赖 L2
storage_backend=simple + transport=rdma 拒绝启动:参数无意义
rdma_mode=required + 探测失败 非零退出
rdma_mode=off + transport=rdma transport 被忽略,打印说明
L4 覆盖多模态 + dedup=true 去重域收缩到 DP 分片(§4.5)
L4 覆盖多模态 + host_staging=false 警告并要求显式确认:L4 把读侧固定开销乘 TP×PP×CP,缺 L5 时净收益可能为负(§4.6)。上游 L5a/L5b 未合入时 L4 默认保持关闭
host_staging=true + storage_backend=simple host_staging 被忽略,打印说明(注册池只对 MooncakeStore 有意义)
batch_bytes_budget>0 + normalize=false 警告:L5c 的收益前提是 L2 把大字段摊成多 key;未开 L2 时基本无效

三种执行模式下的适用性:

模式 数据面是否过 TQ 本方案是否生效
Colocate (sync) 否(同 GPU 时分复用) 不生效,配置被忽略
Fully Async 全部生效
Hybrid 是(Actor/Rollout 分 PG) 全部生效

回滚粒度就是这张开关表。 回到当前 main 行为只需去掉全部 --tq-*
--dataplane-accounting,不需要回滚代码,也不需要重启 Ray 集群以外的任何设施。

4.8 回退与生命周期

Image

三条硬规则:

  1. 回退只发生在第一次成功 tq.init 之前,之后不再切 backend。理由:一次 in-flight RDMA put 可能已在远端落盘而仅丢了 ack,把同一批数据改写到
    另一个 backend 会造成重复数据与 production_status 不一致。运行中失败由 TQ 自身
    MAX_RETRIES = 3 处理;重试耗尽则整个 job 失败,而不是降级续跑。
  2. job 级唯一决策。 探测结果按 AND 归约成一个 {backend, protocol, device}
    三元组,广播给所有 worker。禁止各节点自行决定——否则一部分 worker 用 RDMA、一部分
    用 TCP,put/get 的 key 空间会跨 backend 割裂。
  3. 每次降级留一条结构化记录 {from, to, reason, node, check},同时进日志与
    dataplane_fallback_total{reason}。静默降级是本任务最危险的失败模式——它会让
    benchmark 结论作废。

资源归属与清理:

资源 归属 清理时机 方式
TransferQueueController 命名 actor job job 结束 / 回退前 ray.kill(name, namespace)
mooncake master 进程 运维(never)或 job(job-owned job-owned 由 job 清理 按 job id 精确 kill,禁止宽匹配
MR / staging buffer TQ client put/get 结束(默认);开 L5 后为进程退出 TQ 内部 unregister_buffer;池化后走 HostStaging.close(),与 GdrStaging.close()mooncake_utils.py:147-151)同构
TQ 中的数据 Relax controller 现有 tq.close()controller.py:871-875 不变

异常退出(SIGTERM / OOM / 节点失联)下必须保证 job-owned master 与命名 actor 都不
残留,由 relax/entrypoints/train.py 现有的信号处理路径接管。

4.9 分阶段 PR

任务书要求"设计 Issue + 分阶段 PR(协议/实现、集成、基准与文档)"。拆成 6 个 PR,
每个独立可回滚、独立可评审:

PR 内容 依赖 合并门槛
PR1 数据面会计与分段计时 + 指标 + 单测 L0 CI 绿;在一次真实 fully-async 跑上产出字节会计报告,实测 F3/F4 的量级
PR2 配置面 + 能力探测 + 回退状态机 + 多租户安全 + 运维文档 L1 PR1 §5.2 CI 用例全过;多机的连接/超时/断连/清理/多租户全过
PR3 payload 正规化 + put 侧去重 + 一致性单测 L2, L3b PR2 一致性单测 + 单机"快慢路归属断言"通过
PR4 宿主常驻注册池 + get 目标缓冲复用(上游 TQ PR)+ 字节预算切批 L5 PR3(L5c 依赖 L2) 注册次数随 step 增长斜率趋于 0;缓冲复用哨兵测试通过;dataplane_mr_register_* 上报
PR5 per-rank 直取覆盖多模态 + benchmark 脚本 + A/B 报告 L4 PR4(L5a/L5b 未合入则本 PR 的开关默认关闭) §5.4 的 A 组与 B 组报告齐备,且 C6 − C5(L5)与 C7 − C6(L4)分别为正
PR6 组内共享 processor 结果 + 只读性审计 L3a 只读性审计清单逐项签字;image_processor_time 下降;训练数值与共享前逐 step 一致

PR4 的 L5a/L5b 落在 transfer_queue/storage/clients/mooncake_client.py,属上游改动,走 TransferQueue 仓库的 upstream PR;上游未合并时 Relax 侧只能落 L5c,PR5 的 L4 开关必须保持默认关闭(§4.6)。

PR6 只改 relax/engine/rollout/,与 TQ 数据面无耦合,因此不依赖 PR1–PR5,可以最先合——它也是全套改动里收益/风险比最好的一个:省 7 次 HF processor 调用,代价只是一次 is 判定。

4.10 风险

# 风险 影响 应对
R1 集群无 RDMA / memlock 受限,无法验证 L1 净收益 核心指标拿不到 提前用前置条件清单核对环境;required 模式暴露而非掩盖;单机 loopback 结论不得作为跨机结论
R2 GDR 前置不满足 C4 − C3 为 0 分级降级(§4.2.3);GDR 是可选加成,不是方案支点
R3 L2 改变 payload 形状,下游有隐式依赖字段名的代码 训练报错 全仓 grep 该字段名;开关默认关;消费侧按字段存在性自适应;端到端数值比对
R4 哈希碰撞导致错误共享 静默数据污染 首选 L3a 的 is 身份判定,根本不哈希;仅兜底路径用 blake2b-128 + shape/dtype/nbytes 三重校验 + negative 单测
R4b L3a 共享引用后,某个下游消费点 in-place 改写 multimodal_train_inputs 组内 N 个样本被串写,静默数值错误 合并前完成只读性审计清单(§4.4);开关默认关;与共享前逐 step 数值比对
R5 L3b 去重域与 L4 分片域不一致 引用样本在本 rank 取不到首份 显式约束:L4 打开时去重域收缩到 DP 分片(§4.5)+ 组合单测
R6 auto_init 误伤同集群其他 job 他人训练中断 默认 neverjob-owned 按 job id 精确清理;禁止宽匹配 kill;多租户回归测试
R7 回退路径本身触发 F10 挂死 无报错卡死,极难排查 探测前置 + 回退前 ray.kill + test_fallback_state_machine 锁死行为
R8 静默降级污染 benchmark 结论不可复现 required 模式 + 启动摘要行 + dataplane_fallback_total 三重暴露
R9 Relax CI 是 CPU-only 且 stub transfer_queue 关键路径无 CI 覆盖 L2/L3/配置校验/回退状态机全部设计为可 mock(§5.2);真机测试作为合并前人工门槛并附日志
R10 上游 TQ 后续版本改变 put 分流或 metrics 逻辑 收益前提失效 版本 pin(_MIN_TQ_VERSION + commit);"快慢路归属断言"会在升级时立刻失败报警
R11 L5 缓冲复用后消费侧仍持有已归还缓冲的引用 静默脏读,表现为偶发数值异常 归还点必须晚于数据被拷出/被上层接管;哨兵测试(§5.1)在归还后写入毒值并断言消费侧数据未变;--tq-host-staging 默认关
R12 L5a/L5b 的上游 TQ PR 未被接收 L4 缺前置,净收益可能为负 Relax 侧只落 L5c;L4 开关默认关闭并在启动摘要中说明(§4.6、§4.9);A 组结论按"L5 缺位"标注
R13 常驻注册池长期占用 pinned 宿主内存 节点内存压力、其他进程 pin 失败 池容量上界可配 + 未命中回落 per-call 路径;上报池水位;与 memlock 上限一起纳入 L1 的前置条件清单

5 一致性和性能测试

5.1 一致性:把责任分段,而不是端到端比一次

任务书要求"payload 逐字节一致"。如果只在端到端做一次比对,一旦不等就无法定位是
形状变换、去重还是网络出的问题。因此拆成三段各自可测:

保证方式 在哪测
L2 形状变换 decode(encode(x)) == x,dtype / shape / bytes 全等 单测,纯 CPU,不依赖网络
L3a 共享 processor 结果 组内共享 vs 逐 sample 各跑一次,N 份结果两两 torch.equal;且断言 processor 只被调用一次(mock 计数) 单测,纯 CPU
L3b 去重 引用样本还原后与首份 torch.equalnbytes 相同 单测,纯 CPU
L5 缓冲复用 归还缓冲后向其写入毒值,断言消费侧已持有的数据不变;同一缓冲连续两轮 get 的结果互不污染 单测,可用假 store mock register_buffer/get_into
L1 transport TQ put → get 往返 bytes 全等,四条路径两两相等 集成测试,必须真机

L3a 还需一条只读性的负向测试:在共享开启的前提下,若任何消费点对 multimodal_train_inputs 做 in-place 改写,组内其余样本会被串写。测试方式是共享后对首份做一次哨兵写入,断言下游读到的值与哨兵无关(即下游确实各自 clone/cat 出了新张量,见 megatron/data.py:517-544)。

L5 的哨兵测试是同一思路的镜像:L3a 要证明"共享的东西没被写",L5 要证明"被复用的东西没被读"。两者都属于别名(aliasing)类缺陷,只能靠显式毒值断言暴露,不可能靠端到端曲线发现。

5.2 测试矩阵

5.2.1 CI 内(CPU-only,transfer_queue 被 stub)

Relax CI 在 Python 3.10/3.11/3.12 上跑 CPU-only 且 stub 掉 transfer_queue
所以以下全部设计为不依赖真实 TQ:

测试 内容
test_dataplane_accounting_* 字节会计对已知 payload 给出精确值
test_multimodal_normalize_roundtrip decode(encode(x)) 逐字节全等;覆盖 fp32/bf16/fp16/uint8、0/1/多图、None 样本、非张量 raw 值
test_multimodal_dedup_roundtrip N 份相同 + 部分不同的混合组,还原后逐样本 torch.equal
test_dedup_hash_negative 仅 1 字节不同的两张图不被误判为同一份
test_host_staging_reuse 池命中时 register_buffer 调用次数不随 get 轮数增长(mock 计数);未命中时正确回落 per-call 路径
test_host_staging_no_stale_read §5.1 的毒值哨兵:缓冲归还后写入毒值,断言消费侧数据不变
test_batch_bytes_budget 同一批 key 在给定字节预算下的切批数符合预期;预算为 0 时行为与 BATCH_SIZE_LIMIT 完全一致
test_config_validation §4.7 兼容矩阵每一行,含 4 个"拒绝启动"
test_capability_probe_parse 用假的 /sys 树驱动探测逻辑的各分支
test_fallback_state_machine mock tq.init,断言失败路径ray.kill 命名 actor 再重试

test_fallback_state_machine 是这份测试计划里最关键的一个:它锁住的是一个会
表现为"训练卡住不动、无任何报错"的故障(F10),

5.2.2 单机真实 TQ

测试 内容
四路径一致性 SimpleStorage/ZMQ、Mooncake/TCP、Mooncake/RDMA、Mooncake/RDMA+GDR 各跑同一 payload,两两 bytes 全等
快路/慢路归属断言 断言 normalize=on 后最大字段确实走 tensor 路径
档位覆盖 1 / 4 / 16 / 64 / 256 / 512 MiB
dtype / shape fp32/bf16/fp16/uint8/int64;连续与非连续;jagged 长度分布不均

"快路/慢路归属断言"必须做成显式断言而不是靠时延推断。否则一次无意的类型回退
(例如某个字段被重新包成 dict)会静默丢掉 L2 的全部收益,而 benchmark 只会显示
"好像没那么快了"。

5.2.3 多机(≥2 节点,真实 RDMA)

这一节直接对应任务书的六项行为要求:

场景 期望 对应任务书
正常跨机 put/get bytes 全等 逐字节一致
连接建立 handshake 成功;P2PHANDSHAKE 与显式 metadata 地址两种都通 连接建立
背压 生产快于消费时 put 阻塞而非 OOM;staging buffer 不无界增长 背压
超时 对端不响应时按 --tq-init-timeout-seconds 报错退出,不挂死 超时
断连 传输中拔链路 / kill 对端 → 报错并按 §4.8 规则失败(不降级) 断连
重试 注入瞬时失败,验证 MAX_RETRIES=3 生效且不产生重复数据 重试
清理 job 结束后无残留 master 进程、无残留命名 actor、MR 全部注销 清理
required 不降级 人为屏蔽 RDMA 设备 → 非零退出,日志给出具体失败检查项 自动回退
auto 一致降级 同上但 mode=auto全部节点一致落到 TCP,指标记录 reason 自动回退
多租户 同集群两个 job 同时启用 mooncake,互不杀 master F8 回归

5.2.4 端到端

Qwen3-VL 小尺寸 + fully-async / hybrid 各跑若干 step:

  • loss / grad-norm / reward 曲线与 SimpleStorage 基线在数值容差内一致;
  • 逐 step 比对 advantageslogprobs 等文本字段逐字节相同;
  • 多模态字段抽样比对 pixel_values 逐字节相同。

5.3 性能:度量口径

  • 有效字节 = 消费侧实际重建出的 payload 字节数,不含协议头,不含去重省掉的
    重复份
    。这样定义使得 L3b 不会因为"少传了字节"而虚增带宽。
  • 有效带宽 = 有效字节 / 端到端时间(producer 开始 put → consumer 完成重建
    TensorDict)。
  • 延迟 = 同一区间的 wall clock,报 p50 / p95 / p99。
  • 每档每配置 ≥3 次(任务书要求),报中位数与全距;同拓扑、同镜像、同
    n_samples_per_prompt
  • 每份报告必须附 §4.1 的生效配置摘要一行,证明跑的是声称的那条路径。

档位与拓扑:

  • payload 档位:1 / 4 / 16 / 64 / 256 / 512 MiB
  • 真实档位:--image-max-token-num 取小/中/大三档下的实际 batch 体积。
  • 拓扑:2 节点(最小可证跨机)与 ≥4 节点(证明不随规模退化)。
  • 同时记录 CPU 占用、host 内存峰值、GPU 内存峰值(GDR staging 吃显存,
    --tq-gdr-staging-mb 默认 1024 意味着每 client 1 GiB)。

5.4 两组对照:A 组回答任务书,B 组做归因

A 组(对外结论):当前 main 默认配置 vs 全开配置,逐档给出带宽/延迟。
这一组回答"这个改动让流水线快了多少"。

B 组(归因):逐层累加,每一步只开一个开关。层序按依赖排,L5 在 L4 之前。

配置 backend normalize dedup(L3b) staging(L5) per-rank(L4) share-proc(L3a)
C0 基线 Simple off off off off off
C1 Mooncake/TCP off off off off off
C2 Mooncake/RDMA off off off off off
C3 Mooncake/RDMA on off off off off
C4 Mooncake/RDMA+GDR on off off off off
C5 Mooncake/RDMA+GDR on on off off off
C6 Mooncake/RDMA+GDR on on on off off
C7 Mooncake/RDMA+GDR on on on on off
C8 Mooncake/RDMA+GDR on on on on on

从这张表可以直接读出:

差值 含义
C1 − C0 换 KV backend 本身的收益(不含 RDMA,这一项常被误记为 RDMA 收益)
C2 − C1 RDMA transport 的净收益(同 backend,仅换 protocol)——任务书真正要的那个数,也是唯一能证明 L1 有效的数
C3 − C2 快路收益:去掉 pack_into 全量 memcpy + 连续化
C4 − C3 GDR 收益,且只有 normalize=on 时才可能非零(F2)
C5 − C4 put 侧去重收益,随 n_samples_per_prompt 线性放大
C6 − C5 L5 收益:MR 常驻 + get 缓冲复用。此时 L4 仍关闭,读侧只有 1 个 rank,所以这个差值是未被放大的单份收益——它是 L4 是否值得打开的判据
C7 − C6 L4 扇出收益:去掉 pickle 与 CP→TP→PP 三跳广播的串行路径,代价是读侧字节量与固定开销乘 TP×PP×CP。必须同时报 C7 − C5(即不带 L5 直接开 L4),两者之差就是 L5 作为 L4 前置的实际价值
C8 − C7 L3a 收益。不会体现在数据面带宽上(C5 已经把重复字节去掉了),只体现在 rollout 侧 perf_detail/rollout/image_processor_time 与端到端 step 时间上——因此这一项必须单独用这两个指标报,不能混进带宽结论

C8 − C7 是这张表里唯一不属于数据面的一项,把它放进同一张表是为了让读者看到"字节省掉了不等于 CPU 省掉了":L3b 让 8 份字节变 1 份,但 8 次 processor 调用依然发生,只有 L3a 能去掉。

C6/C7 两档必须记录 dataplane_mr_register_totaldataplane_get_alloc_seconds,否则无法区分"L4 没有收益"和"L4 的收益被读侧固定开销吃掉了"——这两种结论对后续工作的指向完全相反。若 L5a/L5b 的上游 PR 尚未合入,C6 无法测得,此时必须在报告中显式标注 C7 是"L5 缺位下的 L4",不得把该数据当作 L4 的最终结论。

5.5 报告规则

如果 C2 − C1 单独不足 +20%,如实写出这个数,并说明总收益的构成。

一个诚实的"RDMA 净收益 12%,配合 L2 后 31%"比一个不可复核的"提升 40%"更有价值:
后者在别人的集群上复现不出来时,没人知道是哪一层没生效。而在一个含 8× 冗余的基线
之上,不做归因的"提升 20%"不构成对任务书的回答。


附录:关键代码位置索引

TransferQueue,commit 58054a33834aadbcf76aacd6b1e32e25c030f2c9

位置 内容 对应
transfer_queue/config.yaml:20-103 backend 选择、Mooncake 全部字段、Yuanrong enable_rdma F1
storage/clients/mooncake_client.py:162-200 putisinstance(value, torch.Tensor) 分流;GDR 只收 tensor_keys F2
同上 :252-261 快路:注册调用方已有缓冲(不额外拷贝)→ upsert → finally 注销:670-681 _preprocess_tensors_for_put 的 D2H 与 .contiguous() F2, F5, L5
同上 :263-293 慢路 encode + pack + 新 region register → finally 注销:266 上游 TODO "switch to a pre-registered buffer" F2, F5, L5
同上 :295-402 getdtype is not None 分流 F2
同上 :403-417 get 张量:allocate_empty_tensors 全量新分配 → register → get_intofinally 注销 F5, L5
同上 :487-507 get 非张量:同构,:502-506 register/unregister,再 batch_decode_from F5, L5
同上 :683-689 _register_all_buffers / _unregister_all_buffers,逐区间调 store.(un)register_buffer F5, L5
同上(切批) BATCH_SIZE_LIMIT 在 put :178,183 与 get :377,388 均按 key 数切批,每批一轮注册 F5, L5c
同上(常量) BATCH_SIZE_LIMIT=400MAX_BATCH_WORKER_THREADS=4MAX_SERIAL_WORKER_THREADS=4MAX_RETRIES=3 F5, §4.8
utils/tensor_utils.py:28-95 allocate_empty_tensors:按 dtype 分组,:84 每组一次 torch.empty(total_elements),视图经 as_strided 切出 → 注册区间数 = dtype 种类数 F5, L5b
utils/tensor_utils.py:133-165 merge_contiguous_memory:按地址排序合并相邻区间,"reduce register_buffer overhead";不相邻则合不掉 F5, L5
utils/mooncake_utils.py:90-101 GdrStaging docstring:One cudaMalloc buffer, registered once for the process lifetimeavoids repeated cudaMalloc/register overheadL5 的上游先例 F5, L5a
utils/mooncake_utils.py:111-143 lazy_initcudaMalloc + 一次 register_buffer:143 L5a
utils/mooncake_utils.py:147-151 close:唯一的 unregister_buffer 调用点(进程级生命周期) L5a, §4.8
utils/serial_utils.py:196-220 张量叶子零拷贝 encode §4.3
utils/serial_utils.py:420-429 pack_into 全量 memcpy F2
storage/managers/base.py:456-471 _generate_keys{idx}@{field} F4
同上 :474-494 _generate_values:nested → unbind(),其余逐行取 F3, L2
同上 :531-607 _merge_tensors_to_tensordictas_nested_tensor(layout=jagged) F7, L2
storage/bootstrap/mooncake_bootstrap.py:55-65 pkill -f "[m]ooncake_master" F8
interface.py:83-87 _maybe_create_tq_storage 失败只 logger.error F10
interface.py:109-118 while conf is None: sleep(1) F10
interface.py:152, 181-195 init 顺序:命名 actor → storage → store_config F10
interface.py:198-215 metrics 注册硬判定 == "SimpleStorage" F9

Relax(main):

位置 内容 对应
relax/core/controller.py:145-192 _initialize_data_system,唯一 tq.init 调用点,只构造 SimpleStorage 字段 F1, L1
relax/core/controller.py:871-875 tq.close() §4.8
relax/utils/arguments.py:34-65 _MIN_TQ_VERSION = "0.1.10.dev0" 与 pin 的 commit R10
relax/utils/arguments.py:206-245 add_transfer_queue_arguments,现有 5 个开关 L1
同上 :213-226 --per-rank-fetchdefault=False;help 里的"rollout_routed_experts 时自动禁用"与"建议 --num-data-storage-units >= TP world size" L4, §4.5
relax/utils/types.py:20 multimodal_train_inputs: dict[str, Any] | None F3
relax/utils/utils.py:153-154 组装为 list[dict] F3
relax/utils/utils.py:276-277 显式不张量化 → NonTensorStack F3
relax/engine/rollout/data_source.py:20-36 _shallow_copy_sample:组内按引用共享 multimodal_inputs,注释断言 multimodal_train_inputs 是被 set 而非被 mutate F4, L3a
同上 :218-228 组展开:n_samples_per_prompt_shallow_copy_sample F4, L3a
relax/engine/rollout/sglang_rollout.py:207-261 _run_image_processor:产出 multimodal_train_inputs 的实际开销所在 F4, L3a
同上 :264-294 _encode_multimodal_inputs:base64,喂 SGLang HTTP 请求体(不是 processor) L3a(对照)
同上 :297-325 per-sample generate:322-325 无组级去重地调 processor F4
同上 :361-369 _pre_encoded_mm 的消费方式(L3a 缓存挂载的参照实现) L3a
同上 :649-658 已有先例:按 multimodal_inputs is 做的组级 base64 去重 F4, L3a
同上 :621-626 perf_detail/rollout/image_processor_time 指标 L3a, §5.4
relax/backends/megatron/data.py:517-544 多模态张量拼接:只 append + cat/pad_and_flatten,不改源 dict(L3a 只读性依据) L3a, R4b
relax/utils/data/stream_dataloader.py:489-517 _encode_multimodal_inputs(spec+tensor 拆分,现成能力) F7, L2
同上 :520-553 _broadcast_multimodal_inputs F6
同上 :566 per_rank_fetch: bool = False F6, L4
同上 :597-605 per_rank_fetch docstring:all TP/PP broadcasts are skippedTrades a single rank-0 pickle + one NCCL bcast for N parallel ZMQ deserialises F6, L4, §4.5
同上 :627 注释里的真实拓扑量级 TP2/PP2/CP8 §4.5
同上 :633-639 per_rank_fetchshould_fetch = True(每 rank 各自 get_meta+get_data,依赖 sampler 按 (partition_id, task_name, dp_rank, batch_index) 缓存保证各 rank 拿到相同 sample id) F6, L4
同上 :640-649 默认路径的拉取 rank 判定:colocate 下 (tp,pp,cp)==(0,0,0),fully-async 下 (tp,cp)==(0,0) F6
同上 :718-741 has_multimodal and not per_rank_fetchdel td[...] + NCCL F6, L4
同上 :750-782 CP / TP / PP 三跳 broadcast_object_list,全部位于 if not per_rank_fetch: 之内 F6, L4, §4.5
同上 :819-820 per_rank_fetch 路径重建多模态 F6, L4
relax/backends/megatron/actor.py:2160 per_rank_fetch = args.per_rank_fetch and "rollout_routed_experts" not in data_fields — 会静默禁用 L4 L4, §4.5
同上 :2238-2250 _put_data_to_transfer_queuetp_rank==0 and is_pipeline_last_stage() and cp_rank==0写侧每 DP 组只有 1 个 rank put F6, L5

术语:

术语 含义
MR (Memory Region) RDMA 中向 HCA 注册的内存区域,注册需 pin page,是慢操作
GDR (GPU Direct RDMA) 网卡直接读写显存,绕过 host 内存
staging buffer GDR 路径下预注册的中转缓冲,默认 1 GiB/client
常驻注册池 L5a 引入的宿主侧对照物:进程级预分配并一次性注册的宿主缓冲集合,跨 step 复用,对应上游 GdrStaging 在 CPU 路径的缺位
快路 / 慢路 MooncakeStore 中 isinstance(value, Tensor) 为真 / 为假的两条 put 分支
jagged NestedTensor PyTorch 不规则嵌套张量,layout=torch.jagged,TQ 用它还原逐样本张量
partition TQ 中一次 put/get 的批次范围,也是 L3b 去重的作用域
DP 分片域 同一 data-parallel rank 负责的 global_index 集合
cohort 同一 DP 组内的全部 TP/PP/CP rank。写侧只有其中 1 个 rank put,读侧在 L4 打开后全体各 get 一次——这就是"放大轴是 TP×PP×CP 而不是 DP"的含义

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions