nemo_curator.tasks
nemo_curator.tasks
Submodules
nemo_curator.tasks.audio_tasknemo_curator.tasks.documentnemo_curator.tasks.file_groupnemo_curator.tasks.imagenemo_curator.tasks.interleavednemo_curator.tasks.lancenemo_curator.tasks.ocrnemo_curator.tasks.sentinelsnemo_curator.tasks.tasksnemo_curator.tasks.utilsnemo_curator.tasks.video
Package Contents
Classes
API
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:
Manifest entry dict (e.g. {"audio_filepath": "...", "text": "..."}).
Optional key whose value is validated as an existing path.
Validate the task data.
Bases: Task[Table | DataFrame]
Task for processing batches of text documents. Documents are stored as a dataframe (PyArrow Table or Pandas DataFrame).
Get the number of documents in this batch.
Get column names from the data.
Convert data to Pandas DataFrame.
Convert data to PyArrow table.
Validate the task data.
Bases: SentinelTask
Seeds a pipeline with task_id="0" — the implicit root every task
descends from (so all ids share the "0" prefix).
Bases: SentinelTask
Marks a slot as failed → retried on resume (counter does NOT decrement).
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.
Number of files in this group.
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).
Validate the task data.
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.
Number of images in this batch.
Validate the task data.
Represents a single image with metadata.
Aesthetic quality score as float
Image embedding vector as numpy array
Raw image pixel data as numpy array (H, W, C) in RGB format
Unique identifier for the image
Path to the image file on disk
Additional metadata associated with the image
NSFW probability score as float
Bases: Task[Table | DataFrame]
Task carrying row-wise multimodal records.
See module docstring for the full schema reference (reserved vs user columns).
Number of unique samples (distinct sample_id values).
Add rows to this task.
Parameters:
New rows to append. Must contain required columns unless overridden by sample_id / auto_position.
If provided, assign this sample_id to all new rows.
If True, auto-assign position values
continuing from the existing maximum per sample.
Build a source_ref JSON locator string.
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
Delete rows where mask is True.
Parameters:
Boolean Series aligned to the data. True marks a row
for deletion.
Parse a source_ref JSON string into a locator dict.
Return a DataFrame copy with parsed source_ref columns added.
Columns: {prefix}path, {prefix}member, {prefix}byte_offset,
{prefix}byte_size, {prefix}frame_index.
Bases: Task[FragmentIds]
Task containing Lance fragment ids assigned to one read partition.
Parameters:
Path or URI of the Lance dataset to read.
Lance dataset version to read.
FragmentIds — Lance fragment ids to read.
Bases: SentinelTask
Marks a slot as intentionally filtered (resumability counter decrements).
Bases: Task[None]
Base for payload-less marker tasks: no data, framework-assigned task_id.
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.
Free-form metadata carried alongside the task.
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.
Per-stage perf stats this task has accumulated.
The task’s payload (modality-specific).
Name of the dataset this task belongs to.
Get the number of items in this task.
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).
Post-initialization hook.
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:
task_id of the parent. An empty string
(an unassigned / EmptyTask parent) is dropped so it doesn’t
contribute a leading "_" to the path.
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.
Add performance stats for a stage.
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.
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.
Validate the task data.