From c3f0f9f6ece6ef7ace2d796b630508ef8df8b6da Mon Sep 17 00:00:00 2001 From: Onur Yilmaz Date: Tue, 17 Mar 2026 20:48:00 -0400 Subject: [PATCH 1/5] Move ray data out of experimental folder Signed-off-by: Onur Yilmaz --- .../scripts/dedup_removal_benchmark.py | 2 +- benchmarking/scripts/utils.py | 2 +- .../process-data/deduplication/semdedup.md | 2 +- .../quality-assessment/heuristic.md | 6 ++--- .../infrastructure/execution-backends.md | 2 +- .../{experimental => }/ray_data/__init__.py | 0 .../{experimental => }/ray_data/adapter.py | 0 .../{experimental => }/ray_data/executor.py | 1 - .../{experimental => }/ray_data/utils.py | 0 nemo_curator/stages/base.py | 4 ++-- .../{experimental => }/ray_data/__init__.py | 0 .../ray_data/test_max_calls_pid.py | 2 +- .../{experimental => }/ray_data/test_utils.py | 22 +++++++++---------- tests/backends/test_integration.py | 2 +- .../deduplication/semantic/test_workflow.py | 2 +- .../deduplication/test_removal_workflow.py | 2 +- .../text/deduplication/test_semantic.py | 2 +- .../stages/text/io/reader/test_integration.py | 2 +- .../ndd_data_generation_example.py | 2 +- 19 files changed, 27 insertions(+), 28 deletions(-) rename nemo_curator/backends/{experimental => }/ray_data/__init__.py (100%) rename nemo_curator/backends/{experimental => }/ray_data/adapter.py (100%) rename nemo_curator/backends/{experimental => }/ray_data/executor.py (98%) rename nemo_curator/backends/{experimental => }/ray_data/utils.py (100%) rename tests/backends/{experimental => }/ray_data/__init__.py (100%) rename tests/backends/{experimental => }/ray_data/test_max_calls_pid.py (99%) rename tests/backends/{experimental => }/ray_data/test_utils.py (85%) diff --git a/benchmarking/scripts/dedup_removal_benchmark.py b/benchmarking/scripts/dedup_removal_benchmark.py index 9726ac1cb6..7bb4aff349 100755 --- a/benchmarking/scripts/dedup_removal_benchmark.py +++ b/benchmarking/scripts/dedup_removal_benchmark.py @@ -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: diff --git a/benchmarking/scripts/utils.py b/benchmarking/scripts/utils.py index fa032cad13..1e187a9f32 100644 --- a/benchmarking/scripts/utils.py +++ b/benchmarking/scripts/utils.py @@ -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 diff --git a/docs/curate-text/process-data/deduplication/semdedup.md b/docs/curate-text/process-data/deduplication/semdedup.md index cbcb581f16..8672b3215d 100644 --- a/docs/curate-text/process-data/deduplication/semdedup.md +++ b/docs/curate-text/process-data/deduplication/semdedup.md @@ -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/", diff --git a/docs/curate-text/process-data/quality-assessment/heuristic.md b/docs/curate-text/process-data/quality-assessment/heuristic.md index 00623e1fb1..8cb81e44a9 100644 --- a/docs/curate-text/process-data/quality-assessment/heuristic.md +++ b/docs/curate-text/process-data/quality-assessment/heuristic.md @@ -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)) ``` ::: @@ -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}, diff --git a/docs/reference/infrastructure/execution-backends.md b/docs/reference/infrastructure/execution-backends.md index 852f5e06db..511c01ead2 100644 --- a/docs/reference/infrastructure/execution-backends.md +++ b/docs/reference/infrastructure/execution-backends.md @@ -150,7 +150,7 @@ For more details, refer to {ref}`Text Deduplication `. - **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) diff --git a/nemo_curator/backends/experimental/ray_data/__init__.py b/nemo_curator/backends/ray_data/__init__.py similarity index 100% rename from nemo_curator/backends/experimental/ray_data/__init__.py rename to nemo_curator/backends/ray_data/__init__.py diff --git a/nemo_curator/backends/experimental/ray_data/adapter.py b/nemo_curator/backends/ray_data/adapter.py similarity index 100% rename from nemo_curator/backends/experimental/ray_data/adapter.py rename to nemo_curator/backends/ray_data/adapter.py diff --git a/nemo_curator/backends/experimental/ray_data/executor.py b/nemo_curator/backends/ray_data/executor.py similarity index 98% rename from nemo_curator/backends/experimental/ray_data/executor.py rename to nemo_curator/backends/ray_data/executor.py index da46ffea9f..63670d20cc 100644 --- a/nemo_curator/backends/experimental/ray_data/executor.py +++ b/nemo_curator/backends/ray_data/executor.py @@ -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) - 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. diff --git a/nemo_curator/backends/experimental/ray_data/utils.py b/nemo_curator/backends/ray_data/utils.py similarity index 100% rename from nemo_curator/backends/experimental/ray_data/utils.py rename to nemo_curator/backends/ray_data/utils.py diff --git a/nemo_curator/stages/base.py b/nemo_curator/stages/base.py index cbf7652ac0..073ad8d8ad 100644 --- a/nemo_curator/stages/base.py +++ b/nemo_curator/stages/base.py @@ -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/experimental/utils.py Returns (dict[str, Any]): Dictionary containing Ray-specific configuration diff --git a/tests/backends/experimental/ray_data/__init__.py b/tests/backends/ray_data/__init__.py similarity index 100% rename from tests/backends/experimental/ray_data/__init__.py rename to tests/backends/ray_data/__init__.py diff --git a/tests/backends/experimental/ray_data/test_max_calls_pid.py b/tests/backends/ray_data/test_max_calls_pid.py similarity index 99% rename from tests/backends/experimental/ray_data/test_max_calls_pid.py rename to tests/backends/ray_data/test_max_calls_pid.py index 6c3a258994..55464d74c4 100644 --- a/tests/backends/experimental/ray_data/test_max_calls_pid.py +++ b/tests/backends/ray_data/test_max_calls_pid.py @@ -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 diff --git a/tests/backends/experimental/ray_data/test_utils.py b/tests/backends/ray_data/test_utils.py similarity index 85% rename from tests/backends/experimental/ray_data/test_utils.py rename to tests/backends/ray_data/test_utils.py index 8958457b35..d4dc961cd9 100644 --- a/tests/backends/experimental/ray_data/test_utils.py +++ b/tests/backends/ray_data/test_utils.py @@ -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, ) @@ -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)) @@ -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.""" @@ -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.""" @@ -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.""" @@ -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.""" @@ -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.""" @@ -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.""" @@ -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.""" @@ -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.""" @@ -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.""" diff --git a/tests/backends/test_integration.py b/tests/backends/test_integration.py index 8605f27bd8..0393f1a884 100644 --- a/tests/backends/test_integration.py +++ b/tests/backends/test_integration.py @@ -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 diff --git a/tests/stages/deduplication/semantic/test_workflow.py b/tests/stages/deduplication/semantic/test_workflow.py index bae4a6c8c8..bb3ba468a9 100644 --- a/tests/stages/deduplication/semantic/test_workflow.py +++ b/tests/stages/deduplication/semantic/test_workflow.py @@ -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() diff --git a/tests/stages/text/deduplication/test_removal_workflow.py b/tests/stages/text/deduplication/test_removal_workflow.py index 5682931675..09f8b0163e 100644 --- a/tests/stages/text/deduplication/test_removal_workflow.py +++ b/tests/stages/text/deduplication/test_removal_workflow.py @@ -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 diff --git a/tests/stages/text/deduplication/test_semantic.py b/tests/stages/text/deduplication/test_semantic.py index d9907e8d90..e06b74f9c8 100644 --- a/tests/stages/text/deduplication/test_semantic.py +++ b/tests/stages/text/deduplication/test_semantic.py @@ -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 diff --git a/tests/stages/text/io/reader/test_integration.py b/tests/stages/text/io/reader/test_integration.py index 19a82e5d3a..dd6e415a1e 100644 --- a/tests/stages/text/io/reader/test_integration.py +++ b/tests/stages/text/io/reader/test_integration.py @@ -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 ( diff --git a/tutorials/synthetic/nemo_data_designer/ndd_data_generation_example.py b/tutorials/synthetic/nemo_data_designer/ndd_data_generation_example.py index edb4efb101..67b1aee2ec 100644 --- a/tutorials/synthetic/nemo_data_designer/ndd_data_generation_example.py +++ b/tutorials/synthetic/nemo_data_designer/ndd_data_generation_example.py @@ -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 From acf301ecf45319b99453725674ece7990935187d Mon Sep 17 00:00:00 2001 From: Onur Yilmaz Date: Tue, 17 Mar 2026 20:50:16 -0400 Subject: [PATCH 2/5] Fix linting issues Signed-off-by: Onur Yilmaz --- tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb b/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb index 9456b6e248..9c473749b2 100644 --- a/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb +++ b/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb @@ -268,8 +268,8 @@ "import time\n", "\n", "import torch\n", - "\n", "from nemo_curator.backends.experimental.ray_data import RayDataExecutor\n", + "\n", "from nemo_curator.core.client import RayClient\n", "\n", "NUM_GPUS = 2\n", From 8c76bd543b829f07ae9ded8889d3b4ce01dbdc27 Mon Sep 17 00:00:00 2001 From: Onur Yilmaz <35306097+oyilmaz-nvidia@users.noreply.github.com> Date: Wed, 18 Mar 2026 13:24:41 -0400 Subject: [PATCH 3/5] Update tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb Co-authored-by: Sarah Yurick <53962159+sarahyurick@users.noreply.github.com> Signed-off-by: Onur Yilmaz <35306097+oyilmaz-nvidia@users.noreply.github.com> --- tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb b/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb index 9c473749b2..5645cb30f7 100644 --- a/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb +++ b/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb @@ -268,7 +268,7 @@ "import time\n", "\n", "import torch\n", - "from nemo_curator.backends.experimental.ray_data import RayDataExecutor\n", + "from nemo_curator.backends.ray_data import RayDataExecutor\n", "\n", "from nemo_curator.core.client import RayClient\n", "\n", From df611b5cb6e0b0a689815765ba0534ad136b6c2b Mon Sep 17 00:00:00 2001 From: Onur Yilmaz <35306097+oyilmaz-nvidia@users.noreply.github.com> Date: Wed, 18 Mar 2026 13:27:02 -0400 Subject: [PATCH 4/5] Update base.py Signed-off-by: Onur Yilmaz <35306097+oyilmaz-nvidia@users.noreply.github.com> --- nemo_curator/stages/base.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/nemo_curator/stages/base.py b/nemo_curator/stages/base.py index 073ad8d8ad..19e0b4f193 100644 --- a/nemo_curator/stages/base.py +++ b/nemo_curator/stages/base.py @@ -297,7 +297,7 @@ 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 backend. - The keys are defined in RayStageSpecKeys in backends/experimental/utils.py + The keys are defined in RayStageSpecKeys in backends/ray_data/utils.py Returns (dict[str, Any]): Dictionary containing Ray-specific configuration From 2dcc3e845489933a2ed6d15a23bac7f2b6728b5e Mon Sep 17 00:00:00 2001 From: Onur Yilmaz Date: Wed, 18 Mar 2026 18:58:16 -0400 Subject: [PATCH 5/5] Fix linting issues Signed-off-by: Onur Yilmaz --- tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb b/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb index 5645cb30f7..dc51b0f2c4 100644 --- a/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb +++ b/tutorials/text/deduplication/fuzzy/fuzzy_e2e.ipynb @@ -268,8 +268,8 @@ "import time\n", "\n", "import torch\n", - "from nemo_curator.backends.ray_data import RayDataExecutor\n", "\n", + "from nemo_curator.backends.ray_data import RayDataExecutor\n", "from nemo_curator.core.client import RayClient\n", "\n", "NUM_GPUS = 2\n",