How Marin Implements Disk Caching and Artifact Fingerprinting

Marin separates disk caching and artifact fingerprinting into two distinct systems: a decorator that persists function results to cloud storage with status-file tracking, and a deterministic hashing pipeline that produces canonical JSON fingerprints for every lazy artifact.

The marin-community/marin repository treats reproducibility as a first-class concern in ML pipelines. By isolating the mechanism that avoids redundant computation (disk caching) from the mechanism that detects configuration drift (artifact fingerprinting), Marin ensures that distributed workloads remain both efficient and deterministic.

Disk Caching Architecture

The disk_cache decorator in marin.execution.disk_cache provides durable memoization for expensive pure functions. Unlike simple in-memory caches, this implementation writes serialized results to a cloud bucket and uses a status file (StatusFile) to guarantee atomicity.

Cache Key Generation

When you apply @disk_cache without an explicit output_path, Marin generates a deterministic cache key from the function’s arguments. In lib/marin/src/marin/execution/disk_cache.py at lines 88‑95, the fingerprint_args helper concatenates the module name, qualified function name, positional arguments, and sorted keyword arguments, then feeds them through cloudpickle and hashlib.sha256. The resulting hash is sliced to 16 hex characters and mapped to a temporary path via rigging.filesystem.cluster_config.marin_temp_bucket.

Cache Lookup and Execution

Before executing the wrapped function, the decorator checks for a StatusFile adjacent to the cached payload (see lib/marin/src/marin/execution/step_status.py).

  1. Cache hit: If the status reads STATUS_SUCCESS, the decorator invokes load_result and returns the deserialized object (log entry at line 110 of disk_cache.py).
  2. Cache miss: If the status file is missing or indicates failure, the wrapped function executes.
  3. Distributed safety: If a StepAlreadyDone exception surfaces from a distributed-lock wrapper, the decorator treats it as a cache hit from another worker.

Persistence Protocol

After a successful execution, the decorator persists the result using either custom serializers or the default cloudpickle dump to <output_path>/data.pkl. It then writes a new status file with STATUS_SUCCESS (lines 129‑130) and logs the cache write. The decorator intentionally omits locking logic; for multi-worker environments, you must compose it with distributed_lock as noted in the docstring at lines 53‑55.

Practical Example: Basic Function Caching

from marin.execution.disk_cache import disk_cache

@disk_cache
def expensive_preprocess(data_path: str) -> dict:
    # Heavy I/O and transformation logic

    return {"processed": data_path}

# First call executes and persists

result = expensive_preprocess("/my/data")

# Second call reads from GCS without recomputation

result_again = expensive_preprocess("/my/data")

Practical Example: Custom Serializers and Explicit Paths

import json
from pathlib import Path
from marin.execution.disk_cache import disk_cache

def save_json(data, path):
    Path(path, "data.json").write_text(json.dumps(data))

def load_json(path):
    return json.loads(Path(path, "data.json").read_text())

@disk_cache(
    output_path="gs://my-bucket/preprocess-cache",
    save_fn=save_json,
    load_fn=load_json,
)
def preprocess(data_path: str) -> dict:
    return {"source": data_path}

Artifact Fingerprinting System

While disk caching prevents redundant work, artifact fingerprinting guarantees that configuration changes invalidate stale dependencies. Every lazy artifact in Marin carries a deterministic identity composed of a human-readable name@version and a cryptographic recipe fingerprint.

Canonical JSON Encoding

The fingerprinting pipeline begins in lib/marin/src/marin/execution/fingerprint.py. The canonical_json function (implemented via _FingerprintEncoder at lines 90‑108) walks the artifact’s configuration and converts supported types—dataclasses, enums, pathlib.Path, datetime.timedelta, NumPy dtypes—into deterministic JSON representations. Unknown types serialize to a stable string derived from their fully-qualified class name, or raise an error in strict mode.

Hashing and Storage

Once canonicalized, the JSON string passes through fingerprint_hash (lines 164‑169), which computes an MD5 digest and returns the first 8 hex characters as the recipe fingerprint. This short hash is stored in the fingerprint field of an ArtifactStep record, while the full canonical JSON payload occupies fingerprint_payload (see lib/marin/src/marin/execution/step_spec.py, lines 52‑56).

Runtime Verification

During materialization, the step runner (in lib/marin/src/marin/execution/step_runner.py) computes the fingerprint of the current configuration and compares it against the expected_fingerprint stored on the artifact. A mismatch triggers a FingerprintMismatchError (lines 316‑329), halting the pipeline before contaminated data propagates downstream. Many transform steps also accept a content_fingerprint argument that must match the source manifest’s fingerprint, as demonstrated in lib/marin/src/marin/transform/structured_text/web_data_commons.py.

Practical Example: Computing Artifact Fingerprints

from marin.execution.fingerprint import canonical_json, fingerprint_hash
from marin.execution.step_spec import ArtifactStep

config = {
    "tokenizer": "gpt2",
    "max_seq_len": 512,
    "special_tokens": {"pad": "<pad>"},
}

# Generate deterministic payload and hash

payload = canonical_json(config)
fp = fingerprint_hash(payload)

artifact = ArtifactStep(
    name="tokenizer",
    version="v1",
    fingerprint=fp,
    fingerprint_payload=payload,
)

Integration of Caching and Fingerprinting

These two systems operate at different granularities but compose cleanly. When a step produces a lazy artifact, the artifact’s fingerprint is computed before execution begins (using StepContext.for_fingerprint). The disk_cache key remains independent of the artifact fingerprint, yet higher-level decorators often combine both:

  • Disk caching skips the computation if the function inputs have not changed.
  • Fingerprinting ensures that if the underlying configuration drifts, the artifact’s identity changes, forcing downstream steps to recognize a new dependency and recompute their own caches.

This separation allows Marin to cache expensive pure functions while still detecting configuration drift through fingerprint mismatches—a critical property for reproducible ML pipelines.

Summary

  • Disk caching uses the @disk_cache decorator in disk_cache.py to persist function results to cloud storage, tracking state via StatusFile objects and 16-character SHA256 hashes of arguments.
  • Artifact fingerprinting relies on canonical_json and fingerprint_hash in fingerprint.py to produce 8-character MD5 identifiers of configuration objects, stored in ArtifactStep records.
  • Cache keys are derived from function arguments via cloudpickle, while fingerprint payloads are derived from configuration objects via deterministic JSON encoding.
  • Verification happens at materialization time in step_runner.py, where mismatches raise FingerprintMismatchError to prevent stale artifacts from propagating.
  • Distributed safety requires composing disk_cache with external locking mechanisms; the decorator itself handles atomic writes via status files but not concurrent access control.

Frequently Asked Questions

How does Marin generate cache keys for the disk cache decorator?

Marin generates cache keys by serializing the function’s module name, qualified name, positional arguments, and sorted keyword arguments using cloudpickle, then computing a SHA256 hash of the resulting bytes. In lib/marin/src/marin/execution/disk_cache.py at lines 88‑95, the fingerprint_args function slices this hash to 16 hexadecimal characters and appends it to a temporary bucket path configured via rigging.filesystem.cluster_config.marin_temp_bucket.

What happens when an artifact’s fingerprint does not match the expected value?

When the step runner materializes an artifact, it recomputes the fingerprint from the current configuration and compares it to the expected_fingerprint stored in the ArtifactStep record. If the values differ, the runner raises a FingerprintMismatchError (see lib/marin/src/marin/execution/step_runner.py, lines 316‑329). This exception prevents the pipeline from proceeding with stale or misconfigured data, ensuring strict reproducibility.

Can I use custom serializers with Marin’s disk cache?

Yes. The disk_cache decorator accepts optional save_fn and load_fn callables. If provided, these functions replace the default cloudpickle serialization. Your custom functions must accept the result object and the output path, allowing you to write JSON, Parquet, or any other format to the underlying storage bucket while still benefiting from the status-file tracking mechanism.

How does canonical JSON encoding handle complex Python types?

The _FingerprintEncoder class in lib/marin/src/marin/execution/fingerprint.py (lines 90‑108) registers encoders for stable types including dataclasses, enums, pathlib.Path, datetime.timedelta, and NumPy dtypes. Each encoder produces a deterministic JSON representation. Types without explicit encoders fall back to a stable string based on their fully-qualified class name, or raise a serialization error if strict mode is enabled, ensuring that the resulting fingerprint always reflects the exact configuration.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →