How Marin Curates Raw Source Data
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, 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_docfunction selectstext_by_page_genfirst, falling back totext_by_page_srconly 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_SEPARATORto form a singletextfield. - Write parquet shards: The transformed data is written to the pipeline output directory using Zephyr’s
write_parquetutility.
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.
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 thexxh3_128hash 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. 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 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:
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:
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:
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_genovertext_by_page_src) and empty-row removal. - Deterministic identifiers are generated using
xxh3_128hashes 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 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. 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.
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 →