# How Marin Curates Raw Source Data

> Learn how Marin curates raw source data through a deterministic download transform normalize pipeline. Ensure reproducibility and traceability with Datakit and StepSpec.

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

---

**Marin treats raw source data as a first-class asset, processing it through a deterministic download-transform-normalize pipeline built on the Datakit framework and orchestrated by StepSpec objects to ensure reproducibility and full traceability.**

Data curation in the marin-community/marin repository follows a strict engineering discipline. Instead of ad-hoc scripts, the `marin.datakit` package defines a declarative pipeline where each stage—downloading, transforming, and normalizing—is cached, versioned, and executed as a node in a directed acyclic graph (DAG). This architecture guarantees that every document in the final training corpus can be traced back to its exact source revision and transformation logic.

## The Three-Stage Curation Pipeline

Marin's curation workflow is implemented as a sequence of immutable **StepSpec** objects that move data from raw Hugging Face repositories to canonical Dolma-shaped documents ready for model training.

### Stage 1: Downloading Raw Sources with Deterministic Revisions

The pipeline begins with `download_hf_step` (or the higher-level `hf_normalize_steps` utility), which creates a **StepSpec** that fetches datasets from Hugging Face or other public repositories. This step:

- Stores raw data under deterministic paths (e.g., `raw/institutional-books-<revision>`)
- Records the exact Git revision hash for provenance
- Registers itself as a dependency for downstream transforms

In [`lib/marin/src/marin/datakit/download/institutional_books.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/datakit/download/institutional_books.py), the Institutional Books dataset is downloaded using this pattern, ensuring that the revision `abcdef1` (or similar) is locked and reproducible.

### Stage 2: Transforming Heterogeneous Records

Raw sources often contain per-page or per-segment fields that are incompatible with downstream tokenizers. The transform step reads raw parquet files and emits uniform records through a series of strict normalization rules implemented in the source-specific transform functions:

- **Prefer OCR-corrected text:** The `row_to_doc` function selects `text_by_page_gen` first, falling back to `text_by_page_src` only when the corrected version is missing or empty.
- **Drop empty rows:** Pages containing no useful text are filtered out and counted via Zephyr counters for auditability.
- **Join pages:** All selected pages are concatenated using a double-newline `PAGE_SEPARATOR` to form a single `text` field.
- **Write parquet shards:** The transformed data is written to the pipeline output directory using Zephyr’s `write_parquet` utility.

The Common-Pile family follows the same pattern but restricts ingestion to the curated subsets listed in `_COMMON_PILE_ENTRIES`, which acts as an authoritative registry of high-quality, pre-filtered data. This registry is defined in [`lib/marin/src/marin/datakit/download/common_pile.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/datakit/download/common_pile.py).

### Stage 3: Normalizing to the Dolma Schema

After transformation, the `normalize_step` reads the intermediate parquet shards and produces the final output:

- Computes a deterministic document identifier (`id`) using the `xxh3_128` hash of the textual content
- Appends required metadata (e.g., source name)
- Writes final Dolma-formatted parquet files consumed by training pipelines

The complete chain is exposed as a single function, such as `institutional_books_normalize_steps()`, which returns a tuple of **StepSpec** objects linking the download, transform, and normalize stages.

## Orchestration and Reproducibility

