Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:NVIDIA NeMo Curator ParquetReader: Difference between revisions

From Leeroopedia
Auto-imported from implementations/NVIDIA_NeMo_Curator_ParquetReader.md
 
Sync from local file
 
Line 176: Line 176:
== Related Pages ==
== Related Pages ==


* [[NVIDIA_NeMo_Curator_BaseReader]] - Abstract base class that ParquetReaderStage extends
* [[Implementation:NVIDIA_NeMo_Curator_BaseReader]] - Abstract base class that ParquetReaderStage extends
* [[NVIDIA_NeMo_Curator_JSONLReader]] - Analogous reader for JSONL format files
* [[Implementation:NVIDIA_NeMo_Curator_JSONLReader]] - Analogous reader for JSONL format files
* [[NVIDIA_NeMo_Curator_ParquetWriter]] - Writer stage for Parquet format output
* [[Implementation:NVIDIA_NeMo_Curator_ParquetWriter]] - Writer stage for Parquet format output
* [[NVIDIA_NeMo_Curator_FilePartitioningStage]] - Stage that handles file discovery and grouping
* [[Implementation:NVIDIA_NeMo_Curator_FilePartitioningStage]] - Stage that handles file discovery and grouping
* [[requires_env::Environment:NVIDIA_NeMo_Curator_Python_Linux_Base]]
* [[requires_env::Environment:NVIDIA_NeMo_Curator_Python_Linux_Base]]


[[Category:Implementations]]
[[Category:Implementations]]

Latest revision as of 10:48, 27 September 2026

Knowledge Sources
Domains Data Ingestion, IO, Parquet, Data Pipeline
Last Updated 2026-02-14 00:00 GMT

Overview

Provides Parquet file reading capabilities with a low-level ParquetReaderStage and a high-level ParquetReader composite stage for ingesting Parquet data into the NeMo Curator pipeline.

Description

This module contains two classes that implement Parquet file reading at different abstraction levels:

ParquetReaderStage extends BaseReader and implements the read_data abstract method. It reads each Parquet file path using pd.read_parquet with PyArrow as the default engine and pyarrow as the default dtype backend. When fields is specified, only those columns are read by passing them as the columns parameter to pd.read_parquet. All resulting DataFrames are concatenated with pd.concat using ignore_index=True.

ParquetReader is a CompositeStage[_EmptyTask, DocumentBatch] dataclass that provides a high-level interface for reading Parquet files. It follows the same two-stage decomposition pattern as the JSONL reader:

  1. FilePartitioningStage - discovers and groups files into partitions
  2. ParquetReaderStage - reads each file group into a DocumentBatch

Parquet is the default and most common format used throughout the NeMo Curator deduplication and embedding workflows.

Usage

Use ParquetReader as the primary entry point for reading Parquet files in a pipeline. Use ParquetReaderStage directly when you already have FileGroupTask inputs from an upstream FilePartitioningStage or custom file grouping logic.

Code Reference

Source Location

  • Repository: NeMo-Curator
  • File: nemo_curator/stages/text/io/reader/parquet.py
  • Lines: 1-131

Signature

@dataclass
class ParquetReaderStage(BaseReader):
    name: str = "parquet_reader"

    def read_data(
        self,
        paths: list[str],
        read_kwargs: dict[str, Any] | None = None,
        fields: list[str] | None = None,
    ) -> pd.DataFrame: ...


@dataclass
class ParquetReader(CompositeStage[_EmptyTask, DocumentBatch]):
    file_paths: str | list[str]
    files_per_partition: int | None = None
    blocksize: int | str | None = None
    fields: list[str] | None = None
    read_kwargs: dict[str, Any] | None = None
    file_extensions: list[str] = field(default_factory=...)
    task_type: Literal["document", "image", "video", "audio"] = "document"
    _generate_ids: bool = False
    _assign_ids: bool = False
    name: str = "parquet_reader"

Import

from nemo_curator.stages.text.io.reader.parquet import ParquetReaderStage, ParquetReader

I/O Contract

Inputs (ParquetReader)

Name Type Required Description
file_paths str or list[str] Yes Path or list of paths to Parquet files or directories containing Parquet files
files_per_partition int or None No Number of files to group per partition (default: None)
blocksize int or str or None No Target block size for file partitioning, e.g. "256MB" (default: None)
fields list[str] or None No If specified, only read these columns from the Parquet files (default: None)
read_kwargs dict[str, Any] or None No Additional keyword arguments passed to pd.read_parquet (default: None)
file_extensions list[str] No File extensions to match when discovering files (default: Parquet extensions)
task_type Literal["document", ...] No Type of task; only "document" is currently supported (default: "document")
_generate_ids bool No Whether to generate monotonically increasing deduplication IDs (default: False)
_assign_ids bool No Whether to assign pre-computed deduplication IDs (default: False)

Inputs (ParquetReaderStage)

Name Type Required Description
FileGroupTask FileGroupTask Yes Task containing a list of Parquet file paths to read

Outputs

Name Type Description
DocumentBatch DocumentBatch Contains the concatenated data from all Parquet files in the group as a Pandas DataFrame with PyArrow dtypes

Usage Examples

Basic Usage with ParquetReader

from nemo_curator.stages.text.io.reader.parquet import ParquetReader

# Read Parquet files from a directory
reader = ParquetReader(
    file_paths="/data/corpus/",
    files_per_partition=5,
    fields=["text", "url", "language"],
)

With Custom Read Settings

from nemo_curator.stages.text.io.reader.parquet import ParquetReader

reader = ParquetReader(
    file_paths="/data/embeddings/",
    blocksize="512MB",
    read_kwargs={
        "engine": "pyarrow",
        "storage_options": {"key": "...", "secret": "..."},
    },
)

Using ParquetReaderStage Directly

from nemo_curator.stages.text.io.reader.parquet import ParquetReaderStage

reader_stage = ParquetReaderStage(
    fields=["text", "embeddings"],
    _generate_ids=True,
)

Implementation Details

Default Read Settings

ParquetReaderStage applies the following defaults when not overridden by read_kwargs:

  • engine: "pyarrow" - uses the Apache Arrow Parquet reader
  • dtype_backend: "pyarrow" - uses PyArrow-backed nullable dtypes in Pandas for better type fidelity and memory efficiency
  • columns: Set to fields when specified, enabling column pruning at the Parquet level for reduced I/O

Composite Decomposition

ParquetReader.decompose() creates a two-stage pipeline following the same pattern as JsonlReader:

  1. FilePartitioningStage handles file discovery, globbing, and grouping
  2. ParquetReaderStage receives FileGroupTask inputs and reads the actual Parquet data

Storage options from read_kwargs are propagated to the FilePartitioningStage for remote filesystem access.

Task Type Restriction

Currently, only task_type="document" is supported. Passing other values raises a NotImplementedError.

Related Pages