Implementation:NVIDIA NeMo Curator ParquetReader: Difference between revisions
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:
- FilePartitioningStage - discovers and groups files into partitions
- 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:
- FilePartitioningStage handles file discovery, globbing, and grouping
- 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
- Implementation:NVIDIA_NeMo_Curator_BaseReader - Abstract base class that ParquetReaderStage extends
- Implementation:NVIDIA_NeMo_Curator_JSONLReader - Analogous reader for JSONL format files
- Implementation:NVIDIA_NeMo_Curator_ParquetWriter - Writer stage for Parquet format output
- Implementation:NVIDIA_NeMo_Curator_FilePartitioningStage - Stage that handles file discovery and grouping
- Environment:NVIDIA_NeMo_Curator_Python_Linux_Base