# Marin-Zephyr Pipeline Stages: A Complete Guide to the 8 Core Stage Types

> Explore the 8 core Marin-Zephyr pipeline stages: MAP, FILTER, FLATMAP, REDUCE, GROUP_BY, SORTED_MERGE_JOIN, WRITE, and CUSTOM. Understand how ZephyrCoordinator orchestrates these steps.

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

---

**The Marin-Zephyr pipeline executes eight core stage types—MAP, FILTER, FLATMAP, REDUCE, GROUP_BY, SORTED_MERGE_JOIN, WRITE, and CUSTOM—each defined as a StepSpec object and orchestrated by the ZephyrCoordinator within a ZephyrContext.**

The marin-community/marin repository provides a flexible, actor-based data-processing engine that processes large-scale datasets through a series of logical operations. Understanding the Marin-Zephyr pipeline stages is essential for constructing efficient ETL workflows, from raw data ingestion to final dataset materialization.

## Core Stage Types in the Marin-Zephyr Pipeline

The **StageType** enumeration in [`lib/zephyr/src/zephyr/plan.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/plan.py) defines the fundamental operations available within the framework. Each stage type corresponds to a specific data transformation pattern executed by **Zephyr workers** under the coordination of the **ZephyrCoordinator**.

### MAP

The **MAP** stage performs element-wise transformations on input rows. This stage applies a function to each record independently, making it ideal for tokenization, embedding generation, and feature extraction.

```python

# Element-wise transformation example

ctx.map(dataset, lambda row: {"tokens": tokenize(row["text"])})

```

### FILTER

The **FILTER** stage drops rows that fail to satisfy a predicate condition. Use this stage to remove low-quality records or apply data quality constraints before downstream processing.

```python

# Predicate-based filtering

ctx.filter(dataset, lambda row: row["score"] > 0.5 and len(row["text"]) > 100)

```

### FLATMAP

The **FLATMAP** stage expands a single input row into multiple output rows. Unlike MAP, which maintains a 1:1 record ratio, FLATMAP handles sentence splitting, paragraph chunking, or exploding nested arrays into separate records.

```python

# Expanding documents into sentences

ctx.flat_map(dataset, lambda doc: [{"sentence": s} for s in split_sentences(doc["text"])])

```

### REDUCE

The **REDUCE** stage aggregates multiple rows into a single result. This stage implements accumulation patterns such as counting, summing, or generating summary statistics across partitions.

```python

# Summing values across dataset

ctx.reduce(dataset, lambda acc, row: acc + row["value"], init=0)

```

### GROUP_BY

The **GROUP_BY** stage shuffles rows by a specified key and writes partitioned output. This serves as the canonical **deduplication** mechanism, allowing the pipeline to group records by document ID or content hash before eliminating duplicates.

```python

# Shuffling by key for deduplication

ctx.group_by(dataset, key=lambda r: r["doc_id"])

```

### SORTED_MERGE_JOIN

The **SORTED_MERGE_JOIN** stage performs an efficient sorted-merge join of two keyed datasets. This operation is critical for consolidation workflows, such as combining embedding vectors with metadata labels or joining attribute tables.

```python

# Joining embeddings with labels

ctx.sorted_merge_join(left_dataset, right_dataset, on="id")

```

### WRITE

The **WRITE** stage persists the final dataset to storage, typically in Parquet format. Implemented in [`lib/zephyr/src/zephyr/writers.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/writers.py), this stage materializes results for downstream training or analytics consumption.

```python

# Materializing output

ctx.write_parquet(processed_dataset, "gs://bucket/output.parquet")

```

### CUSTOM

The **CUSTOM** stage allows user-defined Python callables wrapped as **StepSpec** objects. This extensibility point enables domain-specific logic that does not fit the standard stage taxonomy.

```python

# Custom stage definition

custom_step = StepSpec(
    name="domain_specific_processing",
    fn=my_custom_function,
    deps=[previous_step]
)

```

## Typical Pipeline Execution Flow

A production Marin-Zephyr pipeline typically follows this logical progression through the Marin-Zephyr pipeline stages:

1. **Intake / Normalize** – Read raw data and emit canonical Parquet layouts
2. **Tokenization** – Convert documents to fixed-size token sequences using MAP
3. **Embedding** – Generate vector representations via neural models in MAP stages
4. **Classification** – Attach per-document labels or scores through MAP operations
5. **Deduplication** – Shuffle by document ID using GROUP_BY to drop exact or fuzzy duplicates
6. **Group-by / Sharding** – Partition data for parallel downstream processing
7. **Sorted-Merge-Join** – Combine datasets (e.g., embeddings + labels) using SORTED_MERGE_JOIN
8. **Write** – Materialize the final dataset via WRITE stages for training or consumption

## Building Pipelines with StepSpec and ZephyrContext

The **ZephyrContext** class in [`lib/zephyr/src/zephyr/context.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/context.py) provides the entry point for pipeline construction, while **StepSpec** objects capture stage definitions, dependencies, and execution logic.

```python
from zephyr.context import ZephyrContext
from zephyr.coordinator import ZephyrCoordinator
from zephyr.plan import StepSpec

# Initialize context for local execution

ctx = ZephyrContext(name="demo", max_workers=4)

# Define pipeline stages as StepSpecs with explicit dependencies

read_step = StepSpec(
    name="read_parquet",
    fn=lambda _: ctx.read_parquet("gs://my-bucket/raw_data.parquet")
)

tokenize_step = StepSpec(
    name="tokenize",
    fn=lambda ds: ctx.map(ds, lambda row: {"tokens": tokenize(row["text"])}),
    deps=[read_step],
)

filter_step = StepSpec(
    name="filter_short",
    fn=lambda ds: ctx.filter(ds, lambda row: len(row["tokens"]) > 10),
    deps=[tokenize_step],
)

write_step = StepSpec(
    name="write_parquet",
    fn=lambda ds: ctx.write_parquet(ds, "gs://my-bucket/processed.parquet"),
    deps=[filter_step],
)

# Execute via coordinator

coordinator = ZephyrCoordinator(client=ctx.client)
plan = [read_step, tokenize_step, filter_step, write_step]
coordinator.run_pipeline(plan, run_id="demo-run", task_cost=1, max_task_cost=1)

```

The **ZephyrCoordinator** in [`lib/zephyr/src/zephyr/coordinator.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/coordinator.py) resolves the dependency graph, locks resources, and streams data through worker processes according to the **max_workers** limit specified in the context.

## Resource Management and Monitoring

Each Marin-Zephyr pipeline stage tracks performance metrics through **ZephyrStageStat** (defined in [`lib/zephyr/src/zephyr/stats.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/stats.py)), automatically recording counters and throughput statistics during execution. Resource constraints are enforced via **ZephyrTaskResources** in [`lib/zephyr/src/zephyr/stage_io.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/stage_io.py), allowing operators to specify CPU and memory limits per stage.

The pipeline supports parallel execution across multiple workers per stage, with the coordinator managing task distribution and fault tolerance through the **StageRunner** protocol.

## Summary

- The Marin-Zephyr pipeline defines eight core stage types in [`lib/zephyr/src/zephyr/plan.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/plan.py): MAP, FILTER, FLATMAP, REDUCE, GROUP_BY, SORTED_MERGE_JOIN, WRITE, and CUSTOM
- **StepSpec** objects declaratively define stages, their callable functions, and dependencies
- The **ZephyrCoordinator** orchestrates execution while respecting resource limits defined in **ZephyrTaskResources**
- **ZephyrStageStat** provides automatic per-stage metrics collection and monitoring
- Pipelines progress logically from intake through transformation, deduplication, joining, and final materialization

## Frequently Asked Questions

### What is the difference between MAP and FLATMAP in Marin-Zephyr?

**MAP** maintains a one-to-one record relationship, transforming each input row into exactly one output row. **FLATMAP** expands a single input row into zero or more output rows, making it suitable for operations like sentence splitting or exploding nested structures. Both are executed element-wise by Zephyr workers but differ in their output cardinality.

### How does the ZephyrCoordinator manage pipeline dependencies?

The **ZephyrCoordinator** analyzes the `deps` parameter of each **StepSpec** to build a directed acyclic graph (DAG) of stage dependencies. It schedules stages for execution only after all dependency stages complete, locking resources via **ZephyrTaskResources** and managing data streaming between stages through the **StageRunner** protocol defined in [`lib/zephyr/src/zephyr/stage_io.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/stage_io.py).

### What file defines the StageType enumeration in Marin?

The **StageType** enumeration is defined in [`lib/zephyr/src/zephyr/plan.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/plan.py) alongside the **PhysicalStage** and **StepSpec** structures. This file serves as the canonical reference for available pipeline operations and their serialization formats within the marin-community/marin repository.

### How are resources allocated per stage in the Marin-Zephyr pipeline?

Resources are allocated through **ZephyrTaskResources** specified in [`lib/zephyr/src/zephyr/stage_io.py`](https://github.com/marin-community/marin/blob/main/lib/zephyr/src/zephyr/stage_io.py). The coordinator enforces CPU and memory limits per stage while the **max_workers** parameter in **ZephyrContext** controls parallelism. The system tracks resource utilization via **ZephyrStageStat** to prevent worker overload during pipeline execution.