nemo_curator.tasks

View as Markdown

Submodules

Package Contents

Classes

NameDescription
AudioTaskA single audio manifest entry.
DocumentBatchTask for processing batches of text documents.
EmptyTaskSeeds a pipeline with task_id="0" — the implicit root every task
FailedTaskMarks a slot as failed → retried on resume (counter does NOT decrement).
FileGroupTaskTask representing a group of files to be read.
ImageBatchTask for processing batches of images.
ImageObjectRepresents a single image with metadata.
InterleavedBatchTask carrying row-wise multimodal records.
LanceReadTaskTask containing Lance fragment ids assigned to one read partition.
NoneTaskMarks a slot as intentionally filtered (resumability counter decrements).
SentinelTaskBase for payload-less marker tasks: no data, framework-assigned task_id.
TaskAbstract base class for tasks in the pipeline.

API

class nemo_curator.tasks.AudioTask(
dataset_name: str = '',
data: dict = _AttrDict(),
_metadata: dict[str, typing.Any] = dict(),
task_id: str = '',
filepath_key: str | None = None
)
Dataclass

Bases: Task[dict]

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
dictDefaults to _AttrDict()

Manifest entry dict (e.g. {"audio_filepath": "...", "text": "..."}).

filepath_key
str | NoneDefaults to 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 = ''
nemo_curator.tasks.AudioTask.__post_init__()
nemo_curator.tasks.AudioTask.validate() -> bool

Validate the task data.

class nemo_curator.tasks.DocumentBatch(
dataset_name: str,
data: pyarrow.Table | pandas.DataFrame = pa.Table(),
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: Task[Table | DataFrame]

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.

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

Get column names from the data.

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

Convert data to Pandas DataFrame.

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

Convert data to PyArrow table.

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

Validate the task data.

class nemo_curator.tasks.EmptyTask(
dataset_name: str = 'empty',
data: None = None,
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: 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')
class nemo_curator.tasks.FailedTask(
dataset_name: str = 'failed',
data: None = None,
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: SentinelTask

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

dataset_name
str = 'failed'
class nemo_curator.tasks.FileGroupTask(
dataset_name: str,
data: list[str] = list(),
_metadata: dict[str, typing.Any] = dict(),
reader_config: dict[str, typing.Any] = dict()
)
Dataclass

Bases: Task[list[str]]

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

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

Validate the task data.

class nemo_curator.tasks.ImageBatch(
dataset_name: str,
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: 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.

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

Validate the task data.

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

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

Bases: Task[Table | DataFrame]

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

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

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 | NoneDefaults to None

If provided, assign this sample_id to all new rows.

auto_position
boolDefaults to True

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

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.

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

nemo_curator.tasks.InterleavedBatch.delete_rows(
mask: pandas.Series

Delete rows where mask is True.

Parameters:

mask
pd.Series

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

nemo_curator.tasks.InterleavedBatch.get_columns() -> list[str]
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.

nemo_curator.tasks.InterleavedBatch.to_pandas() -> pandas.DataFrame
nemo_curator.tasks.InterleavedBatch.to_pyarrow() -> pyarrow.Table
nemo_curator.tasks.InterleavedBatch.validate() -> bool
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: {prefix}path, {prefix}member, {prefix}byte_offset, {prefix}byte_size, {prefix}frame_index.

class nemo_curator.tasks.LanceReadTask(
dataset_name: str,
_metadata: dict[str, typing.Any] = dict(),
path: str,
version: int
)
Dataclass

Bases: Task[FragmentIds]

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
FragmentIdsDefaults to list()

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)
nemo_curator.tasks.LanceReadTask.get_deterministic_id() -> str
nemo_curator.tasks.LanceReadTask.validate() -> bool
class nemo_curator.tasks.NoneTask(
dataset_name: str = 'none',
data: None = None,
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: SentinelTask

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

dataset_name
str = 'none'
class nemo_curator.tasks.SentinelTask(
dataset_name: str,
data: None = None,
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: Task[None]

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

data
None = None
num_items
int
nemo_curator.tasks.SentinelTask.__post_init__() -> None
nemo_curator.tasks.SentinelTask.validate() -> bool
class nemo_curator.tasks.Task(
dataset_name: str,
_metadata: dict[str, typing.Any] = dict()
)
DataclassAbstract

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_ids 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).

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

Post-initialization hook.

nemo_curator.tasks.Task.__repr__() -> str
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_ids (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.

nemo_curator.tasks.Task.add_stage_perf(
) -> None

Add performance stats for a stage.

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.

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.

nemo_curator.tasks.Task.validate() -> bool
abstract

Validate the task data.