Skip to main content

Data Versioning in MLOps: Making Every Experiment Reproducible

Calculating read time…

Imagine you bake the most delicious chocolate cake ever. Your family goes crazy for it. Everyone wants the recipe!

But here's the problem: you didn't write down exactly how much flour you used. You don't remember if the chocolate was dark or milk. You changed the oven temperature halfway through but forgot by how much.

You try to bake it again next week. It tastes completely different. You can never recreate that perfect cake. 




💡 That cake = your ML model. The ingredients = your training data.

In Machine Learning, if you train a great model but didn't record exactly which data you used, with exactly which version, you can never reproduce that result — not for your team, not for auditors, not even for yourself next month.

Data Versioning is the discipline of treating your datasets exactly like your code — tracked, versioned, tagged, and reproducible. 🎯

📚 What We'll Cover

  • 🔹 What is Data Versioning? (The Big Idea)
  • 🔹 Why Data Changes — And Why That's a Problem
  • 🔹 The 5 Painful Problems Without Data Versioning
  • 🔹 The ML Reproducibility Crisis — Why It Matters
  • 🔹 Data Versioning vs Code Versioning — What's Different
  • 🔹 DVC — Data Version Control (The Git for Data)
  • 🔹 DVC Deep Dive: Setup, Tracking, Pipelines, Experiments
  • 🔹 LakeFS — Git-Like Branching for Your Entire Data Lake
  • 🔹 LakeFS Deep Dive: Branches, Commits, Merges, CI/CD for Data
  • 🔹 Data Versioning + Feature Store — The Complete MLOps Pipeline
  • 🔹 Other Tools in the Data Versioning Ecosystem
  • 🔹 Best Practices & Anti-Patterns
  • 🔹 High-Level Summary & Next Steps

📌 Section 1: What is Data Versioning? (The Big Idea)

Data Versioning means recording every state that a dataset has ever been in — with a unique identifier, a timestamp, and a way to go back to it.

It's the same idea as version control for code (Git), but applied to datasets, which are large, binary, and constantly changing.

📷 The "Camera Roll" Analogy

Think of your phone's camera roll. Every photo is timestamped. You can scroll back to any day and see exactly what your life looked like on that date.

Data Versioning gives your datasets a "camera roll." Every time data changes — new rows added, errors cleaned, columns renamed — a new snapshot is recorded. You can always go back to any past state. 📸

Without Data Versioning (nightmare):
├── train_data.csv (which version? when?)
├── train_data_v2.csv (what changed?)
├── train_data_FINAL.csv (is this actually final?)
├── train_data_FINAL2.csv (😱)
└── train_data_use_this_one.csv (please no)

With Data Versioning (organized!):
├── dataset → v1.0 (2025-01-10, 50,000 rows, baseline)
│ → v1.1 (2025-01-22, +5,000 rows, fixed null values)
│ → v2.0 (2025-02-14, full refresh, new features added)
└── Any version instantly restorable with one command ✅

🔄 Section 2: Why Data Changes — And Why That's a Problem

Unlike code, which you deliberately change, data changes for many reasons — often outside your control.

🌊 Reasons Data Changes in ML Projects

  • New data collected: Your pipeline adds 10,000 new user events every day. The training set is different today than it was last week.
  • Data cleaning: You discover 3,000 duplicate rows and remove them. Now the dataset is smaller — but which model used which version?
  • Label corrections: Human reviewers relabeled 500 images as "cat" that were previously "dog." Models trained before vs after see fundamentally different ground truth.
  • Feature engineering changes: You add a new column: user_risk_score. Now old models trained without this column are incompatible.
  • Schema changes: A column is renamed from amt to transaction_amount. Code that worked yesterday breaks silently today.
  • Data source updates: Your vendor ships a corrected dataset replacing last month's. You don't know if your current model was trained on the old or new data.
❌ Without data versioning, every data change creates a mystery:

"Was Model A trained on the cleaned or uncleaned dataset?"
"Did the relabeling happen before or after our best experiment?"
"Which dataset was used for the audit report we submitted?"

These questions become impossible to answer reliably. Compliance fails. Experiments can't be reproduced. Good models can't be recovered because no one remembers what data produced them.

💥 Section 3: The 5 Painful Problems Without Data Versioning

Problem 1: You Can't Reproduce a Good Model 🎲

Three months ago your model hit 94% accuracy. You want to rebuild it, but the training data has changed since then. You run the same code on today's data and get 87%. You don't know if the model got worse or the data changed. Without data versioning, you'll never know which it was.

Problem 2: Debugging Becomes Guesswork 🔎

Your model's performance suddenly drops from 92% to 81% in production. Was it the model code? The feature pipeline? Or did someone accidentally overwrite a data file with an incomplete batch? Without data versioning, you're debugging in the dark.

Problem 3: Compliance & Audits Are Impossible 🏛️

Financial, healthcare, and government AI systems are required by law to demonstrate exactly what data was used to train a model, when it was collected, and whether it was free from bias. Without data versioning, you cannot answer these questions. Regulatory approval is denied or revoked.

Problem 4: Experiments Become Unfair 🧪

Data Scientist A ran Experiment 1 on Monday's dataset. Data Scientist B ran Experiment 2 on Friday's dataset. They compare results and B's model appears better. But the Friday dataset had cleaner labels! The experiment comparison is meaningless — they weren't using the same data.

Problem 5: Collaboration Breaks Down 🤝

Ten people on a team each keep their own copy of "the training data." Everyone has a slightly different version. There's no single source of truth. Merging results is impossible. Trust collapses.

✅ Data Versioning solves ALL five of these problems:
  • ✅ Tag your best model with the exact dataset version it used
  • ✅ Instantly trace any performance change to a specific data change
  • ✅ Pass any compliance audit with an immutable data lineage trail
  • ✅ Ensure all experiments use the same dataset version — fair comparison
  • ✅ One shared, versioned data repository for the whole team

🔬 Section 4: The ML Reproducibility Crisis — Why It Matters

Studies suggest that over 70% of ML experiments are not reproducible — even by the original authors. This isn't just an academic problem; it has real consequences.

📊 What "Reproducibility" Actually Means

An ML experiment is reproducible if, given the same:

  • Code (model architecture + training script)
  • Data (exact same dataset, same version)
  • Environment (Python version, library versions, hardware)
  • Configuration (hyperparameters, random seeds)

…you get the same (or statistically identical) results.

Data is the hardest piece to pin down. Code is tracked in Git. Environments are tracked in Docker. But data? It just lives on a shared drive or S3 bucket, silently mutating with every pipeline run. 😬

The 4 Pillars of ML Reproducibility:

┌─────────────────────────────────────────────────────────┐
│ CODE → Git (solved ✅) │
│ ENVIRONMENT → Docker / conda (solved ✅) │
│ CONFIG → MLflow / config files (solved ✅) │
│ DATA → DVC / LakeFS (this guide solves it! 🎯)│
└─────────────────────────────────────────────────────────┘
💡 Industry Trend:
Regulatory frameworks like the EU AI Act and FDA AI/ML guidance for medical devices now explicitly require data provenance and version traceability for AI systems in regulated industries.

Data versioning is no longer a nice-to-have. For any AI product that faces compliance review, it's a hard legal requirement. ⚖️

💡 Section 5: Data Versioning vs Code Versioning — What's Different?

You already know how Git works for code. Why can't you just use Git for data too?

🤔 Why Git Alone Doesn't Work for Data

Aspect Code (Git) Data (needs DVC/LakeFS)
Typical size KB – MB GB – TB – PB
Format Text files Binary (Parquet, CSV, images, audio)
Change frequency Deliberate, human-initiated Automatic, pipeline-driven, continuous
Diff comparison Line-by-line text diff Statistical summary diff (row counts, schema, distributions)
Storage model Full copy per version Content-addressed + deduplication (can't store 100 copies of a 1TB file)
Branching Lightweight (just pointers) Heavyweight (data must be copy-on-write at storage level)

Git stores the full content of every change. If you add 1 billion rows to a 500GB CSV, Git would store the entire 500GB again. That's untenable.

Data versioning tools solve this with content-addressed storage and pointer-based versioning. Instead of storing the data itself in Git, you store a tiny fingerprint (hash) of the data, and the actual data lives in a separate storage system (S3, GCS, Azure Blob).


🔧 Section 6: DVC — Data Version Control (The Git for Data)

DVC (Data Version Control) is the most widely adopted open-source tool for data versioning in ML. It was built from the ground up to solve the reproducibility problem by extending Git to handle large data files and ML pipelines.

🎮 The Video Game Save File Analogy

Remember saving your progress in a video game? You could have multiple save slots: Save 1 (before the boss), Save 2 (after getting the sword), Save 3 (with the dragon defeated). You could reload any save and play from exactly that point.

DVC creates "save slots" for your datasets. Every dataset version is a named, restorable checkpoint. Your code (Git) and your data (DVC) are always kept in sync. Any commit in Git knows exactly which data version it used. 🎮

🏗️ How DVC Works Under the Hood

YOUR PROJECT:

├── .git/ ← Git tracks code changes
├── .dvc/ ← DVC config and cache
│ └── cache/ ← Local content-addressed store
├── data/
│ └── train.csv.dvc ← TINY pointer file (tracked by Git!) ✅
│ Contents: {md5: "a3f7...", size: 524288000, path: "train.csv"}
├── models/
│ └── model.pkl.dvc ← Same idea for model artifacts
└── dvc.yaml ← Pipeline definition

REMOTE STORAGE (S3/GCS/Azure):
└── a3f7.../ ← Actual 500MB CSV stored here by hash
(DVC fetches this when you need the actual data)

The .dvc pointer file is tiny (just a few lines) and lives in Git. The actual data lives in remote storage (S3, GCS, etc.). When you git checkout a past commit, DVC knows which data version that commit points to and can fetch it with dvc pull.

✅ DVC's Core Superpower:
Git and DVC are always in sync. One git log entry corresponds to one exact data version. Checkout any git commit → pull the matching data → run the exact same experiment. True reproducibility. Always. 🔁

💻 Section 7: DVC Deep Dive — Setup, Tracking, Pipelines, Experiments

7.1 Installation & Project Setup

📌 Code Purpose — DVC Installation & Project Initialization

What this code does: Installs DVC with S3 support, initializes a DVC repository inside an existing Git repo, and configures an S3 bucket as the remote storage backend. This is the one-time setup you run at the beginning of any new ML project.

Why it matters: Without this setup, your data files are untracked and invisible to version control. After this setup, every data file you add will have a tiny Git-tracked pointer (.dvc file) that permanently links code commits to specific data versions.

Key concept: dvc init creates the .dvc/ folder inside your Git repo — the control center for all data tracking. dvc remote add tells DVC where to actually store the large files.
# ── STEP 1: Install DVC with your preferred storage backend ──────────
# S3 support
pip install "dvc[s3]"

# Or GCS support
pip install "dvc[gs]"

# Or Azure support
pip install "dvc[azure]"

# Or install all backends
pip install "dvc[all]"


# ── STEP 2: Initialize DVC inside your Git repository ────────────────
# This must be run inside a git repository (git init first if needed)

git init my-ml-project
cd my-ml-project
dvc init

# DVC creates:
# .dvc/                    ← DVC control directory
# .dvc/.gitignore          ← Tells Git to ignore actual data cache
# .dvc/config              ← DVC configuration
# .dvcignore               ← Like .gitignore but for DVC

git add .dvc .dvcignore
git commit -m "Initialize DVC"


# ── STEP 3: Configure Remote Storage ─────────────────────────────────
# Tell DVC where to store your actual data files

# Option A: Amazon S3
dvc remote add -d myremote s3://my-ml-data-bucket/dvc-store

# Option B: Google Cloud Storage
dvc remote add -d myremote gs://my-gcs-bucket/dvc-store

# Option C: Azure Blob Storage
dvc remote add -d myremote azure://my-container/dvc-store

# Option D: SSH server (e.g., your on-prem server)
dvc remote add -d myremote ssh://user@server.com/path/to/store

# Option E: Local path (for development/testing)
dvc remote add -d myremote /mnt/shared-storage/dvc

# Verify the remote was configured
dvc remote list
# Output: myremote    s3://my-ml-data-bucket/dvc-store

git add .dvc/config
git commit -m "Configure DVC remote storage: S3"

echo "✅ DVC setup complete! Ready to version your data."

7.2 Tracking Your First Dataset

📌 Code Purpose — Adding Data to DVC Tracking

What this code does: Walks through the complete workflow of adding a dataset to DVC tracking, pushing it to remote storage, and showing how Git and DVC work together. It also simulates making a data change and committing a new version, showing how you can switch between v1 and v2 with a single command.

Why it matters: This is the fundamental workflow every ML engineer needs to know. After running these commands, your dataset has an immutable version history. Any team member can checkout any past commit and get the exact data used for that experiment — even 2 years later.

Key concept: dvc add creates the .dvc pointer file. dvc push uploads the actual data to S3. dvc pull downloads the data for a given commit. The actual data is never stored in Git.
# ── Assume you have a dataset ready ──────────────────────────────────
ls -lh data/
# train.csv   (450 MB)
# test.csv    (50 MB)


# ── STEP 1: Add data files to DVC tracking ───────────────────────────
dvc add data/train.csv
dvc add data/test.csv

# What happened?
# 1. DVC computed MD5 hash of each file
# 2. DVC copied files to local .dvc/cache (content-addressed)
# 3. DVC created pointer files:

ls data/
# train.csv        ← your actual data (now in .dvc/cache too)
# train.csv.dvc    ← tiny pointer file created by DVC
# test.csv
# test.csv.dvc

# Inspect the pointer file (it's just a few lines!)
cat data/train.csv.dvc
# outs:
# - md5: a3f7b2c19d4e8f6a...    ← unique fingerprint of this exact file
#   size: 471859200              ← 450 MB
#   path: train.csv


# ── STEP 2: Tell Git to track the pointer, not the data ──────────────
# .dvc files go into Git — actual data does NOT
cat data/.gitignore
# /train.csv     ← DVC auto-added this! Git ignores the actual data.
# /test.csv

git add data/train.csv.dvc data/test.csv.dvc data/.gitignore
git commit -m "Add training and test datasets v1.0

Dataset: fraud detection baseline
Rows: 450,000 train / 50,000 test
Source: internal transaction log 2024-Q4
Features: 28 columns"


# ── STEP 3: Push data to remote storage ──────────────────────────────
# This uploads the actual 450MB file to S3
dvc push

# Output:
# 2 files pushed
# Uploading: data/train.csv → s3://my-ml-data-bucket/dvc-store/a3/f7b2c19d...
# Uploading: data/test.csv  → s3://my-ml-data-bucket/dvc-store/b8/d4e7f2c3...


# ── STEP 4: Simulate a data update (new data arrives) ────────────────
# New week — pipeline added 10,000 more transactions

python -c "
import pandas as pd, numpy as np
df = pd.read_csv('data/train.csv')
new_rows = pd.DataFrame({
    'user_id':    np.random.randint(1000, 9999, 10000),
    'amount':     np.random.exponential(50, 10000).round(2),
    'label':      np.random.choice([0,1], 10000, p=[0.97, 0.03])
    # ... other columns
})
df = pd.concat([df, new_rows], ignore_index=True)
df.to_csv('data/train.csv', index=False)
print(f'Updated dataset: {len(df):,} rows')
"

# Track the new version
dvc add data/train.csv

# The .dvc file now has a NEW hash — the old data is preserved in cache!
cat data/train.csv.dvc
# outs:
# - md5: c9f2a4b8e3d1f7b5...   ← different hash = new version
#   size: 482344960              ← slightly larger


# ── STEP 5: Commit the new version to Git ────────────────────────────
git add data/train.csv.dvc
git commit -m "Update training dataset v1.1

Added 10,000 new transactions from 2025-01-20 to 2025-01-26
Total rows: 460,000 (+10,000)
Note: No schema changes, all existing features preserved"

dvc push   # Upload the new version to S3


# ── STEP 6: Switch between versions like magic! ──────────────────────

# Go back to v1.0
git checkout HEAD~1 -- data/train.csv.dvc
dvc pull data/train.csv
echo "Now using: $(wc -l < data/train.csv) rows (v1.0)"

# Come back to v1.1
git checkout HEAD -- data/train.csv.dvc
dvc pull data/train.csv
echo "Now using: $(wc -l < data/train.csv) rows (v1.1)"

7.3 DVC Pipelines — Reproducible End-to-End Workflows

DVC doesn't just version data files. It lets you define your entire ML pipeline as a Directed Acyclic Graph (DAG) of stages — each stage with defined inputs, outputs, and commands. When inputs change, only affected stages re-run. Smart! 🧠

📌 Code Purpose — DVC Pipeline Definition (dvc.yaml)

What this code does: Defines a complete 4-stage ML pipeline: data preparation → feature engineering → model training → model evaluation. Each stage declares its inputs (deps) and outputs (outs), so DVC can automatically detect what needs to re-run when something changes.

Why it matters: If only the evaluation metric threshold changes, DVC skips data prep and training — it only re-runs the evaluation stage. If the raw data changes, DVC re-runs all downstream stages automatically. This is the equivalent of a Makefile, but for ML with built-in data versioning.

Key concept: dvc repro runs only the stages whose inputs have changed since the last successful run. It's caching + dependency tracking in one command.
# dvc.yaml — The complete ML pipeline definition
# Each stage has: cmd (command to run), deps (inputs), outs (outputs)
# DVC tracks the hashes of all deps and outs.

stages:

  # ── Stage 1: Prepare raw data ──────────────────────────────────────
  prepare:
    cmd: python src/prepare.py
    deps:
      - data/raw/transactions.csv       # Raw input data
      - src/prepare.py                  # Script (code change = re-run)
      - configs/prepare_config.yaml     # Config (config change = re-run)
    outs:
      - data/prepared/train.csv
      - data/prepared/val.csv
      - data/prepared/test.csv
    metrics:
      - reports/prepare_stats.json:     # Track row counts, null rates etc.
          cache: false

  # ── Stage 2: Feature Engineering ───────────────────────────────────
  featurize:
    cmd: python src/featurize.py
    deps:
      - data/prepared/train.csv         # Output of previous stage
      - data/prepared/val.csv
      - src/featurize.py
      - configs/feature_config.yaml
    outs:
      - data/features/train_features.pkl
      - data/features/val_features.pkl
      - data/features/feature_names.json

  # ── Stage 3: Model Training ─────────────────────────────────────────
  train:
    cmd: python src/train.py
    deps:
      - data/features/train_features.pkl
      - data/features/val_features.pkl
      - src/train.py
      - configs/train_config.yaml       # Hyperparameters
    outs:
      - models/fraud_model.pkl
      - models/model_metadata.json
    metrics:
      - reports/training_metrics.json:
          cache: false

  # ── Stage 4: Model Evaluation ───────────────────────────────────────
  evaluate:
    cmd: python src/evaluate.py
    deps:
      - data/features/val_features.pkl
      - models/fraud_model.pkl
      - src/evaluate.py
    metrics:
      - reports/eval_metrics.json:
          cache: false
    plots:
      - reports/confusion_matrix.json
      - reports/roc_curve.json
# ── Run the full pipeline ─────────────────────────────────────────────
dvc repro

# DVC output (smart caching in action):
# Running stage 'prepare':    > python src/prepare.py
# Running stage 'featurize':  > python src/featurize.py
# Running stage 'train':      > python src/train.py
# Running stage 'evaluate':   > python src/evaluate.py
# Use `dvc push` to send to remote.


# ── Now change ONLY hyperparameters (train_config.yaml) ──────────────
# Edit configs/train_config.yaml: learning_rate: 0.01 → 0.001

dvc repro

# DVC smart caching output:
# Stage 'prepare' didn't change, skipping          ← cached! 🚀
# Stage 'featurize' didn't change, skipping        ← cached! 🚀
# Running stage 'train': > python src/train.py     ← only this re-runs
# Running stage 'evaluate': > python src/evaluate.py


# ── Show the pipeline as a visual DAG ────────────────────────────────
dvc dag

# Output:
#           +----------+
#           | prepare  |
#           +----------+
#                *
#                *
#           +-----------+
#           | featurize |
#           +-----------+
#                *
#                *
#            +-------+
#            | train |
#            +-------+
#                *
#                *
#           +----------+
#           | evaluate |
#           +----------+


# ── Compare metrics between experiments ──────────────────────────────
# Check current experiment metrics
dvc metrics show

# Output:
# Path                          roc_auc    precision  recall
# reports/eval_metrics.json     0.9412     0.8823     0.7956


# Compare two git commits' metrics
dvc metrics diff HEAD~1 HEAD

# Output:
# Path                    Metric     HEAD~1    HEAD     Change
# eval_metrics.json       roc_auc    0.9180    0.9412   +0.0232  ✅
# eval_metrics.json       precision  0.8701    0.8823   +0.0122  ✅
# eval_metrics.json       recall     0.7812    0.7956   +0.0144  ✅

7.4 DVC Experiments — Hyperparameter Search with Data Versioning

📌 Code Purpose — DVC Experiment Tracking

What this code does: Runs multiple ML experiments with different hyperparameters using dvc exp run, each experiment is automatically linked to the exact data version and code commit active at that moment. Then compares all experiments in a single table to find the best one.

Why it matters: Unlike MLflow or Weights & Biases which only track model metrics, DVC experiment tracking also records which exact dataset version each experiment used — the complete reproducibility package in one tool.

Key concept: dvc exp run --set-param overrides a config value for one run without permanently changing your files. dvc exp show shows all experiments in a comparison table.
# ── Run a grid of experiments with different hyperparameters ─────────
# Each experiment automatically records which data version it used!

# Experiment 1: baseline
dvc exp run --name "baseline-lr01"

# Experiment 2: higher learning rate
dvc exp run --name "lr-sweep-001" \
    --set-param train.learning_rate=0.001

# Experiment 3: more trees
dvc exp run --name "lr-sweep-005" \
    --set-param train.learning_rate=0.005

# Experiment 4: different model depth
dvc exp run --name "depth-sweep-6" \
    --set-param train.learning_rate=0.001 \
    --set-param train.max_depth=6

# Run experiments in parallel (queue mode)
dvc exp run --queue \
    --set-param train.learning_rate=0.01,0.001,0.005 \
    --set-param train.n_estimators=100,200,300

dvc queue start --jobs 3   # Run 3 in parallel


# ── Compare all experiments ───────────────────────────────────────────
dvc exp show

# Output (truncated for readability):
# ┏━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┳━━━━━━━━━┳━━━━━━━━━┳━━━━━━━━━┳━━━━━━━━━━┓
# ┃ Experiment       ┃ Data Hash     ┃ roc_auc ┃  lr     ┃  depth  ┃ n_est    ┃
# ┡━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━╇━━━━━━━━━╇━━━━━━━━━╇━━━━━━━━━╇━━━━━━━━━━┩
# │ baseline-lr01    │  a3f7b2c1...  │  0.9180 │  0.01   │  4      │  100     │
# │ lr-sweep-001     │  a3f7b2c1...  │  0.9412 │  0.001  │  4      │  100     │ ← best
# │ lr-sweep-005     │  a3f7b2c1...  │  0.9334 │  0.005  │  4      │  100     │
# │ depth-sweep-6    │  a3f7b2c1...  │  0.9388 │  0.001  │  6      │  100     │
# └──────────────────┴───────────────┴─────────┴─────────┴─────────┴──────────┘
# Note: Same data hash (a3f7b2c1) means ALL experiments used identical data ✅


# ── Promote the best experiment to a permanent branch ────────────────
dvc exp branch lr-sweep-001 "best-model-jan-2025"
git checkout best-model-jan-2025
dvc pull   # Get the exact data and model for this experiment

🌊 Section 8: LakeFS — Git-Like Branching for Your Entire Data Lake

While DVC is excellent for individual ML projects, LakeFS operates at a much larger scale — it brings Git-like operations (branch, commit, merge, revert) to your entire data lake (S3, GCS, Azure Blob).

🌲 The Forest Management Analogy

Imagine a vast forest that dozens of rangers manage simultaneously. Without coordination, one ranger might accidentally clear-cut a section while another is trying to survey it. Changes collide. Data is lost.

LakeFS puts a management layer over the forest. Each ranger works on their own "branch" of the forest — a safe copy where they can make changes without affecting anyone else. When their work is ready, it's reviewed and merged into the main forest. The original forest state is always preserved and restorable. 🌲

🏗️ LakeFS Architecture

TRADITIONAL S3 (no versioning):
s3://my-bucket/data/train.parquet ← one file, constantly overwritten 😬

WITH LAKEFS (full versioning):
lakefs://my-repo/main/data/train.parquet ← production branch
lakefs://my-repo/experiment-1/data/train.parquet ← safe sandbox
lakefs://my-repo/experiment-2/data/train.parquet ← another sandbox

LakeFS speaks S3 API — existing tools (Spark, Pandas, Athena) work unchanged!
Just swap s3:// for lakefs:// ← zero code changes needed ✅

LAKEFS OBJECT MODEL:
├── Repository (like a GitHub repo, but for data)
│ ├── Branch: main (production-ready data)
│ ├── Branch: staging (tested, waiting for approval)
│ └── Branch: dev (work in progress)
│ ├── Commit: c3f8a... "Add Jan 2025 transactions"
│ ├── Commit: b2d9e... "Remove duplicate rows"
│ └── Commit: a1c7b... "Initial data load"
└── Tags: v1.0, v2.0, experiment-A (immutable named checkpoints)

🔑 What Makes LakeFS Different from DVC

Feature DVC LakeFS
Scope Per-project, alongside code Entire data lake / org-wide
Storage layer Works on top of existing storage Sits as a proxy in front of S3/GCS
API compatibility DVC-specific CLI commands Fully S3-API compatible — no code changes needed
Branching model Via Git branches + DVC pointers Native data branches (copy-on-write)
Team model Individual ML engineers Entire data teams (DE + DS + Analytics)
Pipeline focus ML experiment pipelines Data engineering + ETL + ML together
Best use case Individual/small team ML projects Enterprise data lake governance

💻 Section 9: LakeFS Deep Dive — Branches, Commits, Merges & CI/CD

9.1 LakeFS Setup & First Repository

📌 Code Purpose — LakeFS Repository Setup & First Commit

What this code does: Sets up LakeFS using Docker (local dev mode), creates a data repository, configures a branch, uploads data, and commits it — all using the LakeFS Python SDK. Demonstrates the core Git-like workflow for data.

Why it matters: After this setup, every data file in LakeFS is versioned with an immutable commit hash. You can point any ML experiment at a specific commit (like main@c3f8a2) and be guaranteed to always get exactly the same data — forever.

Key concept: Unlike DVC which works alongside your storage, LakeFS sits in front of your S3 bucket as an S3-compatible proxy. Your Spark jobs, Pandas readers, and Athena queries just change s3:// to s3://lakefs-endpoint/repo/branch/ — no other code changes needed.
# ── OPTION 1: Start LakeFS locally with Docker ───────────────────────
# Quick local setup for development and learning

docker run --pull always \
    -p 8000:8000 \
    -e LAKEFS_BLOCKSTORE_TYPE=local \
    -e LAKEFS_AUTH_ENCRYPT_SECRET_KEY=some-secret-key \
    treeverse/lakefs run --local-settings

# LakeFS is now running at: http://localhost:8000
# Default credentials: access_key=AKIAIOSFINEXAMPLE, secret=wJalrXUtnFEMI...


# ── OPTION 2: LakeFS on Kubernetes (production) ──────────────────────
# helm repo add lakefs https://charts.lakefs.io
# helm install my-lakefs lakefs/lakefs \
#   --set lakefsConfig.blockstore.type=s3 \
#   --set lakefsConfig.blockstore.s3.region=us-east-1
"""
lakefs_setup.py
Complete LakeFS workflow: create repo, commit data, branch, merge.
"""

import lakefs_client
from lakefs_client import models
from lakefs_client.client import LakeFSClient
import boto3
import pandas as pd
import io
from datetime import datetime


# ── STEP 1: Connect to LakeFS ─────────────────────────────────────────
configuration = lakefs_client.Configuration()
configuration.host        = "http://localhost:8000"
configuration.username    = "AKIAIOSFINEXAMPLE"   # Your LakeFS access key
configuration.password    = "wJalrXUtnFEMI..."    # Your LakeFS secret key

client = LakeFSClient(configuration)
print("✅ Connected to LakeFS")


# ── STEP 2: Create a Repository ───────────────────────────────────────
# A repository is the top-level container — like a GitHub repo but for data

try:
    repo = client.repositories.create_repository(
        repository_creation=models.RepositoryCreation(
            name              = "fraud-detection-data",
            storage_namespace = "s3://my-data-lake/fraud-detection/",
            default_branch    = "main",
        )
    )
    print(f"✅ Repository created: {repo.id}")
    print(f"   Storage: {repo.storage_namespace}")

except lakefs_client.exceptions.ApiException as e:
    if "already exists" in str(e):
        print("Repository already exists — continuing...")
    else:
        raise


# ── STEP 3: Upload Initial Dataset ────────────────────────────────────
# Create a synthetic fraud dataset for demonstration

def create_fraud_dataset(n_rows: int = 50000) -> pd.DataFrame:
    import numpy as np
    np.random.seed(42)

    return pd.DataFrame({
        "transaction_id":  range(1, n_rows + 1),
        "user_id":         np.random.randint(1000, 9999, n_rows),
        "amount":          np.random.exponential(50, n_rows).round(2),
        "merchant_country": np.random.choice(["US","UK","DE","NG","FR"], n_rows,
                                              p=[0.60, 0.15, 0.12, 0.08, 0.05]),
        "hour_of_day":     np.random.randint(0, 24, n_rows),
        "is_fraud":        np.random.choice([0, 1], n_rows, p=[0.97, 0.03]),
    })


# Convert DataFrame to CSV bytes for upload
df_train = create_fraud_dataset(50000)
train_bytes = df_train.to_csv(index=False).encode("utf-8")

# Upload to LakeFS main branch
client.objects.upload_object(
    repository = "fraud-detection-data",
    branch     = "main",
    path       = "datasets/train.csv",
    content    = io.BytesIO(train_bytes),
)
print(f"✅ Uploaded train.csv ({len(df_train):,} rows) to main branch")


# Upload test dataset too
df_test = create_fraud_dataset(10000)
test_bytes = df_test.to_csv(index=False).encode("utf-8")

client.objects.upload_object(
    repository = "fraud-detection-data",
    branch     = "main",
    path       = "datasets/test.csv",
    content    = io.BytesIO(test_bytes),
)
print(f"✅ Uploaded test.csv ({len(df_test):,} rows) to main branch")


# ── STEP 4: Commit the initial data ───────────────────────────────────
# Like 'git commit' but for your data lake!

commit_response = client.commits.commit(
    repository = "fraud-detection-data",
    branch     = "main",
    commit_creation = models.CommitCreation(
        message = "Initial dataset: fraud detection baseline v1.0",
        metadata = {
            "source":         "internal-transaction-log",
            "date_range":     "2024-Q4",
            "train_rows":     "50000",
            "test_rows":      "10000",
            "schema_version": "1.0",
            "prepared_by":    "data-engineering-team",
        }
    )
)

print(f"\n✅ Data committed!")
print(f"   Commit ID: {commit_response.id}")
print(f"   Message:   {commit_response.message}")
print(f"   Timestamp: {datetime.fromtimestamp(commit_response.creation_date)}")
# Output:
# Commit ID: a3f7b2c19d4e8f6a1b2c3d4e5f6789ab
# This ID permanently identifies this exact dataset state

9.2 Branching — Safe Data Experimentation

📌 Code Purpose — LakeFS Data Branching & Safe Experimentation

What this code does: Creates a feature-branch from the main data branch, applies data transformations (removing outliers, rebalancing classes), commits those changes to the branch, compares the branch against main, and then merges it back after validation.

Why it matters: Without branching, every data cleaning operation directly modifies production data. One mistake by a junior data scientist overwrites the dataset for the whole team. With LakeFS branches, changes are isolated until explicitly merged — exactly like pull requests in software development.

Key concept: LakeFS branches are copy-on-write — creating a branch is instantaneous and free (no data is actually copied). Storage is only consumed for the delta (what actually changed).
"""
lakefs_branching.py
Demonstrate safe data experimentation using LakeFS branches.
"""

# ── STEP 1: Create a branch for data cleaning work ────────────────────
# This is instantaneous — no data is copied!
# LakeFS uses copy-on-write under the hood.

client.branches.create_branch(
    repository = "fraud-detection-data",
    branch_creation = models.BranchCreation(
        name   = "feature/remove-outliers-and-rebalance",
        source = "main",            # Branch from current main state
    )
)
print("✅ Branch created: feature/remove-outliers-and-rebalance")
print("   (This was instantaneous — no data was copied!)")


# ── STEP 2: Read data from the branch, transform it ───────────────────
# Use LakeFS S3 gateway — just change s3://bucket to lakefs://repo/branch/

import boto3
import pandas as pd

# Configure boto3 to talk to LakeFS instead of real S3
s3_client = boto3.client(
    "s3",
    endpoint_url          = "http://localhost:8000",  # LakeFS S3 gateway
    aws_access_key_id     = "AKIAIOSFINEXAMPLE",
    aws_secret_access_key = "wJalrXUtnFEMI...",
)

# Read from the feature branch (same API as reading from real S3!)
response = s3_client.get_object(
    Bucket = "fraud-detection-data",     # repository name
    Key    = "feature/remove-outliers-and-rebalance/datasets/train.csv"
)
df = pd.read_csv(response["Body"])
print(f"Loaded {len(df):,} rows from branch")


# ── STEP 3: Apply data transformations on the branch ──────────────────

original_rows = len(df)

# Transformation 1: Remove extreme outliers
df = df[df["amount"] <= df["amount"].quantile(0.999)]
removed_outliers = original_rows - len(df)
print(f"Removed {removed_outliers} extreme outlier transactions")

# Transformation 2: Report class balance
fraud_rate_before = df["is_fraud"].mean()
print(f"Fraud rate: {fraud_rate_before:.2%} (before rebalancing)")

# Transformation 3: Oversample fraud cases (simple upsampling)
fraud_cases    = df[df["is_fraud"] == 1]
legit_cases    = df[df["is_fraud"] == 0]
fraud_upsampled = fraud_cases.sample(
    n          = int(len(legit_cases) * 0.1),
    replace    = True,
    random_state = 42
)
df_balanced = pd.concat([legit_cases, fraud_upsampled]).sample(frac=1, random_state=42)
fraud_rate_after = df_balanced["is_fraud"].mean()
print(f"Fraud rate: {fraud_rate_after:.2%} (after rebalancing)")
print(f"Final dataset: {len(df_balanced):,} rows")


# ── STEP 4: Write the transformed data back to the branch ─────────────
transformed_bytes = df_balanced.to_csv(index=False).encode("utf-8")
import io

client.objects.upload_object(
    repository = "fraud-detection-data",
    branch     = "feature/remove-outliers-and-rebalance",
    path       = "datasets/train.csv",
    content    = io.BytesIO(transformed_bytes),
)

# Commit the transformation to the branch
client.commits.commit(
    repository = "fraud-detection-data",
    branch     = "feature/remove-outliers-and-rebalance",
    commit_creation = models.CommitCreation(
        message = "Data cleaning: remove outliers + rebalance classes",
        metadata = {
            "removed_outliers":    str(removed_outliers),
            "original_rows":       str(original_rows),
            "final_rows":          str(len(df_balanced)),
            "fraud_rate_before":   f"{fraud_rate_before:.4f}",
            "fraud_rate_after":    f"{fraud_rate_after:.4f}",
            "transformation_by":   "alice@company.com",
            "reviewed":            "false",  # Needs review before merge
        }
    )
)
print("✅ Transformation committed to branch")


# ── STEP 5: Compare branch vs main (diff) ─────────────────────────────
diff_response = client.refs.diff_refs(
    repository = "fraud-detection-data",
    left_ref   = "main",
    right_ref  = "feature/remove-outliers-and-rebalance",
)

print(f"\n📊 Diff: main ← feature/remove-outliers-and-rebalance")
print(f"   Changed objects: {len(diff_response.results)}")
for change in diff_response.results[:5]:
    print(f"   {change.type:8s} {change.path}")
# Output:
# Changed objects: 1
# changed  datasets/train.csv


# ── STEP 6: Merge the branch back into main (after review) ────────────
# This is like approving a Pull Request — but for data!

merge_result = client.refs.merge_into_branch(
    repository = "fraud-detection-data",
    source_ref = "feature/remove-outliers-and-rebalance",
    destination_branch = "main",
    merge   = models.Merge(
        message = "Merge: data cleaning pipeline approved by data-engineering lead",
    )
)

print(f"\n✅ Branch merged into main!")
print(f"   Merge commit: {merge_result.reference}")
print(f"   Main branch now points to cleaned, balanced dataset")


# ── STEP 7: Tag the new version for experiments ───────────────────────
client.tags.create_tag(
    repository   = "fraud-detection-data",
    tag_creation = models.TagCreation(
        id  = "v1.1-cleaned-balanced",
        ref = "main",
    )
)
print("✅ Tagged as: v1.1-cleaned-balanced")
print("   Any experiment can now point to this exact tag and be reproducible!")

9.3 CI/CD for Data — Automated Data Quality Gates

📌 Code Purpose — Automated Data Quality Checks (Data CI/CD)

What this code does: Builds a data validation pipeline that runs automatically whenever new data is committed to a branch. It checks schema integrity, null rates, statistical properties, and business rule compliance — and only allows the merge to main if ALL checks pass.

Why it matters: This is the data equivalent of a CI/CD pipeline for code. Just as bad code doesn't get merged without passing tests, bad data doesn't get merged without passing quality checks. This pattern prevents corrupt or drifted data from ever reaching production models — the root cause of most silent model failures.

Key concept: This function would be triggered by a LakeFS webhook or run in a GitHub Actions / Jenkins / Airflow pipeline whenever a branch is proposed for merging into main.
"""
data_quality_gate.py
Automated data quality checks that run before any branch can merge to main.
This is the 'CI/CD pipeline' for your data.
"""

import pandas as pd
import numpy as np
from dataclasses import dataclass
from typing import List, Dict, Any
from datetime import datetime


@dataclass
class QualityCheckResult:
    check_name:  str
    passed:      bool
    severity:    str      # "CRITICAL", "WARNING", "INFO"
    message:     str
    actual_value: Any = None
    expected:    Any = None


class DataQualityGate:
    """
    Runs a battery of quality checks on a dataset.
    All CRITICAL checks must pass for the merge to be allowed.
    """

    def __init__(self, config: Dict):
        self.config  = config
        self.results = []

    def check_schema(self, df: pd.DataFrame) -> None:
        """Verify all required columns exist with correct dtypes."""
        required_cols = self.config.get("required_columns", {})

        for col, expected_dtype in required_cols.items():
            if col not in df.columns:
                self.results.append(QualityCheckResult(
                    check_name  = f"schema.column_exists.{col}",
                    passed      = False,
                    severity    = "CRITICAL",
                    message     = f"Required column '{col}' is MISSING from dataset",
                ))
            else:
                actual_dtype = str(df[col].dtype)
                dtype_ok     = expected_dtype in actual_dtype

                self.results.append(QualityCheckResult(
                    check_name   = f"schema.dtype.{col}",
                    passed       = dtype_ok,
                    severity     = "CRITICAL" if not dtype_ok else "INFO",
                    message      = f"Column '{col}': dtype={'✅' if dtype_ok else '❌'}",
                    actual_value = actual_dtype,
                    expected     = expected_dtype,
                ))

    def check_row_count(self, df: pd.DataFrame) -> None:
        """Verify row count is within expected bounds."""
        min_rows = self.config.get("min_rows", 0)
        max_rows = self.config.get("max_rows", float("inf"))
        n        = len(df)

        passed = min_rows <= n <= max_rows

        self.results.append(QualityCheckResult(
            check_name   = "volume.row_count",
            passed       = passed,
            severity     = "CRITICAL" if not passed else "INFO",
            message      = (f"Row count {n:,} "
                            f"{'within' if passed else 'OUTSIDE'} "
                            f"expected range [{min_rows:,}, {max_rows:,}]"),
            actual_value = n,
            expected     = f"{min_rows} – {max_rows}",
        ))

    def check_null_rates(self, df: pd.DataFrame) -> None:
        """Alert if any column's null rate exceeds threshold."""
        max_null = self.config.get("max_null_rate", 0.05)

        for col in df.columns:
            null_rate = df[col].isnull().mean()
            passed    = null_rate <= max_null

            if null_rate > 0:
                self.results.append(QualityCheckResult(
                    check_name   = f"quality.null_rate.{col}",
                    passed       = passed,
                    severity     = "CRITICAL" if null_rate > 0.20 else
                                   "WARNING"  if not passed else "INFO",
                    message      = (f"Column '{col}': null rate = {null_rate:.2%} "
                                    f"(threshold: {max_null:.0%})"),
                    actual_value = null_rate,
                    expected     = f"<= {max_null}",
                ))

    def check_value_ranges(self, df: pd.DataFrame) -> None:
        """Verify numeric columns stay within expected ranges."""
        ranges = self.config.get("value_ranges", {})

        for col, bounds in ranges.items():
            if col not in df.columns:
                continue

            min_val  = bounds.get("min", float("-inf"))
            max_val  = bounds.get("max", float("inf"))
            col_min  = df[col].min()
            col_max  = df[col].max()

            out_of_range = ((df[col] < min_val) | (df[col] > max_val)).sum()
            passed       = out_of_range == 0

            self.results.append(QualityCheckResult(
                check_name   = f"quality.range.{col}",
                passed       = passed,
                severity     = "CRITICAL" if out_of_range > len(df) * 0.01 else "WARNING",
                message      = (f"Column '{col}': {out_of_range} values out of range "
                                f"[{min_val}, {max_val}]. "
                                f"Actual: [{col_min:.2f}, {col_max:.2f}]"),
                actual_value = out_of_range,
            ))

    def check_class_balance(self, df: pd.DataFrame) -> None:
        """Check that target class is not too imbalanced."""
        label_col     = self.config.get("label_column")
        if not label_col or label_col not in df.columns:
            return

        min_minority  = self.config.get("min_minority_rate", 0.01)
        fraud_rate    = df[label_col].mean()
        passed        = fraud_rate >= min_minority

        self.results.append(QualityCheckResult(
            check_name   = "distribution.class_balance",
            passed       = passed,
            severity     = "WARNING" if not passed else "INFO",
            message      = (f"Minority class rate: {fraud_rate:.3%} "
                            f"(min required: {min_minority:.1%})"),
            actual_value = fraud_rate,
            expected     = f">= {min_minority}",
        ))

    def run_all_checks(self, df: pd.DataFrame) -> bool:
        """Run all configured checks. Returns True if all CRITICAL checks pass."""
        self.results = []

        self.check_schema(df)
        self.check_row_count(df)
        self.check_null_rates(df)
        self.check_value_ranges(df)
        self.check_class_balance(df)

        # Determine overall pass/fail (only CRITICAL matters for gate)
        critical_failures = [r for r in self.results
                             if r.severity == "CRITICAL" and not r.passed]
        return len(critical_failures) == 0

    def print_report(self) -> None:
        """Print a formatted quality gate report."""
        critical_failures = [r for r in self.results
                             if r.severity == "CRITICAL" and not r.passed]
        warnings          = [r for r in self.results
                             if r.severity == "WARNING" and not r.passed]

        print("\n" + "="*60)
        print("  🔍 DATA QUALITY GATE REPORT")
        print(f"  Generated: {datetime.utcnow().strftime('%Y-%m-%d %H:%M UTC')}")
        print("="*60)

        for result in self.results:
            icon = "✅" if result.passed else ("🚨" if result.severity == "CRITICAL" else "⚠️")
            print(f"  {icon} {result.check_name}")
            if not result.passed:
                print(f"       → {result.message}")

        print("="*60)

        if critical_failures:
            print(f"\n  🚨 GATE: FAILED — {len(critical_failures)} critical issue(s)")
            print(f"  ❌ Merge to main is BLOCKED until issues are resolved.")
        else:
            print(f"\n  ✅ GATE: PASSED")
            if warnings:
                print(f"  ⚠️  {len(warnings)} warning(s) — review recommended")
            print(f"  ✅ Branch is APPROVED for merge to main.")

        print("="*60 + "\n")


# ─── DEMO ──────────────────────────────────────────────────────────────

# Quality gate configuration
gate_config = {
    "required_columns": {
        "transaction_id":  "int",
        "user_id":         "int",
        "amount":          "float",
        "is_fraud":        "int",
    },
    "min_rows":       10000,
    "max_rows":       10000000,
    "max_null_rate":  0.05,
    "value_ranges": {
        "amount":     {"min": 0.0,  "max": 100000.0},
        "is_fraud":   {"min": 0,    "max": 1},
        "hour_of_day":{"min": 0,    "max": 23},
    },
    "label_column":       "is_fraud",
    "min_minority_rate":  0.01,
}

# Simulate a dataset to validate
import numpy as np
np.random.seed(42)
n = 50000

test_df = pd.DataFrame({
    "transaction_id": range(1, n + 1),
    "user_id":        np.random.randint(1000, 9999, n),
    "amount":         np.random.exponential(50, n).round(2),
    "merchant_country": np.random.choice(["US","UK","DE"], n),
    "hour_of_day":    np.random.randint(0, 24, n),
    "is_fraud":       np.random.choice([0, 1], n, p=[0.97, 0.03]),
})

gate    = DataQualityGate(config=gate_config)
passed  = gate.run_all_checks(test_df)
gate.print_report()

if passed:
    print("🚀 Proceeding with branch merge to main...")
else:
    print("🛑 Merge blocked. Fix data issues and re-run quality gate.")

🔗 Section 10: Data Versioning + Feature Store — The Complete MLOps Pipeline

Data Versioning and Feature Stores are complementary pieces of the same puzzle. Here's how they fit together in a complete MLOps pipeline.

COMPLETE MLOPS DATA PIPELINE:

┌──────────────┐ ┌─────────────────┐ ┌────────────────────┐
│ RAW DATA │ │ DATA LAKE │ │ FEATURE STORE │
│ SOURCES │───▶│ (LakeFS) │───▶│ (Feast/Tecton) │
│ │ │ │ │ │
│ - Events │ │ Versioned & │ │ - Feature Views │
│ - DB dumps │ │ branched raw │ │ - Online Store │
│ - APIs │ │ + cleaned data │ │ - Offline Store │
└──────────────┘ └────────┬────────┘ └──────────┬─────────┘
│ │
▼ ▼
┌─────────────────┐ ┌────────────────────┐
│ DVC PIPELINE │ │ ML EXPERIMENTS │
│ │ │ │
│ - Prepare │ │ Model A: v1 data │
│ - Featurize │ │ Model B: v2 data │
│ - Train │ │ Each exp links to │
│ - Evaluate │ │ exact data version│
└─────────────────┘ └────────────────────┘

KEY PRINCIPLE: Every ML model can trace its exact data lineage
from raw source → cleaned lake → feature store → training dataset ✅
📌 Code Purpose — Linking Data Version to Model Training

What this code does: Shows the complete handoff from LakeFS (raw data versioning) to Feast (feature store) to DVC (experiment tracking), ensuring every model training run records the exact data commit it used — creating an end-to-end data lineage chain.

Why it matters: This closes the reproducibility loop completely. Given any model in production, you can trace back: which experiment produced it → which DVC pipeline run → which feature store version → which LakeFS commit → which raw data source. Full audit trail from prediction back to original byte. 🔍
"""
end_to_end_pipeline.py
Complete pipeline showing data versioning at every stage:
LakeFS (raw data) → Feature Store → DVC (model training)
"""

import subprocess
import json
from datetime import datetime
import pandas as pd


def get_lakefs_commit_id(repo: str, branch: str) -> str:
    """Fetch the current HEAD commit from LakeFS branch."""
    import lakefs_client
    configuration = lakefs_client.Configuration()
    configuration.host     = "http://localhost:8000"
    configuration.username = "AKIAIOSFINEXAMPLE"
    configuration.password = "wJalrXUtnFEMI..."

    client = lakefs_client.client.LakeFSClient(configuration)
    branch_info = client.branches.get_branch(
        repository  = repo,
        branch      = branch,
    )
    return branch_info.commit_id


def get_feature_store_snapshot(feature_view: str) -> dict:
    """Record which feature store version is being used."""
    from feast import FeatureStore
    store = FeatureStore(repo_path="./feature_repo")
    fv    = store.get_feature_view(feature_view)

    return {
        "feature_view":      feature_view,
        "feature_view_hash": fv.created_timestamp.isoformat()
                             if fv.created_timestamp else "unknown",
        "features":          [f.name for f in fv.features],
        "entities":          [e.name for e in fv.entities],
    }


def train_model_with_full_lineage(
    lakefs_repo:    str,
    lakefs_branch:  str,
    feature_view:   str,
    model_params:   dict,
) -> dict:
    """
    Complete training run with full data lineage recording.

    Returns a lineage record linking:
    - LakeFS raw data commit
    - Feature store snapshot
    - Training parameters
    - Model metrics
    """

    print("=" * 60)
    print("  🔗 TRAINING WITH FULL DATA LINEAGE")
    print("=" * 60)

    # ── Step 1: Record LakeFS data version ────────────────────────────
    lakefs_commit = get_lakefs_commit_id(lakefs_repo, lakefs_branch)
    print(f"\n📦 Raw Data Version (LakeFS)")
    print(f"   Repository: {lakefs_repo}")
    print(f"   Branch:     {lakefs_branch}")
    print(f"   Commit ID:  {lakefs_commit}")

    # ── Step 2: Record Feature Store version ──────────────────────────
    fs_snapshot = get_feature_store_snapshot(feature_view)
    print(f"\n🏪 Feature Store Snapshot")
    print(f"   Feature View: {fs_snapshot['feature_view']}")
    print(f"   Features:     {', '.join(fs_snapshot['features'][:3])}...")

    # ── Step 3: Record DVC data version ───────────────────────────────
    dvc_md5 = subprocess.run(
        ["dvc", "status", "--show-json"],
        capture_output=True, text=True
    ).stdout

    dvc_commit = subprocess.run(
        ["git", "rev-parse", "HEAD"],
        capture_output=True, text=True
    ).stdout.strip()

    print(f"\n🔧 DVC + Git Version")
    print(f"   Git commit: {dvc_commit[:8]}...")

    # ── Step 4: Run the actual training ───────────────────────────────
    # (Simplified — in practice this calls dvc repro or your training script)
    import numpy as np
    from sklearn.ensemble import GradientBoostingClassifier
    from sklearn.metrics import roc_auc_score

    np.random.seed(model_params.get("random_seed", 42))
    X_train = np.random.randn(10000, 10)
    y_train = np.random.choice([0,1], 10000, p=[0.97, 0.03])
    X_val   = np.random.randn(2000, 10)
    y_val   = np.random.choice([0,1], 2000, p=[0.97, 0.03])

    model = GradientBoostingClassifier(
        n_estimators  = model_params.get("n_estimators", 100),
        learning_rate = model_params.get("learning_rate", 0.1),
        max_depth     = model_params.get("max_depth", 3),
        random_state  = model_params.get("random_seed", 42),
    )
    model.fit(X_train, y_train)
    auc = roc_auc_score(y_val, model.predict_proba(X_val)[:, 1])

    print(f"\n🎯 Training Results")
    print(f"   ROC-AUC: {auc:.4f}")

    # ── Step 5: Build complete lineage record ─────────────────────────
    lineage = {
        "training_run_id":   f"run_{datetime.utcnow().strftime('%Y%m%d_%H%M%S')}",
        "timestamp":         datetime.utcnow().isoformat(),
        "data_lineage": {
            "lakefs_repository": lakefs_repo,
            "lakefs_branch":     lakefs_branch,
            "lakefs_commit_id":  lakefs_commit,      # Raw data fingerprint
            "feature_store":     fs_snapshot,         # Feature version
            "git_commit":        dvc_commit,           # Code version
        },
        "model_params":  model_params,
        "metrics": {
            "roc_auc": round(auc, 4),
        },
    }

    # Save lineage to file (tracked by DVC as a metric)
    with open("reports/model_lineage.json", "w") as f:
        json.dump(lineage, f, indent=2)

    print(f"\n✅ Full lineage record saved to reports/model_lineage.json")
    print(f"   Run ID: {lineage['training_run_id']}")
    print(f"\n📋 LINEAGE SUMMARY:")
    print(f"   Raw data:      LakeFS commit {lakefs_commit[:8]}...")
    print(f"   Features:      {len(fs_snapshot['features'])} features from {feature_view}")
    print(f"   Code:          Git commit {dvc_commit[:8]}...")
    print(f"   Model metric:  ROC-AUC = {auc:.4f}")
    print(f"\n🔍 This model is FULLY REPRODUCIBLE from end to end!")

    return lineage


# ─── Run the complete pipeline ──────────────────────────────────────

lineage = train_model_with_full_lineage(
    lakefs_repo   = "fraud-detection-data",
    lakefs_branch = "main",
    feature_view  = "user_transaction_features",
    model_params  = {
        "n_estimators":  200,
        "learning_rate": 0.001,
        "max_depth":     4,
        "random_seed":   42,
    }
)

🛠️ Section 11: Other Data Versioning Tools

DVC and LakeFS are the two most important tools to know. But the data versioning ecosystem is rich. Here's your complete map.

Tool Best For Key Feature Open Source?
DVC ⭐ ML project data versioning Git-like, pipeline DAG, experiments ✅ Yes
LakeFS ⭐ Entire data lake versioning S3-compatible, org-wide branching ✅ Yes (community)
Delta Lake Spark / Databricks pipelines ACID transactions, time travel, schema enforcement ✅ Yes
Apache Iceberg Petabyte-scale table versioning Time travel, hidden partitioning, multi-engine ✅ Yes
Apache Hudi Streaming data lakes Upserts, incremental queries, record-level versioning ✅ Yes
Pachyderm Data-driven pipelines Auto re-runs pipelines when data changes ✅ Community
Weights & Biases Artifacts ML experiment artifact tracking Tight integration with W&B experiment tracking UI Freemium
MLflow + datasets MLflow-centric teams Log dataset fingerprints alongside experiment metrics ✅ Yes
Vertex AI Dataset Versioning GCP-native ML teams Managed, integrated with Vertex AI pipelines ❌ Managed service
💡 Decision Guide — Which Tool to Use?

Just starting out / small project: DVC — easiest to learn, works locally.
Enterprise data lake on S3/GCS: LakeFS — org-wide governance.
Already using Databricks/Spark: Delta Lake — already built in.
Petabyte analytics engine: Apache Iceberg — best at massive scale.
Streaming pipeline with upserts: Apache Hudi.
Using W&B for experiment tracking: W&B Artifacts — seamless integration.

✅ Section 12: Best Practices & Anti-Patterns

✅ DOs — Data Versioning Best Practices:
  • ✅ Version data from Day 1 — retroactively adding versioning is painful
  • ✅ Write commit messages for data like you write them for code — be descriptive
  • ✅ Always tag dataset versions used for major model releases (e.g., v2.0-production)
  • ✅ Record the data version hash in every model's metadata file
  • ✅ Run automated data quality checks before every merge to main
  • ✅ Separate raw, cleaned, and feature-engineered data as distinct versioned layers
  • ✅ Link your Git commits to your DVC data versions — they should always be in sync
  • ✅ Use LakeFS branches for any experimental data transformation — never transform main directly
  • ✅ Test that you can actually reproduce an experiment from 3 months ago — do this quarterly
  • ✅ Store data quality statistics (row count, null rates, distributions) with each version
❌ DON'Ts — Anti-Patterns That Destroy Reproducibility:
  • ❌ Don't manually rename files as versioning ("train_v2_FINAL_use_this.csv")
  • ❌ Don't overwrite data files in place — always create new versions
  • ❌ Don't commit large data files directly into Git — use DVC pointers
  • ❌ Don't delete old dataset versions "to save space" — use lifecycle policies instead
  • ❌ Don't run experiments without recording which data version was used
  • ❌ Don't mix data cleaning and feature engineering in the same pipeline stage — version each separately
  • ❌ Don't skip data versioning for "small" projects — you will regret it when the project grows
  • ❌ Don't assume a file at the same S3 path is the same data — paths lie, content hashes don't

🏆 High-Level Summary: Everything You've Mastered!

  • 🔹 Data Versioning = treating datasets like code. Every state recorded, tagged, and restorable.
  • 🔹 Why it's needed: Data changes constantly. Without versioning, experiments are irreproducible, debugging is guesswork, and compliance is impossible.
  • 🔹 4 Pillars of Reproducibility: Code (Git) + Environment (Docker) + Config (MLflow) + Data (DVC/LakeFS).
  • 🔹 DVC: Git for ML data. Pointer files in Git, data in S3. Pipeline DAG with smart caching. Experiment tracking with data version links.
  • 🔹 LakeFS: Git branching for your entire data lake. S3-compatible. Branch → transform → validate → merge. Full organization-wide data governance.
  • 🔹 Data CI/CD: Automated quality gates that block bad data from merging to main — just like PR checks block bad code.
  • 🔹 Complete MLOps chain: LakeFS (raw data) → Feature Store (computed features) → DVC (experiment tracking) → Model Registry (deployment).
  • 🔹 Other key tools: Delta Lake (Spark), Apache Iceberg (petabyte scale), Apache Hudi (streaming), W&B Artifacts (experiment tracking).
  • 🔹 Reality: EU AI Act and healthcare/finance regulations now legally require data provenance. Data versioning is a compliance requirement, not optional.
The best ML teams in the world treat data with the same engineering discipline they apply to code. Every commit matters. Every version is traceable. Every experiment is reproducible.

You now have the knowledge and tools to build that discipline from day one — whether you're a solo data scientist on a laptop or an engineering team managing petabytes of production data.

Keep versioning, keep reproducing, keep shipping trustworthy AI! 🐼✨

Comments