How to Use Marin-Zephyr for Dataset Processing: A Complete Guide
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.
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.
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.
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.
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.
Execution Flow
Every Marin-Zephyr pipeline follows a five-stage execution model.
- Create a
ZephyrContext– specifymax_workersor let Zephyr auto-detect the Iris client. - Define a
Dataset– point it at a glob pattern or list of files. - Compose a pipeline – chain
map,filter,shuffle, orgroupbyusing the functional API provided by theDatasetobject. - Call
ctx.execute(pipeline)– the coordinator builds a plan (viaplan.py), splits the dataset into shards, and dispatches each shard to a worker. - Collect results –
executereturns aZephyrExecutionResultcontaining 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), memory budgeting via 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.
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), the map implementation (via plan.py), and write_jsonl (in 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.
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.
Sharing Large Objects via MemoryStore
Use upload_shared to distribute large immutable objects to workers without serializing them per-task.
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, 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.
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 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– Central context, client detection, and worker-pool handling.lib/zephyr/src/zephyr/dataset.py– Immutable dataset abstraction, glob resolution, and shard creation.lib/zephyr/src/zephyr/coordinator.py– Execution plan builder, task dispatcher, and result aggregator.lib/zephyr/src/zephyr/worker.py– Worker actor implementing the read-transform-write loop.lib/zephyr/src/zephyr/readers.py– Built-in readers for JSONL, Parquet, and CSV.lib/zephyr/src/zephyr/writers.py– Output writers for the same formats.lib/zephyr/src/zephyr/memory_store.py– Shared-object cache for large immutable assets.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.
ZephyrContextmanages worker pools and shared resources, whileDatasetoffers 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 incontext.pyanddataset.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 and 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 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, 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). 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.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →