Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion benchmarking/scripts/dedup_removal_benchmark.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ def run_removal_benchmark( # noqa: PLR0913
# Setup executor
# TODO: refactor utils.setup_executor to support this and remove this code
if executor == "ray_data":
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor

executor_obj = RayDataExecutor()
if use_ray_data_settings:
Expand Down
2 changes: 1 addition & 1 deletion benchmarking/scripts/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
import pyarrow.parquet as pq

from nemo_curator.backends.experimental.ray_actor_pool.executor import RayActorPoolExecutor
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor
from nemo_curator.backends.xenna import XennaExecutor
from nemo_curator.utils.file_utils import get_all_file_paths_and_size_under

Expand Down
2 changes: 1 addition & 1 deletion docs/curate-text/process-data/deduplication/semdedup.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ Get started with semantic deduplication using the following example of identifyi

```python
from nemo_curator.stages.text.deduplication.semantic import TextSemanticDeduplicationWorkflow
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor

workflow = TextSemanticDeduplicationWorkflow(
input_path="input_data/",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -432,7 +432,7 @@ pipeline.add_stage(JsonlWriter(path="filtered_output/"))
pipeline.run()

# Or use Ray for distributed processing (see Performance Tuning section)
# from nemo_curator.backends.experimental.ray_data import RayDataExecutor
# from nemo_curator.backends.ray_data import RayDataExecutor
# pipeline.run(RayDataExecutor(ignore_head_node=True))
```
:::
Expand Down Expand Up @@ -463,11 +463,11 @@ results = pipeline.run(executor)
If no executor is specified, `pipeline.run()` uses `XennaExecutor` with default settings.
:::

:::{tab-item} RayDataExecutor (Experimental)
:::{tab-item} RayDataExecutor
`RayDataExecutor` provides distributed processing using Ray Data. It has shown performance improvements for filtering workloads compared to the default executor.

```python
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor

executor = RayDataExecutor(
config={"ignore_failures": False},
Expand Down
2 changes: 1 addition & 1 deletion docs/reference/infrastructure/execution-backends.md
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,7 @@ For more details, refer to {ref}`Text Deduplication <text-process-data-dedup>`.
- **Scalable transformations**: Efficient map-batch operations across distributed workers

```python
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor

executor = RayDataExecutor()
results = pipeline.run(executor)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@ class RayDataExecutor(BaseExecutor):

def __init__(self, config: dict[str, Any] | None = None, ignore_head_node: bool = False):
super().__init__(config, ignore_head_node)
Comment on lines 41 to 43

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 No backward-compatibility shim for the old import path

The nemo_curator.backends.experimental.ray_data module is completely removed without a deprecation redirect. Any downstream code that depended on the old import path — including anything not covered by this PR — will immediately receive an ImportError with no migration guidance.

While the experimental label signals that breaking changes are expected, providing a thin shim in the old location is still considered good practice:

# nemo_curator/backends/experimental/ray_data/__init__.py  (re-created as a shim)
import warnings
warnings.warn(
    "nemo_curator.backends.experimental.ray_data is deprecated. "
    "Use nemo_curator.backends.ray_data instead.",
    DeprecationWarning,
    stacklevel=2,
)
from nemo_curator.backends.ray_data import RayDataExecutor  # noqa: F401, E402

__all__ = ["RayDataExecutor"]

This allows existing users to get a clear deprecation warning rather than a hard crash.

Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!

logger.warning("Ray Data executor is experimental and might not work as expected.")

def execute(self, stages: list["ProcessingStage"], initial_tasks: list[Task] | None = None) -> list[Task]:
"""Execute the pipeline stages using Ray Data.
Expand Down
4 changes: 2 additions & 2 deletions nemo_curator/stages/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -296,8 +296,8 @@ def get_config(self) -> dict[str, Any]:

def ray_stage_spec(self) -> dict[str, Any]:
"""Get Ray configuration for this stage.
Note : This is only used for Ray Data which is an experimental backend.
The keys are defined in RayStageSpecKeys in backends/experimental/ray_data/utils.py
Note : This is only used for Ray Data backend.
The keys are defined in RayStageSpecKeys in backends/ray_data/utils.py

Returns (dict[str, Any]):
Dictionary containing Ray-specific configuration
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,8 @@
import ray
from loguru import logger

from nemo_curator.backends.experimental.ray_data.executor import RayDataExecutor
from nemo_curator.backends.experimental.utils import RayStageSpecKeys
from nemo_curator.backends.ray_data.executor import RayDataExecutor
from nemo_curator.core.client import RayClient
from nemo_curator.stages.base import ProcessingStage, Resources
from nemo_curator.tasks import DocumentBatch, EmptyTask
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@

import pytest

from nemo_curator.backends.experimental.ray_data.utils import (
from nemo_curator.backends.ray_data.utils import (
calculate_concurrency_for_actors_for_stage,
get_available_cpu_gpu_resources,
)
Expand Down Expand Up @@ -63,7 +63,7 @@ def test_get_available_cpu_gpu_resources_mock_no_resources(self, mock_available_
class TestCalculateConcurrencyForActorsForStage:
"""Test class for calculate_concurrency_for_actors_for_stage function."""

@patch("nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources")
@patch("nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources")
def test_calculate_concurrency_explicit_num_workers(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency when num_workers is explicitly set."""
mock_stage = Mock(num_workers=lambda: 4, resources=Resources(cpus=2.0, gpus=0.0))
Expand All @@ -72,7 +72,7 @@ def test_calculate_concurrency_explicit_num_workers(self, mock_get_resources: Ma
mock_get_resources.assert_not_called()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
)
def test_calculate_concurrency_explicit_num_workers_zero_or_negative(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency when num_workers is explicitly set to 0 or negative."""
Expand All @@ -81,7 +81,7 @@ def test_calculate_concurrency_explicit_num_workers_zero_or_negative(self, mock_
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
)
def test_calculate_concurrency_cpu_only_constraint(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency with CPU-only constraint."""
Expand All @@ -90,7 +90,7 @@ def test_calculate_concurrency_cpu_only_constraint(self, mock_get_resources: Mag
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 4.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 4.0)
)
def test_calculate_concurrency_gpu_only_constraint(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency with GPU-only constraint."""
Expand All @@ -99,7 +99,7 @@ def test_calculate_concurrency_gpu_only_constraint(self, mock_get_resources: Mag
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 4.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 4.0)
)
def test_calculate_concurrency_both_cpu_gpu_constraints(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency with both CPU and GPU constraints."""
Expand All @@ -108,7 +108,7 @@ def test_calculate_concurrency_both_cpu_gpu_constraints(self, mock_get_resources
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(4.0, 8.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(4.0, 8.0)
)
def test_calculate_concurrency_cpu_more_limiting(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency when CPU is more limiting than GPU."""
Expand All @@ -117,7 +117,7 @@ def test_calculate_concurrency_cpu_more_limiting(self, mock_get_resources: Magic
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(16.0, 2.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(16.0, 2.0)
)
def test_calculate_concurrency_gpu_more_limiting(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency when GPU is more limiting than CPU."""
Expand All @@ -126,7 +126,7 @@ def test_calculate_concurrency_gpu_more_limiting(self, mock_get_resources: Magic
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
)
def test_calculate_concurrency_no_resource_requirements(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency when stage has no resource requirements."""
Expand All @@ -137,7 +137,7 @@ def test_calculate_concurrency_no_resource_requirements(self, mock_get_resources
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(1.0, 0.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(1.0, 0.0)
)
def test_calculate_concurrency_insufficient_resources(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency when there are insufficient resources."""
Expand All @@ -146,7 +146,7 @@ def test_calculate_concurrency_insufficient_resources(self, mock_get_resources:
mock_get_resources.assert_called_once()

@patch(
"nemo_curator.backends.experimental.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
"nemo_curator.backends.ray_data.utils.get_available_cpu_gpu_resources", return_value=(8.0, 2.0)
)
def test_calculate_concurrency_fractional_resources(self, mock_get_resources: MagicMock):
"""Test calculate_concurrency with fractional resource requirements."""
Expand Down
2 changes: 1 addition & 1 deletion tests/backends/test_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@

from nemo_curator.backends.base import BaseExecutor
from nemo_curator.backends.experimental.ray_actor_pool import RayActorPoolExecutor
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor
from nemo_curator.backends.xenna import XennaExecutor
from nemo_curator.tasks import FileGroupTask
from nemo_curator.tasks.utils import TaskPerfUtils
Expand Down
2 changes: 1 addition & 1 deletion tests/stages/deduplication/semantic/test_workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,7 @@ def test_semantic_deduplication_with_duplicate_identification(

executor = XennaExecutor()
elif executor_type == "ray_data":
from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor

executor = RayDataExecutor()

Expand Down
2 changes: 1 addition & 1 deletion tests/stages/text/deduplication/test_removal_workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
import pandas as pd
import pytest

from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor
from nemo_curator.backends.xenna import XennaExecutor
from nemo_curator.pipeline.workflow import WorkflowRunResult
from nemo_curator.stages.deduplication.id_generator import CURATOR_DEDUP_ID_STR
Expand Down
2 changes: 1 addition & 1 deletion tests/stages/text/deduplication/test_semantic.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
import pytest
from huggingface_hub import snapshot_download

from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor
from nemo_curator.backends.xenna import XennaExecutor
from nemo_curator.pipeline.workflow import WorkflowRunResult

Expand Down
2 changes: 1 addition & 1 deletion tests/stages/text/io/reader/test_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
import pandas as pd
import pytest

from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor
from nemo_curator.backends.xenna import XennaExecutor
from nemo_curator.pipeline import Pipeline
from nemo_curator.stages.deduplication.id_generator import (
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -273,7 +273,7 @@ def main() -> None: # noqa: PLR0915

import torch

from nemo_curator.backends.experimental.ray_data import RayDataExecutor
from nemo_curator.backends.ray_data import RayDataExecutor
from nemo_curator.core.client import RayClient

NUM_GPUS = 4 # noqa: N806
Expand Down
2 changes: 1 addition & 1 deletion tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb
Original file line number Diff line number Diff line change
Expand Up @@ -269,7 +269,7 @@
"\n",
"import torch\n",
"\n",
"from nemo_curator.backends.experimental.ray_data import RayDataExecutor\n",
"from nemo_curator.backends.ray_data import RayDataExecutor\n",
"from nemo_curator.core.client import RayClient\n",
"\n",
"NUM_GPUS = 2\n",
Expand Down
Loading