Skip to content

vllm_omni.diffusion.offloader.distributed_layerwise_backend

Distributed Layerwise Offload backend with double-buffered H2D.

This module implements the RFC-1 "Distributed Layerwise Offload" mechanism that:

  • Optionally shards model weights across DP ranks and reconstructs each block with AllGather, or streams a complete rank-local block without a collective.
  • Can retain compatible checkpoint tensors as node-shared mmap sources instead of creating a persistent private host copy in every rank.
  • Uses a fixed double-buffer scheme that keeps only two layers' worth of weights on each device at any time.
  • Pipelines H2D transfers and, when enabled, AllGather communications on dedicated streams, overlapping them with computation.
  • Is hardware-agnostic, supporting both NVIDIA GPU (CUDA) and Ascend NPU (CANN) platforms via vLLM-Omni's platform abstraction layer.

logger module-attribute

logger = init_logger(__name__)

DistributedLayerwiseOffloadBackend

Bases: OffloadBackend

Distributed layer-wise (block-level) offloading backend.

Supports both GPU (CUDA) and NPU (CANN) platforms. Device type is determined by the device passed to enable().

Each rank stores only a shard of each block's weights on host memory. Two device slots alternate: one holds current weights, one holds next weights. H2D and AllGather run asynchronously on dedicated streams, overlapped with computation.

comm_stream instance-attribute

comm_stream = current_omni_platform.Stream()

copy_stream instance-attribute

copy_stream = current_omni_platform.Stream()

dp_group instance-attribute

dp_group: ProcessGroup | None = None

dp_size instance-attribute

dp_size = config.dp_size

host_weight_plan instance-attribute

host_weight_plan = host_weight_plan

rank instance-attribute

rank = 0

disable

disable() -> None

enable

enable(pipeline: Module) -> None

Enable DLO and make partial startup failures transactional.

load_resident_layers

load_resident_layers() -> None

Load the model-declared leading blocks for the denoise stage.

offload_resident_layers

offload_resident_layers() -> None

Release leading blocks before VAE decode to bound peak HBM.

shutdown

shutdown() -> None

DistributedLayerwiseOffloadHook

Bases: ModelHook

Hook for distributed layerwise offloading with fixed double-buffer.

Each rank stores only a shard of each block's weights on host memory. Two device slots alternate: one holds current weights, one holds next weights. H2D and AllGather run asynchronously on dedicated streams, overlapped with computation.

Supports both NVIDIA GPU (CUDA) and Ascend NPU (CANN) platforms.

block_buffers instance-attribute

block_buffers: dict[str, Tensor] = {}

block_parameters instance-attribute

block_parameters: dict[str, Parameter] = {}

comm_stream instance-attribute

comm_stream = comm_stream or current_omni_platform.Stream()

copy_stream instance-attribute

copy_stream = copy_stream or current_omni_platform.Stream()

cpu_shards instance-attribute

cpu_shards: dict[dtype, Tensor] = {}

cpu_sources instance-attribute

cpu_sources: dict[dtype, list[dict[str, Any]]] = {}

cpu_staging_buffers instance-attribute

cpu_staging_buffers: list[dict[dtype, Tensor] | None] = [
    None,
    None,
]

cpu_staging_events instance-attribute

cpu_staging_events: list[Any | None] = [None, None]

current_slot instance-attribute

current_slot = 0

device instance-attribute

device = device

dp_group instance-attribute

dp_group = dp_group

dp_size instance-attribute

dp_size = dp_size

gpu_buffers instance-attribute

gpu_buffers: list[dict[dtype, Tensor] | None] = (
    shared_buffers
)

gpu_shard_buffers instance-attribute

gpu_shard_buffers: list[dict[dtype, Tensor] | None] = [
    None,
    None,
]

is_materialized property

is_materialized: bool

Check whether this block's parameters hold real data on device.

metadata instance-attribute

metadata: dict[dtype, list[dict[str, Any]]] = {}

next_block instance-attribute

next_block = next_block

next_block_buffers instance-attribute

next_block_buffers: dict[str, Tensor] = {}

next_block_parameters instance-attribute

next_block_parameters: dict[str, Parameter] = {}

pin_memory instance-attribute

pin_memory = pin_memory

rank instance-attribute

rank = rank

rank_local_mmap instance-attribute

rank_local_mmap = rank_local_mmap

ready_events instance-attribute

ready_events: list[Any | None] = [None, None]

registered_mmap instance-attribute

registered_mmap = False

tensor_transforms instance-attribute

tensor_transforms = tensor_transforms or {}

initialize_hook

initialize_hook(module: Module) -> Module

offload_layer

offload_layer() -> None

Free GPU memory for current block by replacing tensors with placeholders.

post_forward

post_forward(module: Module, output: Any) -> Any

pre_forward

pre_forward(
    module: Module, *args: Any, **kwargs: Any
) -> tuple[tuple, dict]

prefetch_layer

prefetch_layer(
    slot: int, non_blocking: bool = True
) -> None

Prepare next block's weights into the shared device buffer for slot.

Uses the pre-allocated self.gpu_buffers[slot] instead of allocating fresh tensors every layer. This enforces the fixed double-buffer memory bound (exactly 2 blocks on device).

restore_next_block_to_cpu

restore_next_block_to_cpu() -> None

Restore hook-owned master weights before removing the hook.

Every circular hook owns the host backing for its next_block. Dropping that hook while the module points at rotating device buffers or placeholders would make a later enable shard invalid tensors.

PinnedResidentLayerGroup

Keep selected layers available for stage-scoped device residency.

TODO(offload): Extract this alongside PinnedModuleStager after the distributed shard-and-pin operation becomes a shared storage primitive. It currently remains here because it depends on DLO's local-shard layout. Unlike module.to(device)/module.to("cpu"), this group retains a pinned CPU master copy (or mmap source plus bounded staging) and never copies generated device weights back to host. Entering the denoise stage performs one asynchronous H2D pass; leaving it only restores zero-sized placeholders and releases the device buffers. This lets the following VAE stage reuse the same HBM.

The resident path is intentionally local-shard only. With tensor parallelism, the regular model loader has already produced each rank's TP shard, so no DP AllGather is required or desirable here.

copy_stream instance-attribute

copy_stream = copy_stream

device instance-attribute

device = device

loaded instance-attribute

loaded = False

pin_memory instance-attribute

pin_memory = pin_memory

rank_local_mmap instance-attribute

rank_local_mmap = rank_local_mmap

registered_mmap instance-attribute

registered_mmap = False

load

load() -> None

offload

offload() -> None

restore_to_cpu

restore_to_cpu() -> None

Materialize the persistent host masters back into the module.

Stage-scoped offload() deliberately leaves placeholders in the module while this group owns the CPU backing. disable() discards the group, so it must first restore ordinary CPU tensors to make a later enable cycle safe.

apply_distributed_block_hook

apply_distributed_block_hook(
    module: Module,
    next_block: Module,
    device: device,
    dp_group: ProcessGroup | None,
    dp_size: int,
    rank: int,
    copy_stream: Any | None = None,
    comm_stream: Any | None = None,
    pin_memory: bool = True,
    shared_buffers: list[dict[dtype, Tensor] | None]
    | None = None,
    rank_local_mmap: bool = False,
    tensor_transforms: dict[int, Any] | None = None,
    materialization_probe_tensor: Tensor | None = None,
) -> DistributedLayerwiseOffloadHook

Register a DistributedLayerwiseOffloadHook on module.

remove_distributed_block_hook

remove_distributed_block_hook(module: Module) -> None

Remove the distributed layerwise offload hook from module.