How Zephyr Processes and Normalizes Raw Datasets for Training in Marin

Zephyr implements a lazy, declarative dataset API that turns raw file collections into a normalized, pipeline-ready representation through glob expansion, input-file specification, and format-specific reader dispatch.

The marin-community/marin repository includes Zephyr, a high-performance data loading library designed to streamline ML workflows. Understanding how Zephyr processes and normalizes raw datasets for training in Marin reveals a sophisticated architecture built on lazy evaluation, immutable transformations, and efficient I/O operations.

Lazy File Discovery via Glob Expansion

Zephyr begins normalization with lazy file discovery through the Dataset.from_files() entry point. This method accepts a glob pattern and instantiates a GlobSource object without immediate I/O.

At plan time, the internal resolve_glob() helper—located in lib/zephyr/src/zephyr/dataset.py (lines 54-72)—expands brace patterns using braceexpand and issues a single bulk list-objects RPC via fsspec.glob(detail=True). This operation builds a sorted list of FileEntry objects, each carrying both an InputFileSpec and file size metadata. If the glob matches no files and empty_glob_ok is set to false, Zephyr raises a FileNotFoundError (lines 73-75).

Declarative Input-File Specification

Each discovered file is wrapped in an InputFileSpec defined in lib/zephyr/src/zephyr/input_file.py (lines 14-27). This dataclass captures:

  • The file path and optional format hint (parquet, jsonl, vortex, or auto)
  • Column projection lists for selective reading
  • Row-range slicing parameters
  • Optional filter expressions for predicate pushdown

The auto format resolves to the correct reader based on file extension during the dispatch phase.

Reader Dispatch and Record Normalization

When a pipeline invokes load_file() or format-specific shortcuts like load_parquet(), Zephyr normalizes the specification via _as_spec() and determines the concrete format through _resolve_read_format() in lib/zephyr/src/zephyr/readers.py (lines 16-28).

Parquet Reading with Row-Group Iteration

For Parquet files, load_parquet() utilizes iter_parquet_row_groups to read column-projected, row-range-sliced data. This approach maintains O(row-group) memory usage rather than loading entire files, as implemented in lib/zephyr/src/zephyr/readers.py (lines 38-98).

JSONL Streaming

The load_jsonl() function streams files line-by-line, applying optional filter expressions and column projections during ingestion (lines 24-73). This ensures constant memory overhead regardless of file size.

Vortex Integration

For Vortex files, load_vortex() leverages the Vortex PyArrow dataset interface to enable push-down filtering and column projection before data reaches Python (lines 70-94).

Standardized Output Format

After reading, every record is yielded as a plain Python dict. If include_file_paths=True is specified, Zephyr injects the source file path under the configurable column name __file_path by default (lines 72-95).

Immutable Pipeline Construction

The Dataset class stores the source (GlobSource or iterable) alongside a list of logical operations including MapOp, FilterOp, and SelectOp. Each transformation method—map(), filter(), select(), or load_file()—returns a new Dataset instance with the operation appended, preserving immutability and enabling full plan inspection prior to execution (lines 59-66 and 78-98 in lib/zephyr/src/zephyr/dataset.py).

Execution and Write Normalization

The assembled plan executes via ZephyrContext.execute(dataset). The context iterates over shards, applies queued operations, and writes output through the unified WriteOp infrastructure. This ensures consistent normalization across all stages, converting arbitrary raw collections from local storage, GCS, S3, HuggingFace Hub, or zip archives into clean, column-aligned datasets.


# Create a dataset from a remote parquet glob

ds = (
    Dataset.from_files("gs://my-bucket/data/**/*.parquet")
           .load_parquet(columns=["text", "label"], approx_shard_bytes=128_000_000)
           .filter(lambda r: r["label"] == 1)          # keep only class 1

           .select("text")                             # projection

)

# Add the source file path for debugging

ds = ds.load_file(include_file_paths=True)

# Write normalized data as JSONL shards for training

ds = ds.write_jsonl("gs://my-bucket/processed/train-{shard:05d}.jsonl.gz")
ctx.execute(ds)   # triggers Zephyr's lazy execution

Summary

  • Lazy discovery: Dataset.from_files() uses resolve_glob() in lib/zephyr/src/zephyr/dataset.py to expand patterns efficiently via fsspec.
  • Structured specs: InputFileSpec in lib/zephyr/src/zephyr/input_file.py normalizes format hints, projections, and filters.
  • Format dispatch: _resolve_read_format() in lib/zephyr/src/zephyr/readers.py routes to optimized readers for Parquet, JSONL, and Vortex.
  • Memory efficiency: Parquet reading uses row-group iteration to bound memory usage.
  • Immutability: All transformations return new Dataset instances, enabling plan optimization before execution.

Frequently Asked Questions

What file formats does Zephyr support for dataset normalization?

Zephyr natively supports Parquet, JSONL, and Vortex formats through dedicated loaders in lib/zephyr/src/zephyr/readers.py. The auto format hint automatically detects the correct reader based on file extensions, simplifying pipeline configuration across heterogeneous data sources.

How does Zephyr handle large Parquet files efficiently?

Zephyr processes Parquet files using iter_parquet_row_groups (lines 38-98 in lib/zephyr/src/zephyr/readers.py), which iterates over row groups rather than loading entire files into memory. This ensures memory usage scales with row group size rather than file size, enabling processing of terabyte-scale datasets on modest hardware.

Can I track the source file path for each record during normalization?

Yes. Set include_file_paths=True when calling load_file() or format-specific loaders. Zephyr injects the source path under the __file_path column (configurable) for every record, as implemented in lib/zephyr/src/zephyr/readers.py (lines 72-95). This aids debugging and data lineage tracking.

What happens if my glob pattern matches no files?

If the glob pattern resolves to an empty set and empty_glob_ok is false (the default), Zephyr raises a FileNotFoundError during the planning phase (lines 73-75 in lib/zephyr/src/zephyr/dataset.py). Setting empty_glob_ok=True allows the pipeline to proceed with an empty dataset instead.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →