How P2P RDMA Weight Transfer Enables Fast In-Loop Model Updates in Miles
P2P RDMA weight transfer in Miles moves updated model parameters directly from trainer ranks to rollout engine ranks via RDMA write operations over a shared CPU-pinned buffer, reducing per-step update latency from hundreds of milliseconds to approximately 30-50 milliseconds.
Miles, an open-source distributed training and inference framework developed by radixark, supports two weight-transfer back-ends for disaggregated training/rollout scenarios: NCCL broadcast (default) and P2P RDMA. The P2P RDMA path eliminates collective synchronization overhead by establishing direct point-to-point connections over high-throughput InfiniBand or RoCE networks. This article examines the six-stage pipeline that makes fast in-loop model updates possible, with direct reference to the implementation in miles/backends/training_utils/weight_update/protocols/.
P2P RDMA vs. NCCL Broadcast: Performance Overview
The default NCCL broadcast path performs an all-gather collective operation across all trainer and engine ranks. While robust, this approach introduces significant latency—typically 200-300 milliseconds per update—due to kernel launch overhead, collective synchronization, and PCIe/NIC traversal through CPU memory.
The P2P RDMA path achieves 5-10× lower latency through three key architectural decisions:
- Direct NIC-to-NIC transfers bypass CPU memory copies and the collective scheduler
- Connection pre-planning eliminates runtime negotiation overhead
- Asynchronous pipelining overlaps data conversion with transfer operations
These optimizations keep rollout generators continuously fed with fresh weights without stalling the training loop.
The Six-Stage P2P RDMA Pipeline
The core mechanism spans six tightly-coupled stages across p2p_transfer_utils.py and p2p.py.
Stage 1: Remote Transfer Planning
Each trainer rank builds a remote transfer plan that maps its data-parallel (DP) replica rank to one or more rollout engine ranks. The planning uses a round-robin heuristic to balance RDMA session count per source rank.
In p2p_transfer_utils.py, the RemoteTransferPlan.plan_p2p() method implements this logic:
# From p2p_transfer_utils.py lines 63-84
def plan_p2p(self) -> List[TransferTaskP2PMeta]:
"""
Build a mapping from this trainer's DP rank to target engine ranks.
Uses round-robin distribution to balance RDMA session load.
"""
# Implementation distributes targets evenly across source ranks
This planning occurs once during initialization, not per-step.
Stage 2: Handshake and Remote Memory Registration
Before any data moves, the trainer queries every target engine for remote memory registration info: address, size, and element size. This metadata is cached in RemoteWeightInfo objects to avoid repeated RPCs.
The query_remote_weight_infos() function handles this (lines 12-35):
# From p2p_transfer_utils.py lines 12-35
def query_remote_weight_infos(
target_engine_ranks: List[int],
transfer_engine: TransferEngine
) -> Dict[int, RemoteWeightInfo]:
"""
Query each engine for its memory registration info and parallelism config.
Results are cached for the lifetime of the training job.
"""
Stage 3: Shared CPU-Pinned Buffer Allocation
A single CPU-pinned buffer is allocated once per trainer process via register_cpu_memory(). All model parameters register with the Mooncake TransferEngine, enabling zero-copy NIC access.
From p2p_transfer_utils.py lines 86-100:
# One-time registration; O(1) memory overhead regardless of engine count
_weight_memory_registry = register_cpu_memory(
shared_params_dict, transfer_engine
)
This design is critical: per-engine buffer allocation would scale linearly with cluster size, but the shared buffer provides constant memory overhead.
Stage 4: Bucketed Weight Conversion
Model weights traverse a conversion pipeline before transfer:
- Gather across tensor-parallel groups
- Convert from HuggingFace layout to SGLang layout
- Accumulate into fixed-size buckets (default 512 MiB)
The _get_transfer_ready_params() method in p2p.py (lines 52-66) implements sharding and bucket readiness detection. Bucketing amortizes per-transfer fixed overhead across many tensors.
Stage 5: Pipelined RDMA Writes
For each engine rank, the trainer:
- Calls
model_replica.load_weights()to copy the bucket into the pinned buffer - Invokes
TransferEngine.batch_transfer_sync_write()for the RDMA operation
The send_bucket() and _do_p2p_write_one_session() methods implement this with pipeline overlap (lines 80-122):
# From p2p.py lines 80-108, 107-122
def send_bucket(self, ...):
transfer_ready_params, ready_hf_tensors = self._get_transfer_ready_params(...)
for i, (model_replica, remote_weight_infos) in enumerate(
self._transfer_engine_meta_list
):
model_replica.load_weights(ready_hf_tensors) # Copy to pinned buffer
if i == last_idx:
# Last engine: fire-and-forget to background thread pool
for remote_session in remote_weight_infos:
self.transfer_manager.submit(
self._do_p2p_write_one_session,
remote_session,
transfer_ready_params
)
else:
# Non-last engines: synchronous completion
futures = [
self.transfer_manager.submit_returning_future(
self._do_p2p_write_one_session,
remote_session,
transfer_ready_params
)
for remote_session in remote_weight_infos
]
for f in futures:
f.result()
The last engine's writes execute in a background ThreadPoolExecutor while the trainer begins loading the next bucket—hiding transfer latency under computation.
Stage 6: Completion Synchronization
After dispatching all buckets, after_base_weights() (lines 61-70) waits for background futures:
# From p2p.py lines 61-70
def after_base_weights(self):
"""Wait for all pending P2P transfers before next training step."""
self.transfer_manager.wait_transfers()
This guarantees every engine has received the latest weights before training proceeds.
Why P2P RDMA Achieves Sub-50ms Latency
Four implementation details explain the performance advantage:
| Mechanism | Implementation Location | Benefit |
|---|---|---|
| Zero-copy, zero-GIL | p2p_transfer_utils.py lines 51-55 |
ensure_started() creates ThreadPoolExecutor with explicit comment: "RDMA ops won't be affected by the python GIL" |
| One-shot buffer reuse | register_cpu_memory() |
Single pinned buffer shared across all engine ranks eliminates per-engine allocation |
| Batch transfer | TransferEngine.batch_transfer_sync_write() |
One NIC-side RPC per bucket reduces doorbell overhead |
| Pipeline overlap | send_bucket() last-engine special case |
Background submission overlaps next bucket's load phase with current transfer |
Enabling and Configuring P2P RDMA
Users activate the P2P RDMA path via command-line flag:
python -m miles.scripts.run_qwen3_5_35B_A3B \
--weight-transfer-mode p2p \
--p2p-transfer-num-workers 8 \
--update-weight-buffer-size 512 \
...
Key configuration parameters:
--weight-transfer-mode p2p: SelectsUpdateWeightP2Pprotocol implementation--p2p-transfer-num-workers: Thread pool size for background writes--update-weight-buffer-size: Bucket size in megabytes (trade transfer granularity vs. overhead)
Code Walkthrough: Internal Flow
The UpdateWeightP2P class lifecycle demonstrates the complete flow:
# Initialization: build transfer plan
class UpdateWeightP2P(WeightTransferProtocol):
def __init__(self, args):
self.transfer_plan = RemoteTransferPlan(args)
self.targets = self.transfer_plan.plan_p2p()
def begin_sync(self):
# One-time buffer registration
if self.is_sender and not self._model_registered:
self._weight_memory_registry = register_cpu_memory(
self._shared_params_dict, self._transfer_engine
)
def send_bucket(self, ...):
# Bucket conversion and pipelined RDMA writes
# (implementation shown above)
def after_base_weights(self):
# Final synchronization
self.transfer_manager.wait_transfers()
Key Source Files
| File | Role |
|---|---|
miles/backends/training_utils/weight_update/protocols/p2p.py |
UpdateWeightP2P protocol: bucket handling, RDMA orchestration |
miles/backends/training_utils/weight_update/protocols/p2p_transfer_utils.py |
RemoteTransferPlan, register_cpu_memory, handshake logic |
docs/advanced/p2p-weight-transfer.md |
User documentation and performance benchmarks |
Summary
- P2P RDMA weight transfer replaces NCCL collectives with direct NIC-to-NIC writes, achieving ~30-50ms per-step latency versus 200-300ms for broadcast
- Six-stage pipeline: planning → handshake → buffer registration → bucketed conversion → pipelined writes → completion sync
- Shared CPU-pinned buffer provides O(1) memory overhead across any number of rollout engines
- Background thread pool enables pipeline overlap, hiding transfer latency under computation
- Zero-GIL design ensures Python runtime never blocks RDMA operations
Frequently Asked Questions
How do I enable P2P RDMA weight transfer in Miles?
Add --weight-transfer-mode p2p to your training launch command. This sets args.weight_transfer_mode = "p2p" which instantiates UpdateWeightP2P instead of the default NCCL-based protocol. Ensure your cluster has InfiniBand or RoCE networking configured.
What hardware is required for P2P RDMA weight transfer?
The system requires RDMA-capable NICs (InfiniBand HCAs or RoCE-enabled Ethernet adapters) with appropriate kernel drivers and user-space libraries. The Mooncake TransferEngine abstracts vendor specifics, but the underlying fabric must support RDMA writes.
Why is a single shared buffer more efficient than per-engine buffers?
Per-engine buffer allocation would require O(N) memory and O(N) registration operations with the NIC, where N is the number of rollout engines. The shared buffer design in register_cpu_memory() maintains constant memory regardless of scale, and registration occurs exactly once per trainer process.
What happens if an RDMA write fails during weight transfer?
The transfer_manager.wait_transfers() call in after_base_weights() propagates exceptions from background threads. Failed transfers surface as Future.result() exceptions, which halt training for investigation. The system does not silently proceed with stale weights.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →