Marin-Zephyr Pipeline Stages: A Complete Guide to the 8 Core Stage Types
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 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.
# 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.
# 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.
# 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.
# 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.
# 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.
# 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, this stage materializes results for downstream training or analytics consumption.
# 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.
# 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:
- Intake / Normalize – Read raw data and emit canonical Parquet layouts
- Tokenization – Convert documents to fixed-size token sequences using MAP
- Embedding – Generate vector representations via neural models in MAP stages
- Classification – Attach per-document labels or scores through MAP operations
- Deduplication – Shuffle by document ID using GROUP_BY to drop exact or fuzzy duplicates
- Group-by / Sharding – Partition data for parallel downstream processing
- Sorted-Merge-Join – Combine datasets (e.g., embeddings + labels) using SORTED_MERGE_JOIN
- 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 provides the entry point for pipeline construction, while StepSpec objects capture stage definitions, dependencies, and execution logic.
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 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), automatically recording counters and throughput statistics during execution. Resource constraints are enforced via ZephyrTaskResources in 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: 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.
What file defines the StageType enumeration in Marin?
The StageType enumeration is defined in 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. 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.
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 →