> This page is for version 26.07 · v1.3.0.
> For other versions, use one of these documentation indexes:
> - Latest · v1.4.0 (26.09) (default): https://docs.nvidia.com/nemo/curator/latest/llms.txt
> - Main · preview: https://docs.nvidia.com/nemo/curator/main/llms.txt
> - 26.09 · v1.4.0: https://docs.nvidia.com/nemo/curator/v26.09/llms.txt
> - 26.07 · v1.3.0: https://docs.nvidia.com/nemo/curator/v26.07/llms.txt
> - 26.04 · v1.2.0: https://docs.nvidia.com/nemo/curator/v26.04/llms.txt
> - 26.02 · v1.1.0: https://docs.nvidia.com/nemo/curator/v26.02/llms.txt
> - 25.09 · v1.0.0: https://docs.nvidia.com/nemo/curator/v25.09/llms.txt

> For clean Markdown of any page, append .md to the page URL.
> For a complete documentation index, see https://docs.nvidia.com/nemo/curator/llms.txt.
> For AI client integration (Claude Code, Cursor, etc.), connect to the MCP server at https://docs.nvidia.com/nemo/curator/_mcp/server.

# nemo_curator.backends.ray_actor_pool

## Submodules

* **[`nemo_curator.backends.ray_actor_pool.adapter`](/nemo/curator/nemo-curator/nemo_curator/backends/ray_actor_pool/adapter)**
* **[`nemo_curator.backends.ray_actor_pool.executor`](/nemo/curator/nemo-curator/nemo_curator/backends/ray_actor_pool/executor)**
* **[`nemo_curator.backends.ray_actor_pool.raft_adapter`](/nemo/curator/nemo-curator/nemo_curator/backends/ray_actor_pool/raft_adapter)**
* **[`nemo_curator.backends.ray_actor_pool.shuffle_adapter`](/nemo/curator/nemo-curator/nemo_curator/backends/ray_actor_pool/shuffle_adapter)**
* **[`nemo_curator.backends.ray_actor_pool.utils`](/nemo/curator/nemo-curator/nemo_curator/backends/ray_actor_pool/utils)**

## Package Contents

### Classes

