# How Zephyr Processes and Normalizes Raw Datasets for Training in Marin

> Learn how Zephyr processes and normalizes raw datasets for Marin training using a declarative API. Explore glob expansion, file specification, and format readers for efficient data pipelines.

- Repository: [The Marin Project/marin](https://github.com/marin-community/marin)
- Tags: how-to-guide
- Published: 2026-08-28

---

**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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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.

```python

# 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`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/dataset.py) to expand patterns efficiently via `fsspec`.
- **Structured specs**: `InputFileSpec` in [`lib/zephyr/src/zephyr/input_file.py`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/dataset.py)). Setting `empty_glob_ok=True` allows the pipeline to proceed with an empty dataset instead.