nemo_curator.backends.ray_actor_pool
nemo_curator.backends.ray_actor_pool
Submodules
nemo_curator.backends.ray_actor_pool.adapternemo_curator.backends.ray_actor_pool.executornemo_curator.backends.ray_actor_pool.raft_adapternemo_curator.backends.ray_actor_pool.shuffle_adapternemo_curator.backends.ray_actor_pool.utils
Package Contents
Classes
API
Bases: BaseExecutor
Ray-based executor using ActorPool for better resource management.
This executor:
- Creates a pool of actors per stage using Ray’s ActorPool
- Uses map_unordered for better load balancing and fault tolerance
- Lets Ray handle object ownership and garbage collection automatically
- Provides better backpressure management through ActorPool
Clean up actors in the pool.
Clean up a list of actors.
Create an ActorPool for a specific stage.
Create a RAFT ActorPool for a specific stage.
Create a RapidsMPFShuffling Actors and setup UCXX communication for a specific stage.
Execute an LSH stage with band iteration.
Parameters:
LSHStage — The LSH stage to execute
list[ Task ] — Input tasks to process
Returns: list[ Task ]
List of output tasks from all band iterations
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
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
Process tasks through the actor pool.
Parameters:
The ActorPool to use for processing
ProcessingStage — The processing stage (for logging/context, unused)
list[ Task ] — List of Task objects to process
Returns: list[ Task ]
List of processed Task objects
Execute the pipeline stages using ActorPool.
Parameters:
list[ ProcessingStage ] — List of processing stages to execute
list[ Task ] | None — Initial tasks to process (can be None for empty start)
Returns: list[ Task ]
List of final processed tasks