Implementation:NVIDIA NeMo Curator RaftAdapter: Difference between revisions
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:
- On construction, it initializes RAFT Comms to obtain a unique NCCL ID. The actor at index 0 is designated as the root.
- 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.
- 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().
- 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
- Environment:NVIDIA_NeMo_Curator_Python_Linux_Base
- Implementation:NVIDIA_NeMo_Curator_Backend_Base_Classes -- Parent base classes
- Implementation:NVIDIA_NeMo_Curator_RayActorPoolAdapter -- Base actor pool adapter
- Implementation:NVIDIA_NeMo_Curator_KMeansStage -- Example stage that uses RAFT for distributed k-means