Skip to content

[Proposal] Multi-rail parallel read: aggregate multiple RDMA paths for single-worker stripe reads #33

Description

@fresshman

Motivation

In KVService, large objects are striped across multiple NVMe devices / nodes, so
aggregate disk bandwidth scales with stripe count. A single client worker,
however, reads through ONE RDMA path (one NIC / QP). Once
stripe_count × disk_bw > single-NIC bandwidth, the network becomes the
bottleneck while the disks still have headroom — adding disks or concurrency on
the storage side no longer shortens read time.

Multi-rail lets one client worker transfer disjoint stripe subsets of the
same object over several independent RDMA paths concurrently, aggregating NIC
bandwidth.

What the codebase already gives us (why this is a generalization, not a rewrite)

  • RdmaClient.get_descriptor_stripes_into / get_descriptor_stripes_sge
    (wire tag 15, server RDMA-WRITEs straight into the caller's scattered
    destination buffers) already supports per-node concurrency: "callers issue
    one call per owning node … disjoint target regions let all nodes transfer
    concurrently".
  • Read completion is signalled over the TCP control channel (GET_RESP), so
    aggregating N rails = spawn one task per rail and join — no server-side
    polling changes needed.
  • Consistency fields are already on the wire: object_generation in
    ObjectDescriptor, chunk_checksums in PlacementChunk / StripingInfo.

So multi-rail = generalize "concurrent per node" to "concurrent per rail
(independent NIC/QP)", plus three missing pieces: path lifecycle management,
cross-rail completion aggregation, and memory-consistency guards.

Proposed design

Worker (vLLM / Dynamo / NIXL)
   │  upper interface unchanged: LookupObject / ReadByDescriptor
   ▼
RailManager  ── N independent rails (NIC + QP + CQ); single NIC = 1 rail
   ▼
Stripe→Rail scheduler ── locality affinity + round-robin (least-loaded later)
   ▼
per rail, concurrently: get_descriptor_stripes_sge(descriptor, stripe subset, SGE list)
   │  each rail writes disjoint regions of the same destination buffer
   ▼
Completion aggregation + integrity layer
   (object_generation check · per-stripe checksum · epoch buffer guard)
   ▼
Bottleneck attribution / observability (per-rail timing, classify
   NIC / PCIe / memory / disk / software)

Invariants (hard constraints):

  • On-disk stripe layout (StripingInfo) is never touched.
  • Upper read interface (LookupObject / ReadByDescriptor) is never touched.
  • Single-NIC deployments keep working unchanged: with one rail the scheduler
    is a pass-through.
  • v1 failure semantics: if any required rail fails, the whole request fails
    safely — no partial data returned, no transparent in-request retry (later
    versions may add retry).

Memory safety (late RDMA WRITE hazard): if a cancelled request's buffer is
freed and reused, a late-arriving RDMA WRITE can corrupt the new owner's data.
Three defences:

  1. pin-until-completion: the read holds the buffer and every rail's
    registration until join; cancel only flips a logical flag, MRs are not
    deregistered early.
  2. epoch buffer guard: late writes only land while live == epoch; bumping
    live on cancel makes the guard skip them.
  3. two-phase commit + per-stripe checksum: stripes are verified against
    chunk_checksums after aggregation; a cancelled request never reaches
    commit.

What is already working (Draft PR: #34)

  • multi_rail.rs — RailReader trait, RailManager, MultiRailReader
    (stripe planning → scoped-thread concurrent reads → aggregation → checksum
    verify), RailReadStats with a bottleneck classifier.

  • mock_rail.rs — Mock transport with configurable per-rail bandwidth and
    fault injection (late_write_after_cancel + epoch guard), so the full
    scheduling / state-machine / buffer-lifecycle path is testable on any
    machine (no RDMA hardware).

  • rdma.rs — RailReader implemented for RdmaClient (real verbs path,
    works for Soft-RoCE and physical NICs).

  • Tests (11 cases, default features, hardware-free): multi-rail aggregation
    with byte-exact reconstruction, single-rail pass-through, A/B demonstration
    of the late-write hazard (without cancel: corruption reproduced; with
    cancel: guard blocks it), stripe-planning round-robin, per-stripe checksum
    failure detection, single-rail disconnect safe-fail, partial-completion /
    under-delivery detection, read timeout with the caller's buffer untouched,
    inflight-budget backpressure rejection, and object version-change
    (generation) rejection.

  • cs-mock-bench — capacity model sweeping rail count (per-rail bandwidth
    configurable), measured on WSL2: 64 MiB object, 4 MiB chunks, 16 stripes,
    125 MB/s per rail:

    rails  agg_MB/s  bal_MB/s  theo_MB/s  speedup  bottleneck
    1      123.5     125.0     125.0      0.99x    network: slowest rail caps the object
    2      244.9     250.0     250.0      1.96x    software: post-read stripe verification (client CPU)
    3      327.4     333.3     375.0      2.62x    software: post-read stripe verification
    4      484.1     500.0     500.0      3.87x    software: post-read stripe verification
    

    Near-linear aggregate scaling (3.87x at 4 rails). bal accounts for stripe
    quantization (the slowest rail carries ceil(stripes/rails) stripes — hence
    the 3-rail dip). Notably, once the NIC bottleneck is removed the classifier
    already flags the next ceiling — client-side checksum CPU — which is
    exactly the per-layer overhead attribution this work instruments, and
    motivates aligning on the server's xxh3 checksums (question 3 below).

Environment honesty

All functional validation so far uses the Mock transport (the 11
hardware-free tests above). A Soft-RoCE (RXE) dual-path setup on WSL2
(two rxe devices, server listening via CS_RDMA_DEVICES=rxe0:...,rxe1:...)
is being wired up for end-to-end verbs validation — the
examples/softroce_dual_rail.rs harness is included in the PR. Even once
complete, Mock/Soft-RoCE prove functionality, failure semantics, scheduling
and resource governance — they do NOT prove hardware aggregate bandwidth,
HCA offload, PCIe/NUMA effects or zero-copy behaviour on physical NICs.
All numbers are labelled accordingly. No GPU is involved.

Questions for maintainers

  1. Where should rail configuration live (client config vs. endpoint list
    syntax)? CS_RDMA_DEVICES already accepts multiple server-side NICs —
    should the client mirror that syntax?
  2. The server can already listen on multiple RDMA devices
    (CS_RDMA_DEVICES=dev0:...,dev1:...). Any objection to using that as the
    multi-rail server shape, with per-stripe-subset GETs pinned to the rail
    that reaches the owning endpoint?
  3. Client-side stripe checksum is currently a lightweight local hash; the
    server's chunk_checksums are xxh3-based. Align the client on xxh3-64 so
    multi-rail verification checks the server's checksums?
  4. Preferred surface for per-rail metrics (throughput / latency / errors /
    in-flight / registered memory)?

Non-goals (v1)

  • No change to disk stripe layout or placement policy.
  • No transparent retry within a single request.
  • No GDS / GPU path changes.
  • No claim of physical multi-NIC aggregate bandwidth from Mock/Soft-RoCE data.

Happy to adjust scope or split this into smaller PRs — feedback welcome.

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