# Marin End-to-End Pipeline: Core Components and Architecture Guide

> Explore the Marin end-to-end pipeline, detailing core components like Datakit, Zephyr, and Iris for data ingestion, distributed training, and model evaluation.

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

---

**Marin's end-to-end pipeline integrates seven specialized sub-projects—Datakit, Zephyr, Levanter, Fray, Iris, VLLM, and Marin CLI—that collectively handle data ingestion, distributed JAX training, cluster orchestration, model evaluation, and result publication.**

The Marin end-to-end pipeline provides a complete workflow for training and evaluating large language models, separating concerns into discrete, composable components. According to the `marin-community/marin` repository, this architecture enables reproducible machine learning experiments by chaining together data processing, distributed computation, and automated reporting into a linear execution flow.

## Data Ingestion and Versioned Storage with Datakit

**Datakit** handles raw data collection, sharding, and versioned storage for the Marin ecosystem. It provides utilities for sampling, hashing, and checkpointing datasets, ensuring that training runs consume immutable, reproducible data artifacts.

In [`tests/datakit/test_harrier_pipeline.py`](https://github.com/marin-community/marin/blob/main/tests/datakit/test_harrier_pipeline.py), Datakit demonstrates its ability to manage dataset shards across distributed storage backends. The library exposes methods for hashing data streams and maintaining versioned checkpoints, which downstream components consume via standardized APIs.

```python
from datakit.dataset import MyDataset

# Load a versioned dataset from Datakit storage

dataset = MyDataset.from_datakit("my_corpus")

```

## Dataset Processing and Normalization via Zephyr

**Zephyr** streams, parses, and prepares raw data for downstream consumption. It specializes in processing large unstructured sources—such as Wikipedia dumps—normalizing formats, extracting features, and writing cleaned tensors to Datakit containers.

This component sits between raw storage and the training loop, ensuring that Levanter receives consistently formatted input tensors. Documentation in [`docs/tutorials/storage-bucket.md`](https://github.com/marin-community/marin/blob/main/docs/tutorials/storage-bucket.md) outlines how Zephyr manages data locality and preprocessing pipelines.

## Training Orchestration with Levanter

**Levanter** serves as the core JAX-based training library, implementing the central `Pipeline` abstraction that stitches together data loading, model definition, and optimization. The `levanter.pipeline.Pipeline` class defined in [`lib/levanter/src/levanter/pipeline.py`](https://github.com/marin-community/marin/blob/main/lib/levanter/src/levanter/pipeline.py) acts as the primary coordination point for training workflows.

The Pipeline class accepts dataset references, model configurations, and trainer objects, orchestrating the forward-backward loops and checkpointing strategies. It integrates with Fray's distributed primitives to scale across multiple accelerators while maintaining deterministic behavior.

```python
from levanter.pipeline import Pipeline
from levanter.trainer import Trainer
from levanter.model import MyLM

def make_pipeline():
    # Assemble training pipeline with Datakit source

    dataset = MyDataset.from_datakit("my_corpus")
    model = MyLM(vocab_size=50257, hidden_dim=4096)
    
    return Pipeline(
        dataset=dataset,
        model=model,
        trainer=Trainer(learning_rate=1e-4, optimizer="adamw"),
        max_steps=100_000,
    )

```

## Distributed Execution through Fray

**Fray** provides scalable, multi-node execution primitives that enable Levanter to run large-scale training jobs. It implements parameter server architectures and mesh-aware sharding strategies, abstracting the complexity of distributed JAX computation across GPU clusters.

Located in `lib/fray`, this component manages inter-node communication and memory sharding, allowing the Levanter pipeline to scale beyond single-node limitations without modifying user-level training code.

## Cluster Management and Job Scheduling with Iris

**Iris** manages job submission, monitoring, and lifecycle orchestration on cloud clusters including CoreWeave and GCP. It ensures reproducible runs by handling resource allocation, log aggregation, and automatic cleanup of ephemeral infrastructure.

Configuration details in [`lib/iris/OPS.md`](https://github.com/marin-community/marin/blob/main/lib/iris/OPS.md) specify how Iris wraps containerized execution, injecting environment variables and mounting Datakit volumes into training pods. The CLI interface enables one-line submission of complex distributed jobs.

```bash

# Submit a distributed training job to an Iris-managed cluster

iris submit \
  --job-name my-lm-train \
  --image marin/levanter:latest \
  --resources gpu=8 \
  -- python -m my_pipeline.run

```

## Evaluation and Inference using VLLM

After training completes, models transition to the **VLLM** integration for fast inference and benchmark evaluation. This component serves trained checkpoints through an optimized inference engine, exposing REST endpoints for language model evaluation pipelines.

The evaluation workflow defined in [`experiments/evaluation/pipeline.py`](https://github.com/marin-community/marin/blob/main/experiments/evaluation/pipeline.py) demonstrates how Iris-managed training artifacts flow into VLLM serving instances for automated benchmarking against standard LM evaluation suites.

```bash

# Serve a trained checkpoint for evaluation

vllm serve \
  --model-checkpoint /mnt/datakit/checkpoints/my_lm \
  --port 8000

```

## Result Publication with Marin CLI

**Marin CLI** handles the final stage of the pipeline, generating HTML and Markdown reports from training metrics and evaluation results. It archives artifacts and publishes analysis sites for downstream consumption and research reproducibility.

Documentation in [`docs/tutorials/publish-analysis-site.md`](https://github.com/marin-community/marin/blob/main/docs/tutorials/publish-analysis-site.md) describes how this component aggregates logs from Iris, checkpoints from Datakit, and metrics from Levanter to produce static sites documenting the complete experimental provenance.

```bash

# Generate and publish analysis report

marin publish-report \
  --source ./reports/run_2024_09_10 \
  --target https://my-analysis-site.org

```

## Linear Pipeline Execution Flow

The Marin end-to-end pipeline executes as a sequential workflow:

1. **Ingest and Store** – Datakit collects raw shards and assigns versioned identifiers
2. **Process and Normalize** – Zephyr parses raw data and writes cleaned tensors to Datakit containers  
3. **Train** – Levanter builds a `Pipeline` object that consumes Datakit data, applies Fray's distributed strategies, and executes on Iris-managed clusters
4. **Evaluate** – VLLM serves the trained model for benchmark testing and downstream inference
5. **Publish** – Marin CLI produces reproducible reports and publishes static analysis sites

## Summary

- **Datakit** provides immutable, versioned data storage and checkpointing for training datasets
- **Zephyr** handles preprocessing and normalization of raw data sources before training
- **Levanter** implements the core `Pipeline` class in [`lib/levanter/src/levanter/pipeline.py`](https://github.com/marin-community/marin/blob/main/lib/levanter/src/levanter/pipeline.py) to orchestrate JAX training loops
- **Fray** enables distributed execution across multi-node GPU clusters via mesh-aware sharding
- **Iris** manages cloud resource allocation and job lifecycle through the CLI submission interface
- **VLLM integration** serves trained models for high-throughput inference and evaluation
- **Marin CLI** automates report generation and publication of experimental results

## Frequently Asked Questions

### What file defines the core Pipeline class in Marin?

The `Pipeline` class is defined in [`lib/levanter/src/levanter/pipeline.py`](https://github.com/marin-community/marin/blob/main/lib/levanter/src/levanter/pipeline.py). This module implements the primary abstraction that connects Datakit datasets, Levanter models, and distributed training configurations into a single executable workflow.

### How does Marin handle distributed training across multiple GPUs?

Marin leverages **Fray** for distributed execution primitives combined with **Iris** for cluster orchestration. Fray provides mesh-aware sharding and parameter server logic, while Iris manages the underlying cloud infrastructure and job scheduling across nodes.

### What component manages dataset versioning and storage?

**Datakit** exclusively handles data ingestion, versioning, and checkpointing. It ensures that training runs consume immutable dataset shards, with references stored via the `from_datakit()` method pattern demonstrated in the test suite at [`tests/datakit/test_harrier_pipeline.py`](https://github.com/marin-community/marin/blob/main/tests/datakit/test_harrier_pipeline.py).

### How are evaluation results published after training completes?

The **Marin CLI** provides the `publish-report` command, which aggregates training logs, evaluation metrics, and model checkpoints into static HTML sites. This process is documented in [`docs/tutorials/publish-analysis-site.md`](https://github.com/marin-community/marin/blob/main/docs/tutorials/publish-analysis-site.md) and typically executes as the final step in the Iris job workflow.