# How to Use Marin-Zephyr for Dataset Processing: A Complete Guide

> Learn how to use Marin Zephyr for efficient dataset processing. This guide covers building scalable ETL pipelines with a high-level API. No worker orchestration needed.

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

---

**Marin-Zephyr is a lightweight, distributed dataset-processing framework built on Marin and Iris that provides a high-level API for scalable ETL pipelines without requiring manual worker orchestration.**

Marin-Zephyr enables engineers to process petabyte-scale datasets across distributed clusters using a simple Python API. Hosted in the `marin-community/marin` repository, it abstracts away complex worker management while preserving fine-grained control over memory budgets and resource allocation. This guide demonstrates how to leverage Zephyr’s immutable datasets and functional transformations to build production-grade data pipelines.

## Core Concepts

Understanding Zephyr’s architecture begins with its five primary abstractions.

**`ZephyrContext`** serves as the entry point for every pipeline. It creates a **coordinator** actor, spawns a pool of **worker** actors via Iris when running inside an Iris job (or a local client otherwise), and drives pipeline execution. It also handles shared-data upload and memory-budget configuration. The implementation resides in [`lib/zephyr/src/zephyr/context.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/context.py).

**`Dataset`** represents an immutable collection of source items—either files or glob patterns. Each item becomes a Zephyr **task** scheduled on a worker. This class abstracts file-system details and provides helpers for common formats via `load_jsonl` and `load_parquet`. See [`lib/zephyr/src/zephyr/dataset.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/dataset.py).

**`Coordinator`** runs on the driver node, builds an execution **plan** (stage graph), distributes tasks to workers, monitors progress, and aggregates results into a `ZephyrExecutionResult`. The logic is defined in [`lib/zephyr/src/zephyr/coordinator.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/coordinator.py).

**`Worker`** is a lightweight Iris actor (or local subprocess) that executes a single task: reading input shards, applying a user-provided transform, and writing output. Workers auto-scale based on `max_workers` and reuse shared objects via the memory store. Source: [`lib/zephyr/src/zephyr/worker.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/worker.py).

**`MemoryStore`** provides an optional on-disk cache that enables workers to share large immutable objects—such as tokenizers or model weights—without re-uploading. Implementation details are in [`lib/zephyr/src/zephyr/memory_store.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/memory_store.py).

## Execution Flow

Every Marin-Zephyr pipeline follows a five-stage execution model.

1. **Create a `ZephyrContext`** – specify `max_workers` or let Zephyr auto-detect the Iris client.
2. **Define a `Dataset`** – point it at a glob pattern or list of files.
3. **Compose a pipeline** – chain `map`, `filter`, `shuffle`, or `groupby` using the functional API provided by the `Dataset` object.
4. **Call `ctx.execute(pipeline)`** – the coordinator builds a plan (via [`plan.py`](https://github.com/marin-community/marin/blob/main/plan.py)), splits the dataset into shards, and dispatches each shard to a worker.
5. **Collect results** – `execute` returns a `ZephyrExecutionResult` containing final outputs, counters, and optional profiling data.

The framework automatically handles **chunking and row-group alignment** for Parquet files (see [`lib/zephyr/src/zephyr/parquet_scan.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/parquet_scan.py)), **memory budgeting** via [`memory_budget.py`](https://github.com/marin-community/marin/blob/main/memory_budget.py) to prevent OOM errors, and **retry with idempotency** through the `skip_existing` flag.

## Practical Code Examples

The following patterns demonstrate common Marin-Zephyr workflows using the actual public API.

### Processing JSON-L Files

This example reads raw JSON-L files, extracts and transforms fields, and writes the results back to cloud storage.

```python
from zephyr import Dataset, ZephyrContext

# 1. Create a context with a generous worker pool.

ctx = ZephyrContext(max_workers=200)

# 2. Build a dataset pointing at a directory of *.jsonl files.

data = Dataset.from_glob("gs://my-bucket/raw/*.jsonl")

# 3. Define a map stage that extracts a field and lower-cases it.

def extract_title(record):
    return {"title": record["title"].lower()}

processed = data.map(extract_title)

# 4. Write the transformed records back to GCS as JSON-L.

processed.write_jsonl("gs://my-bucket/processed/titles.jsonl")

# 5. Run the pipeline.

result = ctx.execute(processed)
print("Wrote", result.counters["records_written"], "records")

```

**Key references:** `Dataset.from_glob` (defined in [`dataset.py`](https://github.com/marin-community/marin/blob/main/dataset.py)), the `map` implementation (via [`plan.py`](https://github.com/marin-community/marin/blob/main/plan.py)), and `write_jsonl` (in [`writers.py`](https://github.com/marin-community/marin/blob/main/writers.py)).

### Shuffling Parquet with Row-Group Alignment

Zephyr’s Parquet reader aligns splits on row groups to maintain data locality during wide transformations like shuffles.

```python
from zephyr import Dataset, ZephyrContext

ctx = ZephyrContext(max_workers=128)

# Load a set of Parquet files; Zephyr aligns splits on row groups.

data = Dataset.from_glob("gs://my-bucket/parquet/*.parquet")

# Shuffle the dataset (wide dependency) and write the result.

shuffled = data.shuffle()
shuffled.write_parquet("gs://my-bucket/shuffled/output.parquet")

ctx.execute(shuffled)

```

**Internals:** The `shuffle` method creates a `StageType.SHUFFLE` node, while the coordinator uses `parquet_scan.iter_parquet_row_groups` to emit row-group-aligned chunks, as implemented in [`lib/zephyr/src/zephyr/parquet_scan.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/parquet_scan.py).

### Sharing Large Objects via MemoryStore

Use `upload_shared` to distribute large immutable objects to workers without serializing them per-task.

```python
from zephyr import Dataset, ZephyrContext
from transformers import AutoTokenizer

ctx = ZephyrContext(max_workers=64)

# Upload the tokenizer once; workers will lazy-load it.

tokenizer = AutoTokenizer.from_pretrained("bert-base-uncased")
ctx.upload_shared("tokenizer", tokenizer)   # internally uses MemoryStore

def tokenize(record):
    toks = ctx.get_shared("tokenizer").encode(record["text"])
    return {"tokens": toks}

data = Dataset.from_glob("gs://my-bucket/text/*.jsonl")
tokenized = data.map(tokenize)
tokenized.write_jsonl("gs://my-bucket/tokenized/output.jsonl")

ctx.execute(tokenized)

```

**Key points:** `ctx.upload_shared` interacts with [`lib/zephyr/src/zephyr/memory_store.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/memory_store.py), while workers retrieve objects via `ctx.get_shared`.

### Fine-Grained Resource Configuration

Specify heterogeneous resource requirements per pipeline stage to optimize for CPU-heavy vs. IO-heavy tasks.

```python
from zephyr import Dataset, ZephyrContext, ResourceConfig

# Define per-stage resources (e.g., more CPU for tokenization, less for shuffling).

stage_cfg = {
    0: ResourceConfig(cpu=8, memory_gb=32),   # tokenization stage

    1: ResourceConfig(cpu=4, memory_gb=16),   # shuffle stage

}
ctx = ZephyrContext(max_workers=200, stage_resources=stage_cfg)

data = Dataset.from_glob("gs://my-bucket/raw/*.parquet")

# ... pipeline definition ...

ctx.execute(pipeline)

```

**Implementation:** `ZephyrContext` forwards `stage_resources` to the coordinator, which configures each Iris worker group accordingly (see [`worker_group_race.py`](https://github.com/marin-community/marin/blob/main/worker_group_race.py) in the test suite for heterogeneous resource examples).

## Key Implementation Files

The following source files define the architecture and API referenced throughout this guide:

- [`lib/zephyr/src/zephyr/context.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/context.py) – Central context, client detection, and worker-pool handling.
- [`lib/zephyr/src/zephyr/dataset.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/dataset.py) – Immutable dataset abstraction, glob resolution, and shard creation.
- [`lib/zephyr/src/zephyr/coordinator.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/coordinator.py) – Execution plan builder, task dispatcher, and result aggregator.
- [`lib/zephyr/src/zephyr/worker.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/worker.py) – Worker actor implementing the read-transform-write loop.
- [`lib/zephyr/src/zephyr/readers.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/readers.py) – Built-in readers for JSONL, Parquet, and CSV.
- [`lib/zephyr/src/zephyr/writers.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/writers.py) – Output writers for the same formats.
- [`lib/zephyr/src/zephyr/memory_store.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/memory_store.py) – Shared-object cache for large immutable assets.
- [`lib/zephyr/src/zephyr/parquet_scan.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/parquet_scan.py) – Row-group-aligned Parquet sharding utilities.

## Summary

- **Marin-Zephyr** provides a functional API for distributed dataset processing on top of the Marin and Iris stacks.
- **`ZephyrContext`** manages worker pools and shared resources, while **`Dataset`** offers immutable, composable transformations.
- **Row-group alignment** for Parquet and **memory budgeting** are handled automatically by the framework internals.
- **Per-stage resource configurations** allow fine-tuning CPU and memory allocation for heterogeneous workloads.
- Source code for all core components resides in `lib/zephyr/src/zephyr/`, with entry points in [`context.py`](https://github.com/marin-community/marin/blob/main/context.py) and [`dataset.py`](https://github.com/marin-community/marin/blob/main/dataset.py).

## Frequently Asked Questions

### What file formats does Marin-Zephyr support?

Marin-Zephyr natively supports **JSON-L**, **Parquet**, and **CSV** through pluggable readers and writers. The `Dataset` class provides `load_jsonl`, `write_parquet`, and similar methods that internally use the modules in [`lib/zephyr/src/zephyr/readers.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/readers.py) and [`lib/zephyr/src/zephyr/writers.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/writers.py). You can extend support by implementing custom reader functions that conform to the Zephyr record API.

### How does Marin-Zephyr handle worker failures?

The coordinator in [`lib/zephyr/src/zephyr/coordinator.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/coordinator.py) monitors task completion via Iris heartbeats. Failed tasks are automatically retried up to a configurable limit, and the `skip_existing` flag ensures idempotent execution—workers can skip shards that have already been processed. This design prevents data duplication during transient infrastructure failures.

### Can I use Marin-Zephyr without the Iris cluster manager?

Yes. While Zephyr is optimized for the Iris actor runtime, it can operate in **local mode** by spawning subprocess workers instead of distributed actors. When `ZephyrContext` detects the absence of an Iris client, it falls back to a local pool defined in [`lib/zephyr/src/zephyr/context.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/context.py), allowing you to test pipelines on single machines before deploying to a cluster.

### How does the MemoryStore improve pipeline performance?

The **MemoryStore** eliminates redundant serialization of large objects across tasks. When you call `ctx.upload_shared()`, the object is written once to an on-disk cache (as implemented in [`lib/zephyr/src/zephyr/memory_store.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/memory_store.py)). Workers then lazy-load the object by reference rather than value, significantly reducing network overhead and memory pressure when sharing tokenizers or model weights across hundreds of distributed tasks.