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

## Submodules

* **[`nemo_curator.tasks.audio_task`](/nemo/curator/nemo-curator/nemo_curator/tasks/audio_task)**
* **[`nemo_curator.tasks.document`](/nemo/curator/nemo-curator/nemo_curator/tasks/document)**
* **[`nemo_curator.tasks.file_group`](/nemo/curator/nemo-curator/nemo_curator/tasks/file_group)**
* **[`nemo_curator.tasks.image`](/nemo/curator/nemo-curator/nemo_curator/tasks/image)**
* **[`nemo_curator.tasks.interleaved`](/nemo/curator/nemo-curator/nemo_curator/tasks/interleaved)**
* **[`nemo_curator.tasks.lance`](/nemo/curator/nemo-curator/nemo_curator/tasks/lance)**
* **[`nemo_curator.tasks.ocr`](/nemo/curator/nemo-curator/nemo_curator/tasks/ocr)**
* **[`nemo_curator.tasks.sentinels`](/nemo/curator/nemo-curator/nemo_curator/tasks/sentinels)**
* **[`nemo_curator.tasks.tasks`](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks)**
* **[`nemo_curator.tasks.utils`](/nemo/curator/nemo-curator/nemo_curator/tasks/utils)**
* **[`nemo_curator.tasks.video`](/nemo/curator/nemo-curator/nemo_curator/tasks/video)**

## Package Contents

### Classes