| Name                                                                                          | Description                                                        |
| --------------------------------------------------------------------------------------------- | ------------------------------------------------------------------ |
| [`RayActorPoolExecutor`](#nemo_curator-backends-ray_actor_pool-executor-RayActorPoolExecutor) | Ray-based executor using ActorPool for better resource management. |

### API

```python
class nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor(
    config: dict | None = None,
    ignore_head_node: bool = False,
    show_progress: bool = True,
    progress_interval: float = 10.0
)
```

**Bases:** [BaseExecutor](/nemo/curator/nemo-curator/nemo_curator/backends/base#nemo_curator-backends-base-BaseExecutor)

Ray-based executor using ActorPool for better resource management.

This executor:

1. Creates a pool of actors per stage using Ray's ActorPool
2. Uses map\_unordered for better load balancing and fault tolerance
3. Lets Ray handle object ownership and garbage collection automatically
4. Provides better backpressure management through ActorPool

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._cleanup_actor_pool(
    actor_pool: ray.util.actor_pool.ActorPool
) -> None
```

Clean up actors in the pool.

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._cleanup_actors(
    actors: list[ray.actor.ActorHandle]
) -> None
```

Clean up a list of actors.

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._create_actor_pool(
    stage: nemo_curator.stages.base.ProcessingStage,
    num_actors: int
) -> ray.util.actor_pool.ActorPool
```

Create an ActorPool for a specific stage.

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._create_raft_actor_pool(
    stage: nemo_curator.stages.base.ProcessingStage,
    num_actors: int,
    session_id: bytes
) -> ray.util.actor_pool.ActorPool
```

Create a RAFT ActorPool for a specific stage.

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._create_rapidsmpf_actors(
    stage: nemo_curator.stages.base.ProcessingStage,
    num_actors: int,
    num_tasks: int
) -> list[ray.actor.ActorHandle]
```

Create a RapidsMPFShuffling Actors and setup UCXX communication for a specific stage.

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._execute_lsh_stage(
    stage: nemo_curator.stages.deduplication.fuzzy.lsh.stage.LSHStage,
    input_tasks: list[nemo_curator.tasks.Task],
    resource_baseline: tuple[float, float],
    reserved_cpus: float,
    reserved_gpus: float,
    resource_wait_timeout: float,
    resource_wait_interval: float
) -> list[nemo_curator.tasks.Task]
```

Execute an LSH stage with band iteration.

**Parameters:**

**`stage`** `LSHStage`

[`LSHStage`](/nemo/curator/nemo-curator/nemo_curator/stages/deduplication/fuzzy/lsh/stage#nemo_curator-stages-deduplication-fuzzy-lsh-stage-LSHStage) — The LSH stage to execute

---

**`input_tasks`** `list[Task]`

`list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `]` — Input tasks to process

---

**Returns:** `list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `]`

List of output tasks from all band iterations

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._generate_task_batches(
    tasks: list[nemo_curator.tasks.Task],
    batch_size: int | None = None,
    num_output_tasks: int | None = None,
    task_weights: list[int] | None = None
) -> list[list[nemo_curator.tasks.Task]]
```

Generate task batches from a list of tasks.
Args:
tasks: List of Task objects to process
batch\_size: The size of the batch
num\_output\_tasks: The number of output tasks to generate.
task\_weights: Optional weights used to balance tasks across output batches.
Either batch\_size or num\_output\_tasks must be provided but not both.
Returns:
List of task batches

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._get_stage_spec(
    stage: nemo_curator.stages.base.ProcessingStage
) -> dict
```

staticmethod

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._process_shuffle_stage_with_rapidsmpf_actors(
    actors: list[ray.actor.ActorHandle],
    tasks: list[nemo_curator.tasks.Task],
    band_range: tuple[int, int] | None = None
) -> list[nemo_curator.tasks.Task]
```

Process Shuffle through the actors.
Args:
actors: The actors to use for processing
tasks: List of Task objects to process
band\_range: Band range for LSH shuffle
Returns:
List of processed Task objects

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor._process_stage_with_pool(
    actor_pool: ray.util.actor_pool.ActorPool,
    _stage: nemo_curator.stages.base.ProcessingStage,
    tasks: list[nemo_curator.tasks.Task]
) -> list[nemo_curator.tasks.Task]
```

Process tasks through the actor pool.

**Parameters:**

**`actor_pool`** `ActorPool`

The ActorPool to use for processing

---

**`_stage`** `ProcessingStage`

[`ProcessingStage`](/nemo/curator/nemo-curator/nemo_curator/stages/base#nemo_curator-stages-base-ProcessingStage) — The processing stage (for logging/context, unused)

---

**`tasks`** `list[Task]`

`list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `]` — List of Task objects to process

---

**Returns:** `list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `]`

List of processed Task objects

```python
nemo_curator.backends.ray_actor_pool.RayActorPoolExecutor.execute(
    stages: list[nemo_curator.stages.base.ProcessingStage],
    initial_tasks: list[nemo_curator.tasks.Task] | None = None
) -> list[nemo_curator.tasks.Task]
```

Execute the pipeline stages using ActorPool.

**Parameters:**

**`stages`** `list[ProcessingStage]`

`list[` [`ProcessingStage`](/nemo/curator/nemo-curator/nemo_curator/stages/base#nemo_curator-stages-base-ProcessingStage) `]` — List of processing stages to execute

---

**`initial_tasks`** `list[Task] | None` — default: None

`list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `] | None` — Initial tasks to process (can be None for empty start)

---

**Returns:** `list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `]`

List of final processed tasks