The execution engine treats each stage as a **StepSpec** object defined in [`lib/marin/src/marin/execution/step_spec.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/execution/step_spec.py). Each specification encodes:

- A **name** (e.g., `processed/institutional-books`)
- **Dependencies** (`deps`) that enforce topological ordering
- The executable **function** (`fn`)
- A **hash_attrs** map capturing logical versioning (e.g., page separator strings or version tags)

The **StepRunner** class in [`lib/marin/src/marin/execution/step_runner.py`](https://github.com/marin-community/marin/blob/main/lib/marin/src/marin/execution/step_runner.py) walks this DAG, caches intermediate results to disk, and guarantees that each raw source is processed exactly once per revision. If a step’s inputs or code change, the hash invalidates the cache and triggers a re-run.

## Practical Implementation Examples

### Running the Institutional Books Pipeline

To execute the full curation workflow for the Institutional Books dataset:

```python
from marin.datakit.download.institutional_books import institutional_books_normalize_steps
from marin.execution.step_runner import StepRunner

# Obtain the two steps: (download + transform, normalize)

download_transform, normalize = institutional_books_normalize_steps()

# Execute the full pipeline

StepRunner().run([download_transform, normalize])

```

### Adding a New Curated Source

New datasets follow a template pattern using `hf_normalize_steps`:

```python
from marin.datakit.download.hf_simple_util import hf_normalize_steps
from marin.execution.step_spec import StepSpec

# 1. Define the dataset registry entry

_MY_DATASET = (
    "mycurated/xyz",
    "my-hf-dataset/xyz_filtered",
    "abcdef1",
    "raw/mycurated/xyz_filtered-abcdef1",
)

# 2. Generate the download+normalize steps

def my_dataset_steps() -> tuple[StepSpec, ...]:
    return hf_normalize_steps(
        marin_name="mycurated/xyz",
        hf_dataset_id="my-hf-dataset/xyz_filtered",
        revision="abcdef1",
        staged_path="raw/mycurated/xyz_filtered-abcdef1",
        file_extensions=(".json.gz",),   # adapt to source format

    )

```

### Auditing Data Quality with Counters

After execution, inspect Zephyr counters to verify the curation filter efficacy:

```python
from zephyr import counters

# Retrieve row statistics

kept = counters.pipeline.get("institutional_books/kept")
dropped = counters.pipeline.get("institutional_books/dropped")
print(f"Kept {kept} rows, dropped {dropped} rows")

```

## Summary

- **Marin curates raw source data** through a immutable three-stage pipeline: download, transform, and normalize.
- **StepSpec objects** declare dependencies and logical versions, enabling the StepRunner to execute a reproducible DAG.
- **Source-specific transforms** enforce quality rules like OCR preference (`text_by_page_gen` over `text_by_page_src`) and empty-row removal.
- **Deterministic identifiers** are generated using `xxh3_128` hashes to support deduplication and lineage tracking.
- **Zephyr counters** provide observable metrics for data quality audits at every stage.

## Frequently Asked Questions

### What is the role of StepSpec in Marin's data curation?

**StepSpec** is the fundamental unit of work in Marin’s execution engine. Each StepSpec defines a named step with explicit dependencies (`deps`), a callable function (`fn`), and a `hash_attrs` dictionary that captures versioning information. The StepRunner uses these specifications to build a DAG, cache results, and ensure that downstream steps only execute after their upstream dependencies complete successfully.

### How does Marin handle OCR quality in scanned documents?

Marin prioritizes machine-corrected text over raw OCR output. During the transform stage, the `row_to_doc` function in [`institutional_books.py`](https://github.com/marin-community/marin/blob/main/institutional_books.py) attempts to read from the `text_by_page_gen` field first. If that field is missing or empty, it falls back to `text_by_page_src`. This ensures that higher-quality corrected text is used whenever available, while maintaining robustness against missing data.

### Why does Marin use xxh3_128 for document identifiers?

The `normalize_step` generates document IDs by computing the `xxh3_128` hash of the document's textual content. This produces deterministic, content-addressable identifiers that are consistent across re-runs and environments. Content-based hashing enables automatic deduplication and guarantees that identical text segments receive identical IDs regardless of which pipeline instance processed them.

### How can I add a custom dataset to the Marin curation pipeline?

Create a new Python module under `marin.datakit.download` that defines a registry entry (similar to `_COMMON_PILE_ENTRIES`) and a step-generation function using `hf_normalize_steps` from [`hf_simple_util.py`](https://github.com/marin-community/marin/blob/main/hf_simple_util.py). This function should specify the Hugging Face dataset ID, target revision, staging path, and file extensions. Return a tuple of StepSpec objects from your function, then invoke them via `StepRunner.run()` to integrate your source into the curated corpus.