Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
163 commits
Select commit Hold shift + click to select a range
de2b136
text curation updates
lbliii Sep 22, 2025
c520dde
concepts
lbliii Sep 22, 2025
0bb66b0
Merge branch 'main' into llane/text-curation-updates
lbliii Sep 22, 2025
62bf887
remove synthetic docs not for this release
lbliii Sep 22, 2025
faaf621
updates
lbliii Sep 22, 2025
faf801a
text concepts and getting started changes
lbliii Sep 22, 2025
c6d85e7
links, concepts
lbliii Sep 22, 2025
f46e9ae
crosslinks
lbliii Sep 22, 2025
54af361
quality assessment updates
lbliii Sep 22, 2025
1e8da76
more cleanup
lbliii Sep 22, 2025
342b14f
semdedup
lbliii Sep 22, 2025
5c8e46f
example import cleanup
lbliii Sep 22, 2025
89fc14d
concepts
lbliii Sep 22, 2025
2d18203
Merge branch 'main' into llane/text-curation-updates
lbliii Sep 23, 2025
c81a102
Merge branch 'main' into llane/text-curation-updates
lbliii Sep 24, 2025
29fa656
Update docs/about/concepts/text/data-acquisition-concepts.md
lbliii Sep 24, 2025
17261a5
feedback batch 1
lbliii Sep 24, 2025
bcff339
feedback batch 2
lbliii Sep 24, 2025
6f6c733
file_paths="/path/to/jsonl_directory",
lbliii Sep 24, 2025
545104e
revert removal of xenna for common crawl executors
lbliii Sep 24, 2025
29b0f6b
Merge branch 'main' into llane/text-curation-updates
lbliii Sep 24, 2025
9099950
quickstart installation steps
lbliii Sep 24, 2025
7775532
Update docs/about/concepts/text/data-acquisition-concepts.md
lbliii Sep 24, 2025
2454ea8
data loading concepts updates / simplification
lbliii Sep 24, 2025
baeb2a7
Merge branch 'main' into llane/text-curation-updates
lbliii Sep 25, 2025
3e1be2b
data processing feedback
lbliii Sep 25, 2025
1c45863
read-existing pg updates
lbliii Sep 25, 2025
98f1d0f
add-id updates
lbliii Sep 25, 2025
7946c4e
dedup updates
lbliii Sep 25, 2025
0c5e675
feedback
lbliii Sep 25, 2025
67d163d
Merge branch 'main' into llane/text-curation-updates
lbliii Sep 25, 2025
a528229
Merge llane/text-curation-updates
lbliii Sep 25, 2025
9cf24e0
gerge branch 'lbliii-llane/text-curation-updates'
lbliii Sep 25, 2025
c950aa9
Merge branch 'main' of github.com:NVIDIA/NeMo-Curator
lbliii Sep 25, 2025
8bc3e2e
ci: Test latest uv and no build group first install (#1117)
thomasdhc Sep 25, 2025
9eb14b7
Audio curation improvements (#1113)
Ssofja Sep 26, 2025
935ee9c
Add iv2 install guide and guard import (#1112)
suiyoubi Sep 26, 2025
a9dae11
Enhance RayClient shutdown process by resetting RAY_ADDRESS (#1101)
abhinavg4 Sep 26, 2025
28421b9
video doc improvement (#1124)
suiyoubi Sep 26, 2025
28f6711
Add torchvision as dependency for image (#1131)
praateekmahajan Sep 26, 2025
b300de3
about section updates (#1122)
lbliii Sep 29, 2025
1f5da7e
Update cherry-pick workflow to use v0.63.0 (#1137)
pablo-garay Sep 29, 2025
5d869c6
changelog (#1141)
lbliii Sep 30, 2025
6ed06a3
Llane/docs image updates (#1087)
lbliii Sep 30, 2025
de3a32e
text curation review feedback (#1129)
lbliii Sep 30, 2025
5e9c249
Fix auto download models for images (#1140)
huvunvidia Sep 30, 2025
bed9806
Llane/docs video feedback (#1146)
lbliii Sep 30, 2025
e347af8
Fix InternVideo2 doc and replace NVIDIA/NeMo-Curator with NVIDIA-NeMo…
suiyoubi Sep 30, 2025
c4b3b1a
Raise informative error when on unsupported system (#1126)
jrbourbeau Oct 1, 2025
370bfbf
Fixing versions1.json to support nested folders (#1157)
aschilling-nv Oct 1, 2025
f96cceb
Resolve a few warnings in the docs build (#1120)
jrbourbeau Oct 2, 2025
fc5e416
Context manager support for ``RayClient`` (#1155)
jrbourbeau Oct 2, 2025
648812f
added hydra params as defaults (#1159)
Ssofja Oct 2, 2025
b07dd68
Increase timeout for CPU tests (#1162)
sarahyurick Oct 3, 2025
1a42754
Llane/docs readme updates (#1163)
lbliii Oct 3, 2025
6227e75
add_codecov_badge (#1172)
pablo-garay Oct 7, 2025
1225061
Support more flags when building docs (#1175)
ayushdg Oct 8, 2025
9274b8e
Fix bug for running filters on empty batches (#1173)
sarahyurick Oct 8, 2025
bf4d2cf
docs: pyrpoject toml classifiers section (#1176)
lbliii Oct 9, 2025
04fe990
docs: migration guide and faq (#1167)
lbliii Oct 10, 2025
a398f84
citation file (#1165)
lbliii Oct 13, 2025
88e071d
add enhanced captioning to rn; bump versions on json files (#1154)
lbliii Oct 13, 2025
a88829b
Add large file splitting script (#1161)
jrbourbeau Oct 14, 2025
23d8f07
Fix ``JsonlWriter`` / ``ParquetWriter`` output directory typo (#1180)
jrbourbeau Oct 16, 2025
b3206cb
new extension: rich metadata for SEO (#1182)
lbliii Oct 17, 2025
1fea56f
Llama Nemotron Data Curation tutorial (#1063)
sarahyurick Oct 17, 2025
23823ad
Fix flaky ``test_split_parquet_file_by_size`` failure (#1185)
jrbourbeau Oct 20, 2025
0d23825
Switch from pynvml to nvidia-ml-py (#1186)
jrbourbeau Oct 20, 2025
89baaf3
Remove ftfy pin (#1189)
sarahyurick Oct 21, 2025
5f8166f
Increase UV_HTTP_TIMEOUT to reduce sporadic test failures (#1196)
jrbourbeau Oct 22, 2025
dd929a7
Link to documentation within text tutorials (#1190)
sarahyurick Oct 22, 2025
d843fa8
Bump ray to the latest version (#1195)
ayushdg Oct 22, 2025
8fc2cf5
remove max_seq_len_to_capture (#1200)
suiyoubi Oct 23, 2025
919771a
GLiNER PII Redaction tutorial (#1208)
sarahyurick Oct 29, 2025
3909b89
SDG Pipeline (#1215)
suiyoubi Nov 3, 2025
954222f
ci: Resolve starlette cve (#1219)
thomasdhc Nov 4, 2025
6af4735
Adds initial benchmarking framework (#1197)
rlratzel Nov 10, 2025
602524e
Small feedback updates to video curation (#1222)
sarahyurick Nov 10, 2025
da4ec43
Bump rapids to 25.10 (#1204)
ayushdg Nov 12, 2025
4f224c2
Update user-facing variable names and add override checks for `name`,…
sarahyurick Nov 13, 2025
0f1b98f
Update field names for classifiers (#1220)
sarahyurick Nov 13, 2025
324f8b1
Bug fix in Text Semantic Dedup Worklfow to allow output_filetype to …
praateekmahajan Nov 14, 2025
63d3b98
Adds optional filename arg to capture ray stdout/stderr (#1239)
rlratzel Nov 19, 2025
a8ba762
Add ability to ignore_head_node for RayDataExecutor and RayActorPoolE…
praateekmahajan Nov 19, 2025
4a82ed7
Fix ParquetReader and *Writer for Cloud I/O (#1249)
praateekmahajan Nov 20, 2025
eb5f118
Fix Pairwise IO and IdentifyDuplicates in SemDedup for Cloud I/O (#1253)
praateekmahajan Nov 20, 2025
f1a0577
Add tests for cloud I/O changes (#1257)
praateekmahajan Nov 20, 2025
131011b
``ScoreFilter`` ``score_fn=`` keyword typo in text filtering docs (#1…
jrbourbeau Nov 20, 2025
24dbe1c
Nemotron-CC SDG stages (#1245)
huvunvidia Nov 20, 2025
7f9abc0
Review and update text curation documentation (#1241)
sarahyurick Nov 21, 2025
f28f812
Update execution-backends.md (#1263)
arhamm1 Nov 25, 2025
df2f5d9
Update container-environments.md (#1262)
arhamm1 Nov 25, 2025
42d8a42
Update get-started/video.md (#1261)
arhamm1 Nov 25, 2025
94db01e
Update audio-tutorials-beginner.md (#1252)
arhamm1 Nov 25, 2025
c3374a0
Docs - Update curate-audio/process-data/quality-assessment/index.md (…
arhamm1 Nov 25, 2025
c8442c0
Update curate-audio/process-data/text-integration/index.md (#1247)
arhamm1 Nov 25, 2025
f0966a0
Update duration-filtering.md (#1243)
arhamm1 Nov 25, 2025
419b8c7
Update get-started/ image.md (#1240)
arhamm1 Nov 25, 2025
a2520f7
Update getting started - text quickstart.md (#1238)
arhamm1 Nov 25, 2025
6a497d4
Update getting started - index.md (#1237)
arhamm1 Nov 25, 2025
264573c
Nemotron-CC SDG section pipelines (#1268)
huvunvidia Dec 1, 2025
0f968bc
ci: Unbound transformers in default dep (#1267)
thomasdhc Dec 1, 2025
8d0b308
feat: Add pre-commit hooks for uv lock and export (#1277)
pablo-garay Dec 2, 2025
7decb5d
Adds option to set object-store size when starting Ray cluster, uses …
rlratzel Dec 3, 2025
96e6fc5
feat: CI perf improvement (manage bottlenecks) (#1285)
pablo-garay Dec 3, 2025
a8219ce
Update Memory Management Guide (#1279)
sarahyurick Dec 4, 2025
98e0277
Adds FuzzyDedup identification and Removal benchmarks (#1233)
rlratzel Dec 5, 2025
cf449c2
ci: Create contrainst for ray for CVE (#1286)
thomasdhc Dec 10, 2025
b87379a
Download Hugging Face datasets for text classification tutorials (#1288)
sarahyurick Dec 11, 2025
08453b9
Update wer-filtering.md (#1246)
arhamm1 Dec 12, 2025
8d8d02e
Enhance ASR documentation and code examples (#1276)
abhinavg4 Dec 12, 2025
d1b6e36
ci: Revert cicd cancellation on failure (#1298)
thomasdhc Dec 12, 2025
11388a5
bump urllib3 due to GHSA-gm62-xv2j-4w53 (#1306)
ayushdg Dec 15, 2025
a597ddf
ci: Disable uv caching for unit test (#1323)
thomasdhc Dec 16, 2025
be3634f
Add Cursor rules (#1294)
sarahyurick Dec 16, 2025
0c788d8
feat: Add dependabot workflow for periodic lock file updates v2 (#1307)
pablo-garay Dec 16, 2025
a9aa663
ci: Update vllm and torch to address CVE (#1287)
thomasdhc Dec 16, 2025
968732d
Inital tutorial for e2e fuzzy deduplication (#1242)
ayushdg Dec 17, 2025
b9126eb
Collision warning code (#1325)
huvunvidia Dec 17, 2025
fffbace
[Doc Review 25.09] Huy - Video Concepts + Text Advanced (#1270)
huvunvidia Dec 17, 2025
8efe410
feat: add install-test v3 (#1335)
pablo-garay Dec 19, 2025
d0d8c1c
ci: Skip video cuda12 and pip for extra all (#1338)
thomasdhc Dec 19, 2025
e6a3caa
Review Image Documentations (#1272)
suiyoubi Jan 5, 2026
7fab03e
Megatron tokenization pipeline (#1259)
asolergi-nv Jan 5, 2026
666e9ec
Fix ID generator to be non-blocking (#1344)
Chyroprase Jan 6, 2026
b4ac62f
Specify `allow_module_level=True` in InternVideo tests (#1354)
sarahyurick Jan 6, 2026
a14dbfb
ci: Always trigger NeMo CICD test stage (#1353)
thomasdhc Jan 6, 2026
1eebb85
[benchmarking] Updates to support multiple nightly runs, code cleanup…
rlratzel Jan 6, 2026
162a768
Remove nvenc/dec for xenna 0.1.6 (#1202)
suiyoubi Jan 7, 2026
a5baf25
Isolate modality tests with PyTest markers (#1293)
sarahyurick Jan 7, 2026
bc85ff2
Add YAML support for text pipelines (#1212)
sarahyurick Jan 9, 2026
d2b98c6
feat: FFmpeg to 8.0.1 (#1362)
suiyoubi Jan 12, 2026
13be895
ci: Bump version to 1.1.0 (#1364)
thomasdhc Jan 12, 2026
44864c6
Adding one worker per partition to FilePartioningStage and URLGenerat…
abhinavg4 Jan 13, 2026
4bf949a
Add workflow results (#1275)
praateekmahajan Jan 13, 2026
e148326
Add vLLM and Sentence Transformers support for embedding generation (…
praateekmahajan Jan 14, 2026
e22dbc4
[benchmarking] Adds image curation benchmark to nightly (#1341)
rlratzel Jan 14, 2026
c0df452
[benchmarking] Adds audio curation benchmark to nightly (#1360)
rlratzel Jan 14, 2026
fbde2cd
Fix bug in SDG example (#1370)
sarahyurick Jan 14, 2026
a2b1863
ferm docs init
lbliii Jan 14, 2026
299bcee
broken link batch
lbliii Jan 14, 2026
209c4fa
link fixes
lbliii Jan 14, 2026
8df9853
updates
lbliii Jan 21, 2026
6ae0906
update
lbliii Jan 21, 2026
ab145c2
update
lbliii Jan 21, 2026
9d774be
update
lbliii Jan 26, 2026
c4eb63a
update
lbliii Jan 29, 2026
aa2e33f
style updates
lbliii Feb 19, 2026
e2a2b9a
updates
lbliii Mar 3, 2026
624fbf2
updates
lbliii Mar 3, 2026
069e000
sanity check
lbliii Mar 3, 2026
7a55224
corrections
lbliii Mar 3, 2026
9559c99
updates
lbliii Mar 16, 2026
7784a58
fix "/index" routes
lbliii Mar 23, 2026
aa55fae
updates
lbliii Mar 23, 2026
66ebe29
updates
lbliii Mar 23, 2026
3879af8
updates
lbliii Mar 23, 2026
0d65e51
updates
lbliii Mar 23, 2026
778b0c7
fix
lbliii Mar 23, 2026
800ef64
fix links
lbliii Mar 23, 2026
3f8ec59
card link fixes
lbliii Mar 23, 2026
12db1e5
mass link fixes
lbliii Mar 23, 2026
d2bfb23
broken link bash
lbliii Mar 23, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
84 changes: 84 additions & 0 deletions .cursor/rules/coding-standards.mdc
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
---
alwaysApply: true
---

# NeMo Curator Coding Standards

## Linting and Formatting

The project uses **Ruff** for linting and formatting with line length of 119 characters.

## Key Style Rules

### Allowed Patterns

- ✅ Print statements (T20 ignored)
- ✅ Boolean arguments in functions (FBT ignored)
- ✅ `df` as variable name for DataFrames (PD901 ignored)
- ✅ TODOs without author/link (TD002, TD003 ignored)
- ✅ Long exception messages (TRY003 ignored)
- ✅ Accessing private attributes (SLF001 ignored)
- ✅ Branching after return (RET505-508 ignored)

### Required Patterns

- ❌ No docstrings required (D ignored)
- ❌ No pathlib enforcement (PTH ignored)
- ❌ No logging enforcement (G ignored)
- ✅ Type annotations for functions (except `*args`, `**kwargs`, special methods)

## File-Specific Exceptions

### `examples/` directory
- No `__init__.py` required (INP001)

### `tests/` directory
- Allow assertions (S101)
- Allow magic values (PLR2004)
- No return type annotations required (ANN201)

### `tutorials/` directory
- No `__init__.py` required (INP001)
- Ignore Unicode complaint (PLE2515)

## Copyright Header

All non-empty Python files must include the NVIDIA copyright header:

```python
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
```

## Python Version Support

Supports Python 3.10, 3.11, and 3.12 (requires >=3.10, <3.13)

## Loguru

NeMo Curator uses loguru for logging:

```python
from loguru import logger
```

Common uses include `logger.info`, `logger.warning`, and `logger.error`.

## PyTest Standards

All changes to NeMo Curator's source code must be accompanied with relevant tests. Tests which require a GPU and cannot be run on a CPU-only machine must be marked with `@pytest.mark.gpu`.

NeMo Curator enforces 80% PyTest coverage within the `nemo_curator/` directory.

For straightforward navigation purposes, the `tests/` directory structure matches the `nemo_curator/` directory structure.
70 changes: 70 additions & 0 deletions .cursor/rules/composite-stage-patterns.mdc
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
---
alwaysApply: true
---

# `CompositeStage` Patterns

## When to Use `CompositeStage`

Use `CompositeStage` for high-level, user-facing stages that decompose into multiple low-level execution stages:

- Provide simplified API while maintaining fine-grained execution control
- Bundle multiple related stages into a single logical operation
- Especially useful when stages require different resources (e.g., a CPU-based stage followed by a GPU-based stage)
- Decomposed during pipeline planning

## Creating a `CompositeStage`

```python
from dataclasses import dataclass

from nemo_curator.stages.base import CompositeStage, ProcessingStage


@dataclass
class MyCompositeStage(CompositeStage[InputTaskType, OutputTaskType]):
param1: str
param2: int

def __post_init__(self) -> None:
super().__init__()

self.stages = [
StageA(param1=self.param1),
StageB(param2=self.param2),
StageC(),
]

def inputs(self) -> tuple[list[str], list[str]]:
# StageA's inputs
return self.stages[0].inputs()

def outputs(self) -> tuple[list[str], list[str]]:
# StageC's outputs
return self.stages[2].outputs()

def decompose(self) -> list[ProcessingStage]:
return self.stages
```

## Configuration with `with_()`

`CompositeStage` use a different `with_()` signature that takes a dictionary:

```python
from nemo_curator.stages.resources import Resources


composite_stage = MyCompositeStage(param1, param2)

# Add a with operation
stage_config = {"StageA": {"resources": Resources(cpus=5.0)}}
updated_composite_stage = composite_stage.with_(stage_config)
```

## Important Rules

- Decomposed stages cannot be `CompositeStage`s themselves
- `inputs()` returns first stage's inputs
- `outputs()` returns last stage's outputs
- All stages in `decompose()` must have unique names for `with_()` to work
98 changes: 98 additions & 0 deletions .cursor/rules/executors.mdc
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
---
alwaysApply: true
---

# Executors

Executors are the runtime engines that execute NeMo Curator pipelines. They handle distributed task orchestration, resource allocation, and worker management.

## Role of Executors

Executors translate the high-level pipeline definition (a sequence of `ProcessingStage` instances) into actual distributed execution:

- **Orchestrate stage execution**: Run stages in sequence, passing tasks between them
- **Distribute work**: Parallelize task processing across workers/nodes
- **Manage resources**: Allocate CPUs, GPUs, memory according to stage requirements
- **Handle batching**: Group tasks into batches for efficient processing
- **Track performance**: Collect timing and throughput metrics
- **Manage lifecycle**: Setup and teardown workers/resources

## Available Executors

### `XennaExecutor` (Default)

Production executor using Cosmos-Xenna for distributed execution:

```python
from nemo_curator.backends.xenna import XennaExecutor

executor = XennaExecutor(config={
"logging_interval": 60,
"ignore_failures": False,
"execution_mode": "streaming", # or "batch"
"cpu_allocation_percentage": 0.95,
"autoscale_interval_s": 180,
})
```

### Experimental Executors

Located in `nemo_curator.backends.experimental`:

- **RayDataExecutor**: Ray Data backend (supports `ignore_head_node`)
- **RayActorPoolExecutor**: Ray Actor Pool backend (supports `ignore_head_node`)

## `BaseExecutor` Interface

All executors inherit from `BaseExecutor` and implement:

```python
class BaseExecutor(ABC):
def __init__(self, config: dict[str, Any] | None = None, ignore_head_node: bool = False):
"""Initialize executor with configuration.

Args:
config: Executor-specific configuration dictionary
ignore_head_node: Whether to exclude head node from execution (not supported by XennaExecutor)
"""

@abstractmethod
def execute(self, stages: list[ProcessingStage], initial_tasks: list[Task] | None = None) -> list[Task]:
"""Execute the pipeline stages.

Args:
stages: List of processing stages to execute
initial_tasks: Initial tasks to start pipeline (defaults to EmptyTask)

Returns:
List of output tasks from final stage
"""
```

## Stage Adapters

Executors use `BaseStageAdapter` to wrap stages for execution:

- Implements batching logic via `process_batch()`
- Calls stage lifecycle methods (`setup_on_node()`, `setup()`, `teardown()`)
- Tracks performance metrics with `StageTimer`
- Attaches performance statistics to output tasks

## Usage in Pipelines

Executors are passed to `pipeline.run()`:

```python
from nemo_curator.pipeline import Pipeline
from nemo_curator.backends.xenna import XennaExecutor


pipeline = Pipeline(name="my_pipeline", stages=[...])

# Explicit executor
executor = XennaExecutor(config={"execution_mode": "streaming"})
results = pipeline.run(executor=executor)

# Default executor (XennaExecutor)
results = pipeline.run()
```
70 changes: 70 additions & 0 deletions .cursor/rules/modality-structure.mdc
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
---
alwaysApply: true
---

# Modality Structure

NeMo Curator supports four data modalities, each with its own stage directory and patterns.

## Directory Structure

```
nemo_curator/stages/
├── text/ # Text/document processing
├── image/ # Image processing
├── audio/ # Audio/speech processing
└── video/ # Video processing
```

Other subdirectories of `nemo_curator/stages/` are `deduplication` and `synthetic`, which contain stages that can be used by two or more modalities.

## Text Processing (`stages/text/`)

Common subdirectories:
- `classifiers/`: Quality, domain, content safety classification
- `deduplication/`: Duplicate text removal
- `download/`: Data ingestion (Common Crawl, Wikipedia, ArXiv)
- `embedders/`: Text embedding generation
- `filters/`: Heuristic and model-based filtering
- `io/`: Readers (Parquet, JSONL) and writers
- `modifiers/`: Text transformations (cleaning, normalization)

Task type: `DocumentBatch`

## Image Processing (`stages/image/`)

Common operations:
- CLIP embedding generation
- Aesthetic quality scoring
- NSFW detection
- Deduplication

Task type: `ImageBatch`

## Audio Processing (`stages/audio/`)

Common operations:
- ASR transcription (NeMo Framework)
- Word Error Rate (WER) calculation
- Duration analysis
- Quality filtering

Task type: `AudioBatch`

## Video Processing (`stages/video/`)

Common operations:
- Scene detection (TransNetV2)
- Clip extraction
- GPU H.264 encoding/decoding
- Motion and aesthetic filtering
- Embeddings (InternVideo2, Cosmos-Embed1)

Task type: `VideoTask`

## Shared Components

All modalities share:
- `stages/base.py`: Base classes (`ProcessingStage`, `CompositeStage`)
- `stages/resources.py`: Resource configuration
- `stages/function_decorators.py`: Decorators for creating `ProcessingStage` instances from simple functions
57 changes: 57 additions & 0 deletions .cursor/rules/pipeline-structure.mdc
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
---
alwaysApply: true
---

# Pipeline Structure

## Creating Pipelines

Pipelines compose `ProcessingStage` instances into an executable workflow:

```python
from nemo_curator.pipeline import Pipeline

pipeline = Pipeline(
name="my_pipeline",
description="My data processing pipeline",
stages=[stage1, stage2, stage3],
)
```

## Adding Stages

Stages can be added during initialization or via `add_stage()`:

```python
# Method 1: During initialization
pipeline = Pipeline(name="my_pipeline", stages=[stage1, stage2])

# Method 2: Add stages individually
pipeline = Pipeline(name="my_pipeline")
pipeline.add_stage(stage1)
pipeline.add_stage(stage2)
```

## Running Pipelines

```python
pipeline.run()
```

The `run()` method accepts 2 optional parameters:

- `executor`: Executor to use. If None, defaults to `XennaExecutor`
- `initial_tasks` Initial `Task`s to start the pipeline with. Defaults to None.

Before executing, the `run()` method calls `self.build()` to build an execution plan from the pipeline (e.g., decomposes composite stages).

## Composite Stage Decomposition

When a pipeline is built, any `CompositeStage` instances are automatically decomposed into their constituent execution stages. The decomposition info is stored in `pipeline.decomposition_info`.

## Pipeline Methods

- `add_stage(stage)`: Add a stage to the pipeline (returns self for chaining)
- `build()`: Decompose composite stages into execution stages
- `describe()`: Get detailed description of pipeline stages and requirements
- `run(executor, initial_tasks)`: Execute the pipeline
Loading
Loading