| Name                                                                   | Description                                                                |
| ---------------------------------------------------------------------- | -------------------------------------------------------------------------- |
| [`AudioTask`](#nemo_curator-tasks-audio_task-AudioTask)                | A single audio manifest entry.                                             |
| [`DocumentBatch`](#nemo_curator-tasks-document-DocumentBatch)          | Task for processing batches of text documents.                             |
| [`EmptyTask`](#nemo_curator-tasks-sentinels-EmptyTask)                 | Seeds a pipeline with `task_id="0"` — the implicit root every task         |
| [`FailedTask`](#nemo_curator-tasks-sentinels-FailedTask)               | Marks a slot as failed → retried on resume (counter does NOT decrement).   |
| [`FileGroupTask`](#nemo_curator-tasks-file_group-FileGroupTask)        | Task representing a group of files to be read.                             |
| [`ImageBatch`](#nemo_curator-tasks-image-ImageBatch)                   | Task for processing batches of images.                                     |
| [`ImageObject`](#nemo_curator-tasks-image-ImageObject)                 | Represents a single image with metadata.                                   |
| [`InterleavedBatch`](#nemo_curator-tasks-interleaved-InterleavedBatch) | Task carrying row-wise multimodal records.                                 |
| [`LanceReadTask`](#nemo_curator-tasks-lance-LanceReadTask)             | Task containing Lance fragment ids assigned to one read partition.         |
| [`NoneTask`](#nemo_curator-tasks-sentinels-NoneTask)                   | Marks a slot as intentionally filtered (resumability counter decrements).  |
| [`SentinelTask`](#nemo_curator-tasks-sentinels-SentinelTask)           | Base for payload-less marker tasks: no data, framework-assigned `task_id`. |
| [`Task`](#nemo_curator-tasks-tasks-Task)                               | Abstract base class for tasks in the pipeline.                             |

### API

```python
class nemo_curator.tasks.AudioTask(
    dataset_name: str = '',
    data: dict = _AttrDict(),
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict(),
    task_id: str = '',
    filepath_key: str | None = None
)
```

Dataclass

**Bases:** [Task\[dict\]](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

A single audio manifest entry.

Represents one line from a JSONL manifest file (e.g. one audio file
with its metadata).  `data` is always a single `dict`, never a list.

Matches the `VideoTask` naming convention used by the video modality.

**Parameters:**

**`data`** `dict` — default: \_AttrDict()

Manifest entry dict (e.g. `&#123;"audio_filepath": "...", "text": "..."&#125;`).

---

**`filepath_key`** `str | None` — default: None

Optional key whose value is validated as an existing path.

---

**`data`** `dict = field(default_factory=_AttrDict)`

---

**`dataset_name`** `str = ''`

---

**`filepath_key`** `str | None = None`

---

**`num_items`** `int`

---

**`task_id`** `str = ''`

---

```python
nemo_curator.tasks.AudioTask.__post_init__()
```

```python
nemo_curator.tasks.AudioTask.validate() -> bool
```

Validate the task data.

```python
class nemo_curator.tasks.DocumentBatch(
    dataset_name: str,
    data: pyarrow.Table | pandas.DataFrame = pa.Table(),
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [Task\[Table | DataFrame\]](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

Task for processing batches of text documents.
Documents are stored as a dataframe (PyArrow Table or Pandas DataFrame).

**`data`** `Table | DataFrame = field(default_factory=(pa.Table))`

---

**`num_items`** `int`

Get the number of documents in this batch.

---

```python
nemo_curator.tasks.DocumentBatch.get_columns() -> list[str]
```

Get column names from the data.

```python
nemo_curator.tasks.DocumentBatch.to_pandas() -> pandas.DataFrame
```

Convert data to Pandas DataFrame.

```python
nemo_curator.tasks.DocumentBatch.to_pyarrow() -> pyarrow.Table
```

Convert data to PyArrow table.

```python
nemo_curator.tasks.DocumentBatch.validate() -> bool
```

Validate the task data.

```python
class nemo_curator.tasks.EmptyTask(
    dataset_name: str = 'empty',
    data: None = None,
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [SentinelTask](#nemo_curator-tasks-sentinels-SentinelTask)

Seeds a pipeline with `task_id="0"` — the implicit root every task
descends from (so all ids share the `"0"` prefix).

**`dataset_name`** `str = 'empty'`

---

**`task_id`** `str = field(init=False, default='0')`

---

```python
class nemo_curator.tasks.FailedTask(
    dataset_name: str = 'failed',
    data: None = None,
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [SentinelTask](#nemo_curator-tasks-sentinels-SentinelTask)

Marks a slot as failed → retried on resume (counter does NOT decrement).

**`dataset_name`** `str = 'failed'`

---

```python
class nemo_curator.tasks.FileGroupTask(
    dataset_name: str,
    data: list[str] = list(),
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict(),
    reader_config: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [Task\[list\[str\]\]](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

Task representing a group of files to be read.
This is created during the planning phase and passed to reader stages.

**`data`** `list[str] = field(default_factory=list)`

---

**`num_items`** `int`

Number of files in this group.

---

**`reader_config`** `dict[str, Any] = field(default_factory=dict)`

---

```python
nemo_curator.tasks.FileGroupTask.get_deterministic_id() -> str
```

Content-based id derived from the sorted file paths. Stable
across runs even if the source stage emits the file group at a
different position (e.g. because new files were added or removed
between runs).

```python
nemo_curator.tasks.FileGroupTask.validate() -> bool
```

Validate the task data.

```python
class nemo_curator.tasks.ImageBatch(
    dataset_name: str,
    data: list[nemo_curator.tasks.image.ImageObject] = list(),
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [Task](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

Task for processing batches of images.
Images are stored as a list of ImageObject instances, each containing
the path to the image and associated metadata.

**`data`** `list[ImageObject] = field(default_factory=list)`

---

**`num_items`** `int`

Number of images in this batch.

---

```python
nemo_curator.tasks.ImageBatch.validate() -> bool
```

Validate the task data.

```python
class nemo_curator.tasks.ImageObject(
    image_path: str = '',
    image_id: str = '',
    metadata: dict[str, typing.Any] = dict(),
    image_data: numpy.ndarray | None = None,
    embedding: numpy.ndarray | None = None,
    aesthetic_score: float | None = None,
    nsfw_score: float | None = None
)
```

Dataclass

Represents a single image with metadata.

**`aesthetic_score`** `float | None = None`

Aesthetic quality score as float

---

**`embedding`** `ndarray | None = None`

Image embedding vector as numpy array

---

**`image_data`** `ndarray | None = None`

Raw image pixel data as numpy array (H, W, C) in RGB format

---

**`image_id`** `str = ''`

Unique identifier for the image

---

**`image_path`** `str = ''`

Path to the image file on disk

---

**`metadata`** `dict[str, Any] = field(default_factory=dict)`

Additional metadata associated with the image

---

**`nsfw_score`** `float | None = None`

NSFW probability score as float

---

```python
class nemo_curator.tasks.InterleavedBatch(
    dataset_name: str,
    data: pyarrow.Table | pandas.DataFrame = (lambda: pa.Table.from_pyli...,
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict(),
    REQUIRED_COLUMNS: frozenset[str] = frozenset(name for name, f ...
)
```

Dataclass

**Bases:** [Task\[Table | DataFrame\]](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

Task carrying row-wise multimodal records.

See module docstring for the full schema reference (reserved vs user columns).

**`REQUIRED_COLUMNS`** `frozenset[str]`

---

**`data`** `Table | DataFrame`

---

**`num_items`** `int`

Number of unique samples (distinct `sample_id` values).

---

```python
nemo_curator.tasks.InterleavedBatch.add_rows(
    rows: pyarrow.Table | pandas.DataFrame | list[dict],
    sample_id: str | None = None,
    auto_position: bool = True
) -> nemo_curator.tasks.interleaved.InterleavedBatch
```

Add rows to this task.

**Parameters:**

**`rows`** `pa.Table | pd.DataFrame | list[dict]`

New rows to append. Must contain required columns unless
overridden by *sample\_id* / *auto\_position*.

---

**`sample_id`** `str | None` — default: None

If provided, assign this `sample_id` to all new rows.

---

**`auto_position`** `bool` — default: True

If `True`, auto-assign `position` values
continuing from the existing maximum per sample.

---

```python
nemo_curator.tasks.InterleavedBatch.build_source_ref(
    path: str | None,
    member: str | None,
    byte_offset: int | None = None,
    byte_size: int | None = None,
    frame_index: int | None = None
) -> str
```

staticmethod

Build a `source_ref` JSON locator string.

```python
nemo_curator.tasks.InterleavedBatch.count(
    modality: str | None = None
) -> int
```

Return row count, optionally filtered by modality.

Examples::

task.count()                    # total rows
task.count(modality="image")    # image rows only
task.count(modality="text")     # text rows only

```python
nemo_curator.tasks.InterleavedBatch.delete_rows(
    mask: pandas.Series
) -> nemo_curator.tasks.interleaved.InterleavedBatch
```

Delete rows where *mask* is `True`.

**Parameters:**

**`mask`** `pd.Series`

Boolean Series aligned to the data. `True` marks a row
for deletion.

---

```python
nemo_curator.tasks.InterleavedBatch.get_columns() -> list[str]
```

```python
nemo_curator.tasks.InterleavedBatch.parse_source_ref(
    source_value: str | None
) -> dict[str, str | int | None]
```

staticmethod

Parse a `source_ref` JSON string into a locator dict.

```python
nemo_curator.tasks.InterleavedBatch.to_pandas() -> pandas.DataFrame
```

```python
nemo_curator.tasks.InterleavedBatch.to_pyarrow() -> pyarrow.Table
```

```python
nemo_curator.tasks.InterleavedBatch.validate() -> bool
```

```python
nemo_curator.tasks.InterleavedBatch.with_parsed_source_ref_columns(
    prefix: str = '_src_'
) -> pandas.DataFrame
```

Return a DataFrame copy with parsed `source_ref` columns added.

Columns: `&#123;prefix&#125;path`, `&#123;prefix&#125;member`, `&#123;prefix&#125;byte_offset`,
`&#123;prefix&#125;byte_size`, `&#123;prefix&#125;frame_index`.

```python
class nemo_curator.tasks.LanceReadTask(
    dataset_name: str,
    data: nemo_curator.tasks.lance.FragmentIds = list(),
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict(),
    path: str,
    version: int
)
```

Dataclass

**Bases:** [Task\[FragmentIds\]](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

Task containing Lance fragment ids assigned to one read partition.

**Parameters:**

**`path`** `str`

Path or URI of the Lance dataset to read.

---

**`version`** `int`

Lance dataset version to read.

---

**`data`** `FragmentIds` — default: list()

[`FragmentIds`](#nemo_curator-tasks-lance-FragmentIds) — Lance fragment ids to read.

---

**`data`** `FragmentIds = field(default_factory=list)`

---

**`num_items`** `int`

---

**`path`** `str = field(kw_only=True)`

---

**`version`** `int = field(kw_only=True)`

---

```python
nemo_curator.tasks.LanceReadTask.get_deterministic_id() -> str
```

```python
nemo_curator.tasks.LanceReadTask.validate() -> bool
```

```python
class nemo_curator.tasks.NoneTask(
    dataset_name: str = 'none',
    data: None = None,
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [SentinelTask](#nemo_curator-tasks-sentinels-SentinelTask)

Marks a slot as intentionally filtered (resumability counter decrements).

**`dataset_name`** `str = 'none'`

---

```python
class nemo_curator.tasks.SentinelTask(
    dataset_name: str,
    data: None = None,
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

**Bases:** [Task\[None\]](/nemo/curator/nemo-curator/nemo_curator/tasks/tasks#nemo_curator-tasks-tasks-Task)

Base for payload-less marker tasks: no data, framework-assigned `task_id`.

**`data`** `None = None`

---

**`num_items`** `int`

---

```python
nemo_curator.tasks.SentinelTask.__post_init__() -> None
```

```python
nemo_curator.tasks.SentinelTask.validate() -> bool
```

```python
class nemo_curator.tasks.Task(
    dataset_name: str,
    data: nemo_curator.tasks.tasks.T,
    _stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
    _metadata: dict[str, typing.Any] = dict()
)
```

Dataclass

Abstract

**Bases:** `Generic[T]`

Abstract base class for tasks in the pipeline.

A task represents a batch of data to be processed. Different modalities
(text, audio, video) can implement their own task types.

**`_metadata`** `dict[str, Any] = field(default_factory=dict)`

Free-form metadata carried alongside the task.

---

**`_source_id`** `str = field(init=False, default='')`

Source (input partition) this task descends from. Stamped at
the source stage, inherited downstream; used only by the opt-in
resumability layer. Empty for pre-source tasks.

---

**`_stage_perf`** `list[StagePerfStats] = field(default_factory=list)`

Per-stage perf stats this task has accumulated.

---

**`data`** `T`

The task's payload (modality-specific).

---

**`dataset_name`** `str`

Name of the dataset this task belongs to.

---

**`num_items`** `int`

Get the number of items in this task.

---

**`task_id`** `str = field(init=False, default='')`

Deterministic identifier for this task. NOT user-settable —
the framework assigns it via `_set_task_id` at every stage
boundary. It is an underscore-joined id path through the pipeline
DAG — the parents' ids plus this task's own segment (e.g.
`"abc123_0_5"` = source `abc123`, then child 0, then
grandchild 5). Using the readable path directly (rather than a
hash of it) keeps task ids easy to debug. Empty string until the
first stage runs; two runs of the same pipeline on the same
inputs produce byte-identical `task_id`s across all tasks.

A `task_id` that starts with `"r"` (followed by a uuid) is a
fallback assigned when the parent→child mapping could NOT be
derived — e.g. a stage that overrides `process_batch` with an
ambiguous batch fan-out (M inputs → K≠M outputs). Such ids are
NON-deterministic (differ across runs).

---

```python
nemo_curator.tasks.Task.__post_init__() -> None
```

Post-initialization hook.

```python
nemo_curator.tasks.Task.__repr__() -> str
```

```python
nemo_curator.tasks.Task._set_task_id(
    parent_task_id: str,
    current_task_id_suffix: str | int
) -> None
```

Assign this task's deterministic `task_id` from its parent.

The `task_id` is the parent id and this task's own segment joined
by `"_"` — e.g. parent `"abc123"` + suffix `0` →
`"abc123_0"`. Always overwrites `task_id`; there is no
idempotency check — each stage transition re-derives it, so the
same physical Python object passing through N stages gets N
distinct `task_id`s (one per stage boundary). The dedup keys
used by resumability are captured BEFORE this method runs on a
given output, so the rewrite is safe.

Only a single parent id is taken: the supported mappings (1→1,
1→N fan-out, N→N positional) each give an output exactly one
parent. N→1 aggregations don't track ancestry — those outputs get
a random `"r"`-prefixed id in the adapter instead of calling this.

**Parameters:**

**`parent_task_id`** `str`

`task_id` of the parent. An empty string
(an unassigned / EmptyTask parent) is dropped so it doesn't
contribute a leading `"_"` to the path.

---

**`current_task_id_suffix`** `str | int`

This task's own segment of the id
path — appended after the parent id. Either a positional
index (`int` → coerced to `str`) for plain emissions,
or a string id (e.g. a content-based hash from
`get_deterministic_id`) for source-stage emissions
where stability across input reordering matters.

---

```python
nemo_curator.tasks.Task.add_stage_perf(
    perf_stats: nemo_curator.utils.performance_utils.StagePerfStats
) -> None
```

Add performance stats for a stage.

```python
nemo_curator.tasks.Task.get_deterministic_id() -> str | None
```

Return a content-based identifier for this task as a source,
or `None` to fall back to the positional index.

Override in subclasses that have stable content. The canonical
example is `FileGroupTask`, which hashes its sorted file
paths so that adding or removing files between runs doesn't shift
the identifiers of unchanged source partitions.

Only called by source-stage adapters; non-source stages ignore
this and always use positional indices.

```python
nemo_curator.tasks.Task.get_source_id() -> str
```

This task's source-partition identity: the trailing segment of
`task_id` (the id-path leaf). At a source stage that segment is the
partition's own id (content id or index); the resumability layer stamps
it onto `_source_id` and inherits it downstream. Kept here next to
`_set_task_id` so the `"_"` id-path encoding lives in one place.

```python
nemo_curator.tasks.Task.validate() -> bool
```

abstract

Validate the task data.