Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:NVIDIA NeMo Curator RaftAdapter: Difference between revisions

From Leeroopedia
Auto-imported from implementations/NVIDIA_NeMo_Curator_RaftAdapter.md
 
Sync from local file
 
Line 143: Line 143:


* [[requires_env::Environment:NVIDIA_NeMo_Curator_Python_Linux_Base]]
* [[requires_env::Environment:NVIDIA_NeMo_Curator_Python_Linux_Base]]
* [[NVIDIA_NeMo_Curator_Backend_Base_Classes]] -- Parent base classes
* [[Implementation:NVIDIA_NeMo_Curator_Backend_Base_Classes]] -- Parent base classes
* [[NVIDIA_NeMo_Curator_RayActorPoolAdapter]] -- Base actor pool adapter
* [[Implementation:NVIDIA_NeMo_Curator_RayActorPoolAdapter]] -- Base actor pool adapter
* [[NVIDIA_NeMo_Curator_KMeansStage]] -- Example stage that uses RAFT for distributed k-means
* [[Implementation:NVIDIA_NeMo_Curator_KMeansStage]] -- Example stage that uses RAFT for distributed k-means


[[Category:Implementations]]
[[Category:Implementations]]

Latest revision as of 10:48, 27 September 2026

Knowledge Sources
Domains Backend Architecture, Ray Integration, Distributed GPU, RAFT
Last Updated 2026-02-14 00:00 GMT

Overview

Extends the base stage adapter to support RAFT-based distributed GPU stages by initializing NCCL communicators and RAFT handles across a pool of Ray actors.

Description

RayActorPoolRAFTAdapter is a specialized adapter that enables multi-GPU collective operations (such as all-reduce and all-gather) within the Ray ActorPool backend. It is required for distributed algorithms like fuzzy deduplication that use cuVS/RAFT.

The adapter follows this initialization sequence:

  1. On construction, it initializes RAFT Comms to obtain a unique NCCL ID. The actor at index 0 is designated as the root.
  2. The root actor broadcasts its unique ID to all other actors in the pool via broadcast_root_unique_id(), which uses ray.get_actor() to look up named actors and calls set_root_unique_id.remote() on each.
  3. During setup(), each actor initializes an NCCL communicator via nccl().init() using the shared root unique ID, then creates a RAFT handle via pylibraft.common.handle.Handle and injects NCCL communicator bindings using inject_comms_on_handle_coll_only().
  4. The RAFT handle (_raft_handle) and pool size (_actor_pool_size) are set on the underlying stage so it can use collective operations during processing.

The teardown() method delegates cleanup to the base adapter and the underlying stage.

Usage

This adapter is used internally by the Ray ActorPool executor when a stage requires RAFT-based multi-GPU communication. Users do not create this adapter directly. The executor determines from the stage specification whether RAFT is needed and constructs RayActorPoolRAFTAdapter instances with the appropriate index, pool size, session ID, and actor name prefix.

Code Reference

Source Location

  • Repository: NeMo-Curator
  • File: nemo_curator/backends/experimental/ray_actor_pool/raft_adapter.py
  • Lines: 1-160

Signature

class RayActorPoolRAFTAdapter(BaseStageAdapter):
    def __init__(
        self,
        stage: ProcessingStage,
        index: int,
        pool_size: int,
        session_id: bytes,
        actor_name_prefix: str = "RAFT",
    ): ...
    def get_batch_size(self) -> int: ...
    def setup_on_node(self) -> None: ...
    def broadcast_root_unique_id(self) -> None: ...
    def set_root_unique_id(self, root_unique_id: int) -> None: ...
    def setup(self, worker_metadata: WorkerMetadata | None = None) -> None: ...
    def teardown(self) -> None: ...

Import

from nemo_curator.backends.experimental.ray_actor_pool.raft_adapter import RayActorPoolRAFTAdapter

I/O Contract

Inputs

Name Type Required Description
stage ProcessingStage Yes The processing stage that requires RAFT collective operations
index int Yes The index of this actor in the pool (0 is root)
pool_size int Yes Total number of actors in the pool
session_id bytes Yes Unique session identifier for this execution
actor_name_prefix str No Prefix for Ray actor names (default: "RAFT")

Outputs

Name Type Description
get_batch_size() int Batch size for this stage (defaults to 1 if not set)
process_batch() list[Task] Processed tasks with performance stats (inherited from BaseStageAdapter)
root_unique_id int or None The NCCL unique ID from the root actor, used for collective initialization

Internal Methods

Method Description
_setup_nccl() Initializes the NCCL communicator using raft_dask.common.nccl.nccl() with the pool size, root unique ID, and actor index
_setup_raft() Creates a RAFT Handle and injects NCCL communicator bindings via inject_comms_on_handle_coll_only()

Usage Examples

RAFT Adapter Lifecycle (Internal)

# This is the sequence performed internally by the executor:
import ray

# 1. Create RAFT adapter actors
adapters = []
for i in range(pool_size):
    adapter = RayActorPoolRAFTAdapter(
        stage=my_raft_stage,
        index=i,
        pool_size=pool_size,
        session_id=session_id,
        actor_name_prefix="RAFT",
    )
    adapters.append(adapter)

# 2. Root broadcasts NCCL unique ID
adapters[0].broadcast_root_unique_id()

# 3. All actors run setup (initializes NCCL + RAFT)
for adapter in adapters:
    adapter.setup()

# 4. Process batches
results = adapter.process_batch(tasks)

# 5. Teardown
for adapter in adapters:
    adapter.teardown()

Related Pages