This release is a pre-release and may not be stable for production use.
Zephyr
Simple data processing library for Marin pipelines. Build lazy dataset pipelines that run on Iris jobs or a local backend.
Quick Start
from zephyr.context import ZephyrContext
from zephyr.dataset import Dataset
from zephyr.readers import load_jsonl
# Read, transform, write
ctx = ZephyrContext(max_workers=100)
pipeline = (
Dataset.from_files("gs://input/", "**/*.jsonl.gz")
.flat_map(load_jsonl)
.filter(lambda x: x["score"] > 0.5)
.map(lambda x: transform_record(x))
.write_jsonl("gs://output/data-{shard:05d}-of-{total:05d}.jsonl.gz")
)
ctx.execute(pipeline)
Key Patterns
Dataset Creation:
Dataset.from_files(path, pattern)- glob filesDataset.from_list(items)- explicit list
Loading Files
.load_{file,parquet,jsonl,vortex}- load rows from a file
Transformations:
.map(fn)- transform each item.flat_map(fn)- expand items (e.g.,load_jsonl).filter(fn)- filter items by function or expression.select(columna, columnb)- select out the given columns.window(n)- group into batches.reshard(n)- redistribute across n shards
Output:
.write_jsonl(pattern)- write JSONL (gzip if.gz).write_parquet(pattern, schema)- write zstd Parquet with page indexes and at most 256 rows per data page.write_vortex(pattern)- write to a Vortex file
Execution (ZephyrContext):
ZephyrContext(max_workers=N)— auto-detects the backend (Iris inside an Iris job, local otherwise) viafray.current_client()ZephyrContext(client=LocalClient())— explicit local backend (testing)ctx.execute(pipeline)— runs the pipeline; returns aZephyrExecutionResult(results, counters)
Read-only memory stores
ZephyrContext.load_memory_store() loads an existing partitioned dataset into
the workers of an entered context. Every worker starts an empty multi-table
service; pipelines that do not load a table perform no source reads. The table
handle is picklable, so later pipelines and child jobs can use get() or
order-preserving get_many() lookups without copying the table into every task.
from fray.types import ResourceConfig
def document_partition(key: tuple[int, str]) -> int:
file_index, _ = key
return file_index
documents = Dataset.from_files("s3://bucket/documents/*.parquet").load_parquet().map(
lambda row: ((row["file_index"], row["id"]), row["text"])
)
with ZephyrContext(
max_workers=16,
resources=ResourceConfig(cpu=2, ram="8g"),
) as ctx:
document_store = ctx.load_memory_store(
documents,
name="documents",
hash_key=document_partition,
recovery_timeout=900,
)
result = ctx.execute(Dataset.from_list(document_keys).map(document_store.get))
document_store.destroy()
For P source shards, every key must already satisfy
hash_key(key) % P == source_shard_index. Construction checks every row and
does not insert a shuffle. Readers and shard-local maps can load directly.
Persist and reload the output of a shuffle, join, reshard, reduce, or write
before constructing a store.
Keys must be hashable and unique. Keys, values, and the hash function must be
picklable for remote calls; Python's salted hash() is not stable enough to
serve as the partition function for string or byte keys. Actors retain the
loaded Python objects directly, and store.stats() reports item counts and
load time. Invalid input fails the load call without consuming actor restart
retries. Multiple tables can share the worker process. Zephyr does not reserve,
limit, or evict table memory; size the context's worker RAM for the combined
tables and pipeline workload.
Iris reconstructs a preempted worker at the same endpoint. The first table
lookup on that replacement reloads its immutable source shards, and all worker
responses and reconstruction share one recovery_timeout deadline. Destroying
one table leaves other tables active. Exiting the creating context stops the
worker pool and invalidates its table handles.
Real Usage
Wikipedia Processing:
from zephyr.context import ZephyrContext
from zephyr.dataset import Dataset
from zephyr.readers import load_jsonl
ctx = ZephyrContext(max_workers=100)
pipeline = (
Dataset.from_list(files)
.load_jsonl()
.map(lambda row: process_record(row, config))
.filter(lambda x: x is not None)
.write_jsonl(f"{output}/data-{{shard:05d}}-of-{{total:05d}}.jsonl.gz")
)
ctx.execute(pipeline)
Dataset Sampling:
from zephyr.context import ZephyrContext
from zephyr.dataset import Dataset
ctx = ZephyrContext(max_workers=1000)
pipeline = (
Dataset.from_files(input_path, "**/*.jsonl.gz")
.map(lambda path: sample_file(path, weights))
.write_jsonl(f"{output}/sampled-{{shard:05d}}.jsonl.gz")
)
ctx.execute(pipeline)
Parallel Downloads:
from zephyr.context import ZephyrContext
from zephyr.dataset import Dataset
tasks = [(config, fs, src, dst) for src, dst in file_pairs]
ctx = ZephyrContext(max_workers=32)
pipeline = Dataset.from_list(tasks).map(lambda t: download(*t))
ctx.execute(pipeline)
Installation
# From Marin monorepo
uv sync
# Standalone
cd lib/zephyr
uv pip install -e .
Running Tests
Zephyr tests run against multiple execution backends to ensure correctness across different environments.
All Tests on Both Backends (Default)
uv run pytest lib/zephyr/tests
# Runs all tests on both Local and Iris backends
# Local Iris cluster is started automatically via ClusterManager
Run Specific Backend Only
uv run pytest lib/zephyr/tests -k "local"
uv run pytest lib/zephyr/tests -k "iris"
The Iris cluster is started once per test session and reused across all tests for efficiency.
Design
Zephyr consolidates ad-hoc distributed and Hugging Face dataset processing patterns in Marin into a simple abstraction.
Key Features:
- Lazy evaluation with operation fusion
- Disk-based inter-stage data flow for low memory footprint
- Chunk-by-chunk streaming to minimize memory pressure
- Distributed execution with bounded parallelism (Iris/local backends)
- Automatic chunking to prevent large object overhead
- fsspec integration (GCS, S3, local)
- Type-safe operation chaining
See AGENTS.md for execution internals and source layout.
Release files for marin-zephyr 0.2.126.dev35987878076
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| marin_zephyr-0.2.126.dev35987878076.tar.gz | 117.0 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| marin_zephyr-0.2.126.dev35987878076-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 247.9 kB
Release files / marin_zephyr-0.2.126.dev35987878076.tar.gz
| Download URL | marin_zephyr-0.2.126.dev35987878076.tar.gz |
|---|---|
| Size | 117.0 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
e70043f77d59dea5619cbbb2c7999d97e7efed425a094a07d8ad8ac380eefd3b
|
|
BLAKE2b-256 checksum How to use checksums |
1e0476d2aa1d80286f3a5ad6c5de8b1d9d3508e184a38e1edb1a52e5314f5c52
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 24, 2026.
Transparency logRelease files / marin_zephyr-0.2.126.dev35987878076-py3-none-any.whl
| Download URL | marin_zephyr-0.2.126.dev35987878076-py3-none-any.whl |
|---|---|
| Size | 131.0 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
483276884b79409451648ef876c040cc70993eb99fc902a483ff37106d6c05f9
|
|
BLAKE2b-256 checksum How to use checksums |
4cc69d1e0a6d51cdfab7aa95f88d612960ccfe5f39c6dad2d0d564405a4a5cf5
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 24, 2026.
Transparency log