跳转至

Memory-Semantic TransferQueue

导言

能否把 MOV A, B 的直觉扩展到跨 device、跨节点甚至跨超节点:应用只给出源地址和目标地址,系统自动选择 HCCS、RDMA 等路径,把大 tensor 直接搬到消费者附近?可以把它做成 verl TransferQueue 的可插拔数据面,但不应把这件事直接称为“MTE 跨超节点”。MTE 是 Ascend 内存层级和部分 HCCS 路径上的搬运引擎;真正向应用提供远端地址、单边 put/get、注册和完成语义的是 SHMEM、MemFabric 一类软件层;样本字段、ready、consumed、staleness 和恢复仍应由 TransferQueue 管理。

小黑用地址控制远端 tensor 搬运

认知示意图。控制面只传很小的远端引用;大 tensor 由数据面按拓扑选择 MTE、RDMA 或其它路径。图中不是一条真实硬件指令,而是分层系统接口。

结论

先给出判断:

  1. 有价值:异步、多模态 RL 恰好具备大 tensor、多生产者、多消费者、字段分批 ready、设备分离等特征,适合用远端内存引用替代 controller 中转大对象。
  2. 有可行性:TransferQueue 已经有 StorageManager 插件面;verl 已经通过 async KV 接口按字段写入和读取。最自然的做法是增加内存语义 backend,不改 trainer 上层数据契约。
  3. 当前没有完整实现:截至本文固定版本,verl 和 TransferQueue 中没有显式的 SHMEM、MemFabric 或 MTE backend。TransferQueue 的 Yuanrong NPU transport 是最近的已有路径和首选 baseline,但不能把它等同于本文设计。
  4. MTE 不应上浮成业务接口:上层只看到 RemoteTensorRef、put/get、completion 和 lease;MTE 仅在 HCCS 拓扑与设备约束满足时成为执行引擎,跨 HCCS group 或跨节点时通常要切到 RDMA 等路径。
  5. 收益不是无条件的:大、已在 NPU 上、会复用、可按字段消费的 tensor 最可能获益;小 metadata、压缩图片/视频、collective、频繁注册的小块和拓扑不支持的路径未必获益。
  6. 正确性比带宽优先:必须满足“copy 完成后才能发布 ready”“任何 pending copy 或 live lease 存在时 pool slot 不得复用”“重试不能混用 generation”。

因此,建议把它定义为 Memory-Semantic TransferQueue:TransferQueue 负责 RL 数据系统语义,SHMEM/MemFabric 负责内存语义,MTE/RDMA 负责物理搬运。

语义分层

MTE:搬运引擎

华为 Ascend AI Core 的 MTE(Memory Transfer Engine)管理不同内存层级之间的数据读写,并可在搬运时完成格式或数据类型转换。典型路径包括:1

  • MTE1:L1 到 L0A/L0B。
  • MTE2:Global Memory 到 L1、L0 或 Unified Buffer。
  • MTE3:Unified Buffer 到 Global Memory 或 L1。

其核心价值不是“远端地址”,而是把数据搬运放进独立队列,与 Vector/Cube 计算形成流水。例如处理一个大 tensor 时,不必先同步搬完全部数据;可以预取 tile_{i+1}、计算 tile_i、回写 tile_{i-1}

def mte_tiled_compute(src_gm, dst_gm, tile_count, tile_bytes):
    input_events = [Event(), Event()]
    output_events = [Event(), Event()]
    ub_buffers = [UnifiedBuffer(tile_bytes), UnifiedBuffer(tile_bytes)]

    for tile_id in range(tile_count):
        slot = tile_id % 2
        src_offset = tile_id * tile_bytes
        dst_offset = tile_id * tile_bytes

        mte2_copy_async(
            dst=ub_buffers[slot],
            src=src_gm.byte_range(src_offset, tile_bytes),
            completion=input_events[slot],
        )
        wait_event(input_events[slot])

        vector_or_cube_compute_inplace(ub_buffers[slot])

        mte3_copy_async(
            dst=dst_gm.byte_range(dst_offset, tile_bytes),
            src=ub_buffers[slot],
            completion=output_events[slot],
        )
        wait_event(output_events[slot])

    return dst_gm

MTE 五联机制图

自绘五联机制图。MTE2/MTE3、计算队列、event 和 tile 生命周期共同构成局部流水;跨节点寻址、注册、租约与恢复不是 MTE 单独提供的能力。

不要把 MTE 等同于跨超节点 MOV

“源地址、目标地址、字节数”很像一条跨设备 MOV,但完整系统至少还要回答:地址在哪个进程和设备上有效、内存是否注册、链路是否 HCCS 直连、何时完成、源缓冲区何时可复用、节点失败后引用是否仍有效。MTE 只回答其中的搬运执行问题。

内存语义:地址后的位置

传统消息语义往往让应用显式组织:

节点 → 设备 → 内存类型 → 本地地址 → 通信协议 → 目标地址

内存语义则希望把位置隐藏在一个可传播的地址或引用后:

GVA / symmetric address + put/get/xcopy + completion

UnifiedBus 1.0 把内存语义与统一寻址作为超节点互联的重要方向;官方超节点报告进一步区分了适合大块异步搬运的 DMA 和适合细粒度访问的 Load/Store。2 但从工程接口看,更直接的落点是:

  • Ascend SHMEM:提供对称内存、PE/rank 语义、单边 put/get、signal 和 stream 接口。版本固定的文档明确写明:跨机 put/get 可在 HCCS 连通时使用 MTE,否则在条件满足时使用 RDMA。3
  • MemFabric:提供 Global Virtual Address、异构内存池、注册、copy、batch copy、stream 和 wait;底层可选择 SDMA、Host RDMA、Device RDMA、共享内存或其它 engine。4

这正是 MOV A, B 类比真正成立的地方:应用不再手工拼接每段传输路径,但完成、生命周期和故障仍然显式。

def memory_semantic_copy(runtime, src_tensor, dst_rank, dst_bytes):
    pool = runtime.symmetric_pool()
    dst_slot = pool.allocate(owner_rank=dst_rank, nbytes=dst_bytes)
    remote_address = runtime.global_address(dst_slot)
    route = runtime.select_route(
        src_device=src_tensor.device,
        dst_rank=dst_rank,
        preference=["MTE_HCCS", "DEVICE_RDMA", "HOST_COPY"],
    )

    request = runtime.putmem_on_stream(
        dst=remote_address,
        src=src_tensor,
        nbytes=dst_bytes,
        route=route,
        stream=runtime.copy_stream(),
    )
    runtime.wait(request)
    runtime.publish_completion(remote_address, request.generation)
    return remote_address, request.generation

内存语义五联机制图

自绘五联机制图。GVA/对称地址是控制对象,远端 pool 是物理对象;MTE 或 RDMA 是按拓扑选择的执行路径,wait/signal 决定可见性。

MemFabric 全局虚拟地址与 xcopy

华为昇腾官方技术文章 Figure 2。MemFabric 用全局虚拟地址组织 CPU/NPU 内存池,写入、复制和读取通过 xcopy 作用于地址空间。该图解释软件语义,不代表 verl 已集成此路径。

TransferQueue:字段生命周期

TransferQueue 不是普通 FIFO,也不只是 rollout token queue。它把 RL 样本组织成“行是 sample、列是 field”的数据系统:

  • key:定位 trajectory、session、partition 和 field。
  • tag:表达 ready、consumed、stale、policy version 等状态。
  • metadata:告诉 worker 应该取哪些样本和字段。
  • storage backend:保存真实 tensor 或远端引用。

因此三层的职责不能混在一起:

负责什么 不负责什么
TransferQueue key、field、ready、采样、消费记录、重试与恢复策略 不固定物理链路
SHMEM / MemFabric 地址、注册、put/get、completion、topology route 不理解 RL policy version 与样本分布
MTE / RDMA / Host copy 真正移动字节 不决定样本何时 train-ready

最重要的跨层不变量是:

\[ \text{ready(field, generation)} \Rightarrow \text{bytes-complete} \land \text{generation-match} \land \text{buffer-live} \]

如果只发布远端地址而不等待 copy 完成,TransferQueue 的 ready 就会比数据本身更早可见;这会把偶发脏读伪装成训练随机性。

当前代码

verl 调用面

本文检查的 verl 上游版本固定为 983cb0f24443f87b3d161fad318445130a620b07。当前代码已经具备面向新 backend 的应用形态:5

  1. requirements.txtrequirements-npu.txt 固定 TransferQueue==0.1.8
  2. verl/trainer/main_ppo.py 在 v1 trainer 启动时设置 config.transfer_queue.enable=True 并调用 tq.init(...)
  3. agent_loop_tq.py 把 response、mask、logprob、reward 等字段通过 async_kv_batch_put 写入 TQ。
  4. 多模态路径移除原始 multi_modal_data,保留模型处理后的 multi_modal_inputs,后者正是更可能适合 device-direct 传输的对象。
  5. replay_buffer.py{uid}_{session_id}_{index} 一类 key 和 tags 组织样本,真实 values 交给 storage unit。
  6. transferqueue_utils.py 通过 TransferQueueClient.async_get_data/async_put 读写真实 tensor,agent loop 写侧还直接使用 batch KV API。

这意味着 trainer 不需要知道底层用了 MTE 还是 RDMA。它只应继续看到:

await tq.async_kv_batch_put(
    keys=keys,
    fields=tensor_fields,
    tags=tags,
    partition_id="train",
)
tensor_fields = await tq.get_client().async_get_data(metadata)

verl 当前没有 MTE backend

在固定版本中没有找到显式的 SHMEMMemFabricUnifiedBusMTE storage backend。ppo_trainer.yaml 公开列出的主要是 SimpleStorage 和实验性的 MooncakeStore 配置,因此“verl 已经用上内存语义”是不成立的。

TransferQueue 插件面

TransferQueue 当前 main 固定为 b75d570d88c50bbfcbe2171baa727fadd7216f76,verl 使用的 v0.1.8 固定为 35bcf198ac087f0a0e41a69a897f48929c313d22。核心扩展点已经存在:6

class StorageManager:
    def put_data(self, data, metadata):
        raise NotImplementedError

    def get_data(self, metadata):
        raise NotImplementedError

    def clear_data(self, metadata):
        raise NotImplementedError

StorageManagerFactory 负责按名称注册和创建 backend;KVStorageManager 负责把 field 与 index 展平为 storage key,并保存 schema、shape、dtype 与 backend metadata。这里就是 Memory-Semantic backend 的正确接缝。

它不应该侵入:

  • PPO/GRPO 算法代码;
  • rollout engine;
  • reward model;
  • trainer 对 TensorDict / DataProto 的语义;
  • sampler 对 ready/staleness 的判断。

Yuanrong 已有路径

当前 TransferQueue 已包含 Yuanrong backend,并把 NPU tensor 路由到 NPUTensorKVClientAdapter。官方文档描述了 host-host TCP/RDMA、device-device HCCS、remote H2D/D2H 和 per-host worker 等能力。7

因此,落地顺序不应从“写一套新 MTE 通信库”开始,而应先回答:

  1. Yuanrong 在目标 A3/超节点拓扑上是否已经消除 controller bounce?
  2. enable_yr_npu_transport 和 remote H2D/D2H 能否覆盖 rollout→trainer 的实际路径?
  3. 对相同 tensor field、相同 batch、相同拓扑,Yuanrong 与新 backend 的增量差异是什么?

TransferQueue v0.1.8 会把调用方配置与自身默认配置合并,因此通过额外 override 选择 Yuanrong 可能可行;但 verl 当前 YAML 没有完整暴露 Yuanrong 参数,本文也没有完成部署侧 E2E 验证,必须把它标记为待测工程推断。

后端设计

架构与对象

建议新增 MemorySemanticStorageManager 和配套 MemorySemanticStorageClient。对上仍实现 put_data/get_data/clear_data;对下可先接 Ascend SHMEM,后续再接 MemFabric,或反过来。

Memory-Semantic TransferQueue 架构五联图

自绘架构五联图。地址只进入 TQ 控制面,大 tensor 进入 backend 自有 pool;completion 先于 ready,lease ack 先于 slot 复用。

关键对象如下:

对象 物理位置 生产者 消费者 生命周期
field_tensor rollout/reward/logprob worker NPU 上游 stage storage backend 创建到安全 copy/adopt 完成
pool_slot 注册的对称/global HBM 或 DRAM pool backend allocator transport 与远端消费者 分配到所有 lease 结束
RemoteTensorRef controller/meta storage backend sampler 与 consumer backend publish 到 clear/stale
completion copy stream/event/request transport publisher/reclaimer submit 到 wait 完成
lease controller/backend metadata consumer acquire reclaimer/recovery pull 到 ack/超时处理
local_dst trainer NPU consumer allocator trainer stage pull 到训练释放

远端引用建议至少包含:

@dataclass(frozen=True)
class RemoteTensorRef:
    owner_pe: int
    global_address: int
    offset: int
    nbytes: int
    shape: tuple[int, ...]
    dtype: str
    generation: int
    checksum: str
    memory_kind: str

global_address 不是裸指针的永久承诺。必须用 generation 防止 pool slot 复用后的 ABA 问题,并用 shape/dtype/nbytes 防止控制面与物理字节布局不一致。

写入伪代码

默认使用 copy-in

  • producer tensor 可以按原 allocator 创建;
  • backend 把它复制进自己管理的 registered pool;
  • completion 后 producer 就能较早释放原 tensor;
  • backend 对 pool slot 的所有权稳定,恢复和回收更容易。

只有 producer tensor 本来就由同一个 backend pool 分配时,才使用 adopt,避免一次本地 copy。

def put_field(self, field_key, tensor, policy_version, mode="copy_in"):
    validate_contiguous_or_packable(tensor)
    generation = self.generation_store.next(field_key)
    checksum = checksum_metadata_and_optional_payload(tensor)
    slot = None
    request = None

    try:
        if mode == "adopt":
            slot = self.pool.lookup_owned_tensor(tensor)
            if slot is None:
                raise ValueError("adopt requires a tensor allocated from this pool")
            self.pool.pin(slot, generation)
        elif mode == "copy_in":
            slot = self.pool.allocate(
                nbytes=tensor.nbytes,
                alignment=self.transport.required_alignment(),
                generation=generation,
            )
            route = self.router.select(
                src_device=tensor.device,
                dst_memory=slot.memory_kind,
                candidates=["MTE_HCCS", "DEVICE_RDMA", "HOST_COPY"],
            )
            request = self.transport.copy_into_pool_async(
                dst_slot=slot,
                src_tensor=tensor,
                route=route,
                stream=self.copy_stream,
            )
            self.transport.wait(request)
        else:
            raise ValueError(f"unsupported put mode: {mode}")

        remote_ref = RemoteTensorRef(
            owner_pe=slot.owner_pe,
            global_address=slot.global_address,
            offset=slot.offset,
            nbytes=tensor.nbytes,
            shape=tuple(tensor.shape),
            dtype=str(tensor.dtype),
            generation=generation,
            checksum=checksum,
            memory_kind=slot.memory_kind,
        )

        self.ref_store.publish(
            field_key=field_key,
            policy_version=policy_version,
            remote_ref=remote_ref,
        )
        self.controller.mark_ready(
            field_key=field_key,
            generation=generation,
            policy_version=policy_version,
        )
        return remote_ref
    except Exception:
        if request is not None:
            self.transport.cancel_or_drain(request)
        if slot is not None:
            self.pool.rollback(slot, generation)
        self.controller.mark_failed(field_key, generation)
        raise

顺序中不能交换的两步是:

wait(copy completion) → publish RemoteTensorRef / mark_ready

先 ready、后 completion 即使平均情况下能工作,也会在拥塞、跨节点 fallback 或重试时产生时间窗。

读取与回收

读取应该按 field batch 执行,使 transport 有机会合并注册查询、route 判断和 copy submission:

def get_fields(self, consumer_id, field_keys, expected_policy_version):
    refs = self.ref_store.batch_get(field_keys)
    leases = []
    destinations = []
    requests = []

    try:
        for field_key, remote_ref in zip(field_keys, refs):
            validate_ref_schema(remote_ref)
            validate_policy_version(field_key, expected_policy_version)

            lease = self.lease_store.acquire(
                field_key=field_key,
                generation=remote_ref.generation,
                consumer_id=consumer_id,
            )
            destination = allocate_local_tensor(
                shape=remote_ref.shape,
                dtype=remote_ref.dtype,
                device=self.local_device,
            )
            route = self.router.select(
                src_memory=remote_ref.memory_kind,
                src_owner_pe=remote_ref.owner_pe,
                dst_device=self.local_device,
                candidates=["MTE_HCCS", "DEVICE_RDMA", "HOST_COPY"],
            )
            request = self.transport.get_async(
                dst_tensor=destination,
                src_ref=remote_ref,
                route=route,
                stream=self.copy_stream,
            )

            leases.append(lease)
            destinations.append(destination)
            requests.append(request)

        self.transport.wait_all(requests)

        for field_key, remote_ref, destination, lease in zip(
            field_keys, refs, destinations, leases
        ):
            validate_generation(field_key, remote_ref.generation)
            validate_destination(destination, remote_ref)
            validate_optional_checksum(destination, remote_ref.checksum)
            self.lease_store.ack(lease)
            self.controller.mark_consumed(
                field_key=field_key,
                generation=remote_ref.generation,
                consumer_id=consumer_id,
            )

        return dict(zip(field_keys, destinations))
    except Exception:
        for request in requests:
            self.transport.cancel_or_drain(request)
        for lease in leases:
            self.lease_store.release_failed(lease)
        raise

清理不能只看 TQ 的逻辑 consumed,还要看物理传输与 lease:

def clear_fields(self, field_keys):
    refs = self.ref_store.batch_get(field_keys)

    for field_key, remote_ref in zip(field_keys, refs):
        self.controller.mark_draining(field_key, remote_ref.generation)
        self.transport.wait_generation(
            owner_pe=remote_ref.owner_pe,
            generation=remote_ref.generation,
        )
        self.lease_store.wait_until_no_live_lease(
            field_key=field_key,
            generation=remote_ref.generation,
        )
        self.ref_store.unpublish(field_key, remote_ref.generation)
        self.pool.free_by_address(
            global_address=remote_ref.global_address,
            generation=remote_ref.generation,
        )
        self.controller.mark_cleared(field_key, remote_ref.generation)

一致性与故障

异步 RL 中,传输正确不等于训练语义正确。至少还要维护:

  • policy version:拒绝把旧 policy rollout 混进当前 batch,或显式记录允许的 staleness。
  • sample generation:同一个 key 重试时,新旧引用不能混用。
  • consumer lease:一个字段可能被 old logprob、ref logprob、critic 和 actor 多次读取。
  • checkpoint boundary:device pool 通常不是 durable storage;恢复时要么重建 ref,要么把 pending 样本标记 stale。
  • 幂等性put(field_key, generation) 重放不能生成两个都可见的版本。
  • backpressure:pool 满时应阻塞、spill 或拒绝新样本,不能静默改变采样分布。

最危险的不是传输失败,而是看似成功

节点重启、slot 复用或超时重试后,一个数值相同的 GVA 可能已指向新内容。没有 generation、lease 和完成状态时,错误常常表现为 reward、logprob 或 advantage 的小概率异常,而不是清晰的通信报错。

RL 场景

异步流水

TransferQueue 已经把 rollout、reward、old logprob、ref logprob、critic 和 actor update 的字段生命周期拆开。内存语义进一步减少“字段已经逻辑 ready,但 tensor 仍要回到 controller/host 再发出去”的路径。

潜在收益包括:

  1. 减少 controller RSS:controller 只保存 key、ref、tag 和 version。
  2. 减少序列化与 bounce:tensor 不必反复变成 Python/Ray 大对象并经 host 中转。
  3. 更早释放 producer tensor:copy-in 完成后,rollout worker 可释放原对象,pool 独立持有数据。
  4. 按字段传输:trainer 只拉取当前阶段需要的 response、mask、logprob、reward 或 multimodal tensor。
  5. 传输与计算重叠:下一个 field 的 get 可在当前 field 计算时进入独立 copy stream。
  6. 拓扑自适应:同超节点走 HCCS/MTE,跨域走 RDMA,功能上不要求上层改代码。

但异步也放大两个代价:

  • producer/consumer 速率不匹配会把压力转移到 registered pool;
  • live lease 变多会延迟回收,pool 容量不能只按单个 batch 估算。

多模态负载

多模态 RL 的对象要分层看:

对象 是否适合内存语义 原因
token ids、mask、position ids 视大小而定 结构规则,但短序列可能太小
multi_modal_inputs 中的 pixel/vision tensors 适合 已 tensor 化、体积大、可直接送模型
image/video encoder hidden states 很适合 已在设备上、体积大、可能复用
logprobs、values、advantages 适合 规则 tensor、会被不同 stage 消费
raw JPEG/PNG/MP4 bytes 通常不优先 压缩对象适合对象存储/CPU path,解码后再进入 device data plane
tokenizer 配置、字符串、采样 metadata 不适合 对象小,控制成本高于传输收益

verl 当前 TQ agent loop 会保留模型处理后的 multi_modal_inputs,而不是继续携带原始 multi_modal_data。这使得多模态路径比纯文本短序列更可能跨过收益临界点。

不占优边界

一次远端搬运可以粗略写成:

\[ T_{\text{move}} = T_{\text{lookup}} + T_{\text{register-amortized}} + \frac{\text{bytes}}{B_{\text{effective}}} + T_{\text{completion}} + T_{\text{queue}} \]

只有当减少的序列化、host bounce 和额外 copy 大于新增的 lookup、注册、完成和排队成本时,方案才有端到端收益。

以下场景应继续走别的路径:

  • 小对象:几十字节到小 KB 的 metadata 直接放 controller/meta store。
  • collective:梯度 AllReduce、参数 broadcast、ReduceScatter 等优先使用 HCCL;它们需要 collective 代数语义,不是普通远端 copy。
  • 高频临时注册:每个小 tensor 都注册/注销会吞掉收益,应使用长期 pool 和 registration cache。
  • 不连通拓扑:部分设备组之间不能直接走 MTE,需要 RDMA 或 host fallback。
  • 强 durability:需要跨进程重启长期保存的数据,应落到持久对象存储,而不是只依赖 HBM GVA。
  • 极端碎片化:大量变长小 tensor 会导致 pool 内部碎片,应分 size class、合并字段或 pack。

内存占用也不是凭空下降。更完整的峰值账本是:

\[ \begin{aligned} M_{\text{peak}} ={}& M_{\text{model-state}} + M_{\text{saved-activations}} + M_{\text{operator-workspace}} \\ &+ M_{\text{communication-buffers}} + M_{\text{live-inputs-outputs}} + M_{\text{allocator-reserved-unused}} + M_{\text{runtime-margin}}. \end{aligned} \]

内存语义主要减少 controller host memory 和中间 M_communication-buffers,同时新增固定的 M_pool-reserved、registration metadata 与碎片。它是“转移并有界化”内存,不是消灭内存。

分阶段落地

Phase 0:Yuanrong 基线

  1. 在目标 NPU 集群启用现有 Yuanrong backend。
  2. 验证 NPU tensor、remote H2D/D2H、HCCS 和 RDMA fallback。
  3. 记录 controller RSS、host bounce bytes、pool/worker 内存、E2E samples/s。
  4. 不改 trainer,只用当前 TQ API。

成功标准是证明“现有 backend 能否已经解决主要瓶颈”,而不是先证明新方案更复杂。

Phase 1:单节点 SHMEM 原型

  1. 实现 MemorySemanticStorageManager/Client
  2. 只支持固定 size class 和 NPU tensor。
  3. 先实现 copy_in,暂不实现 adopt
  4. 同节点或明确 HCCS 连通设备间完成 put/get/clear。
  5. 对每个 field 做 shape、dtype、generation 与 checksum 校验。

Phase 2:跨节点与混合路由

  1. 增加 topology probe。
  2. HCCS 连通走 MTE;其它路径走 Device RDMA 或 Host copy。
  3. 增加 batch put/get、registration cache 和多 copy stream。
  4. 增加 timeout、retry、lease recovery 和 checkpoint stale 处理。

Phase 3:verl 异步多模态

  1. 先迁移 multi_modal_inputs、encoder hidden 和长 response/logprobs。
  2. metadata 与小 tensor 保持原后端。
  3. 按 payload size 和 topology 自动分流。
  4. 开启 adopt 前,先让 producer 能从 TQ pool 分配 tensor。
  5. 在 FullAsync、router replay 与 checkpoint 同时开启的组合上做故障注入。

验证计划

基线矩阵

Backend 作用
SimpleStorage 功能与小规模延迟基线
MooncakeStore 分布式内存存储基线
Yuanrong 当前最接近 NPU/HCCS transport 的基线
Memory-Semantic backend 本文提案

测试维度必须同时覆盖:

  • payload:4 KB64 KB1 MB8 MB64 MB、真实 trajectory field;
  • topology:同 device、同节点跨 device、HCCS group 内、跨 group、跨节点/超节点;
  • fan-out:单消费者和 old/ref/critic/actor 多消费者;
  • 负载:纯文本、图像、视频、encoder hidden;
  • 调度:同步 batch、partial async、FullAsync;
  • pool:稳定态、80% 水位、碎片化、backpressure。

指标

不能只看链路带宽:

维度 指标
端到端 samples/s、tokens/s、step time、rollout idle、trainer idle
时延 field ready→consume p50/p95/p99、put/get completion
内存 controller RSS、worker host RSS、NPU peak、pool reserved/used、碎片率
搬运 host bounce bytes、MTE bytes、RDMA bytes、fallback ratio
资源 copy-stream 占用、registration cache hit、queue depth
正确性 stale drop、generation mismatch、checksum mismatch、duplicate consume
恢复 节点退出、copy timeout、pool exhaustion、checkpoint restart

MemFabric 内存语义与消息语义性能对比

华为昇腾官方技术文章 Figure 7。两节点 A3、约 8.57 MB 离散 block 的合成负载中,文中内存语义路径报告了更低总时延和约 190–215 GB/s 的带宽,消息语义约 40–64 GB/s。它说明地址化批量 copy 的潜力,但不是 verl、TransferQueue 或 RL E2E 结果。

verl 的 TransferQueue 文档还报告了 128×H100 多模态后训练约 49.1% E2E gain。8 这证明了“把大 tensor 从单 controller 数据流中拆走”的价值,但仍不能推导出“增加 MTE/SHMEM 后还会再快 49.1%”。新 backend 的收益必须在同一模型、同一 batch、同一精度、同一拓扑和同一 TQ 版本下测量。

先找 crossover,再谈全量迁移

第一轮实验的目标不是追求最大 GB/s,而是画出 payload size × topology 的 crossover 曲线。只有跨过临界点的 field 才进入 memory-semantic path,其余对象继续走原 storage backend。

总结

华为 MTE 语义可以参与 verl 到 TransferQueue 的设计,但位置必须放对:

verl trainer / agent loop
        ↓  key、field、ready、staleness
TransferQueue StorageManager
        ↓  RemoteTensorRef、lease、completion
SHMEM / MemFabric
        ↓  topology route
MTE(HCCS) / RDMA / Host copy

对异步、多模态 RL,这个方向有结构性优势,也具备清晰的代码接缝;但它现在是“可实现的 backend 设计”,不是“当前 verl 已经具备的功能”,更不是一条可以无条件跨超节点执行的 MTE 指令。

最务实的路线是:先用 Yuanrong 建立 NPU baseline,再实现 backend 自有 pool、completion-before-ready、lease 与 generation,最后只迁移超过 crossover 的大 tensor field。这样既能保留 TransferQueue 的 RL 数据语义,也能让底层内存系统决定什么时候用 MTE、什么时候用 RDMA。

参考资料


  1. Ascend C MTE glossary and Ascend C programming guide

  2. Huawei Research and Innovation: UnifiedBus 1.0, Huawei SuperPoD keynote, and 超节点发展报告

  3. CANN SHMEM repository, fixed at 643bfb73b50d88042234d3f5e0c7c9841cdedbbe; see docs/api/stream_api_usage.md and docs/debug/Troubleshooting_FAQs.md

  4. MemFabric Hybrid repository, fixed at 004d9317289fe99bd6bf13def0500b3fa3795ccc; see README.md, smem_bm_def.h, smem_bm.h, and the Python binding. 

  5. verl main_ppo.py, agent_loop_tq.py, replay_buffer.py, and transferqueue_utils.py, fixed at 983cb0f24443f87b3d161fad318445130a620b07

  6. TransferQueue repository, current inspection fixed at b75d570d88c50bbfcbe2171baa727fadd7216f76; verl's pinned v0.1.8 is 35bcf198ac087f0a0e41a69a897f48929c313d22

  7. TransferQueue Yuanrong backend documentation, plus yuanrong_manager.py and yuanrong_client.py at b75d570d88c50bbfcbe2171baa727fadd7216f76

  8. TransferQueue Data System in verl, PR #5401, and RFC #5400. Reported gains are source-scoped to their stated hardware, workload, and baseline. 

评论