> 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.pipeline

## Submodules

* **[`nemo_curator.pipeline.pipeline`](/nemo/curator/nemo-curator/nemo_curator/pipeline/pipeline)**
* **[`nemo_curator.pipeline.workflow`](/nemo/curator/nemo-curator/nemo_curator/pipeline/workflow)**

## Package Contents

### Classes

| Name                                                   | Description                                                      |
| ------------------------------------------------------ | ---------------------------------------------------------------- |
| [`Pipeline`](#nemo_curator-pipeline-pipeline-Pipeline) | User-facing pipeline definition for composing processing stages. |

### API

```python
class nemo_curator.pipeline.Pipeline(
    name: str,
    description: str | None = None,
    stages: list[nemo_curator.stages.base.ProcessingStage] | None = None,
    config: dict[str, typing.Any] | None = None
)
```

User-facing pipeline definition for composing processing stages.

**`config`** `= config or {}`

---

**`stages`** `list[ProcessingStage] = stages or []`

---

```python
nemo_curator.pipeline.Pipeline.__repr__() -> str
```

String representation of the pipeline.

```python
nemo_curator.pipeline.Pipeline._assign_source_sink_roles() -> None
```

```python
nemo_curator.pipeline.Pipeline._decompose_stages(
    stages: list[nemo_curator.stages.base.ProcessingStage | nemo_curator.stages.base.CompositeStage]
) -> tuple[list[nemo_curator.stages.base.ProcessingStage], dict[str, list[str]]]
```

Decompose composite stages into execution stages.

**Parameters:**

**`stages`** `list[ProcessingStage | CompositeStage]`

`list[` [`ProcessingStage`](/nemo/curator/nemo-curator/nemo_curator/stages/base#nemo_curator-stages-base-ProcessingStage) `|` [`CompositeStage`](/nemo/curator/nemo-curator/nemo_curator/stages/base#nemo_curator-stages-base-CompositeStage) `]` — List of stages that may include composite stages

---

**Returns:** `tuple[list[` [`ProcessingStage`](/nemo/curator/nemo-curator/nemo_curator/stages/base#nemo_curator-stages-base-ProcessingStage) `], dict[str, list[str]]]`

tuple\[list\[ProcessingStage], dict\[str, list\[str]]]: Tuple of (execution stages, decomposition info dict)

**Raises:**

* `TypeError`: If a composite stage is decomposed into another composite stage

```python
nemo_curator.pipeline.Pipeline._run_with_resumability(
    executor: nemo_curator.backends.base.BaseExecutor,
    initial_tasks: list[nemo_curator.tasks.Task] | None,
    checkpoint_path: pathlib.Path
) -> list[nemo_curator.tasks.Task] | None
```

Run with resumability around a pre-existing Ray cluster (e.g. one
started by `RayClient`).

We briefly connect with `with ray.init()` to spawn the detached
checkpoint actor, then disconnect *before* `executor.execute` so the
executor's own `ray.init` runs un-nested — a nested
`ray.init(runtime_env=...)` is silently dropped, so the executor's env
vars wouldn't propagate otherwise. The detached actor lives in the
cluster across the executor's separate Ray session; a final
`with ray.init()` closes and kills it. The cluster must pre-exist: had
we started it, the first `with`-exit shutdown would tear it down and
take the actor with it.

```python
nemo_curator.pipeline.Pipeline.add_stage(
    stage: nemo_curator.stages.base.ProcessingStage
) -> nemo_curator.pipeline.pipeline.Pipeline
```

Add a stage to the pipeline.

**Parameters:**

**`stage`** `ProcessingStage`

[`ProcessingStage`](/nemo/curator/nemo-curator/nemo_curator/stages/base#nemo_curator-stages-base-ProcessingStage) — Processing stage to add

---

**Returns:** [`Pipeline`](#nemo_curator-pipeline-pipeline-Pipeline)

Self (Pipeline) for method chaining

```python
nemo_curator.pipeline.Pipeline.build() -> None
```

Build an execution plan from the pipeline.

**Raises:**

* `ValueError`: If the pipeline has no stages

```python
nemo_curator.pipeline.Pipeline.describe() -> str
```

Get a detailed description of the pipeline stages and their requirements.

```python
nemo_curator.pipeline.Pipeline.run(
    executor: nemo_curator.backends.base.BaseExecutor | None = None,
    initial_tasks: list[nemo_curator.tasks.Task] | None = None,
    checkpoint_path: str | pathlib.Path | None = None
) -> list[nemo_curator.tasks.Task] | None
```

Run the pipeline.

**Parameters:**

**`executor`** `BaseExecutor` — default: None

[`BaseExecutor`](/nemo/curator/nemo-curator/nemo_curator/backends/base#nemo_curator-backends-base-BaseExecutor) — Executor to use

---

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

`list[` [`Task`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task) `]` — Initial tasks to start the pipeline with. Defaults to None.

---

**`checkpoint_path`** `str | Path` — default: None

Resumability directory. Must
be a LOCAL filesystem path (the LMDB state is written locally),
not a remote/cloud URI. When set, completed source partitions are
tracked (in a `.nemo_curator_metadata` subdir) and skipped on
rerun. Multiple runs (e.g. a SLURM array) may share the directory
— each writes its own LMDB file, so there is no contention.

---

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

list\[Task] | None: List of tasks