簡介
您可以使用 Git 對程式碼進行版本控制。您使用 MLflow 追蹤的模型。但是**數據呢? **
資料是機器學習中最重要的部分,但它通常沒有版本控制。結果:“舊模型更好,但舊數據已被覆蓋...”
🎯 本文: 用於資料版本控制的 DVC + 用於特徵儲存的盛宴。
1. 為什麼需要資料版本控制?
Vấn đề thực tế:
❌ Data engineer update CSV → model accuracy tụt
❌ "Rollback data về tuần trước" → không có version
❌ Train/serve data khác nhau (Training-Serving Skew)
❌ Không biết model nào dùng data version nào
Giải pháp:
✅ Version data giống version code
✅ Mỗi data version → hash, metadata
✅ Link data version ↔ model version
✅ Reproducibility: quay lại bất kỳ thời điểm nào
2.DVC——資料版本控制
2.1 設置
# Install
pip install dvc dvc-s3 # hoặc dvc-gcs, dvc-azure
# Init DVC trong Git repo
cd my-ml-project
git init
dvc init
# Config remote storage
dvc remote add -d myremote s3://my-bucket/dvc-storage
# hoặc local: dvc remote add -d myremote /path/to/storage
# hoặc GCS: dvc remote add -d myremote gs://my-bucket/dvc
git add .dvc .dvcignore
git commit -m "Init DVC"
2.2 軌跡數據
# Track data files
dvc add data/raw/train.csv
dvc add data/raw/images/
# → Tạo train.csv.dvc và images.dvc (metadata files)
# Push data lên remote
dvc push
# Git commit metadata
git add data/raw/train.csv.dvc data/raw/images.dvc .gitignore
git commit -m "Add training data v1"
git tag data-v1.0
# Workflow khi data thay đổi
# 1. Update data
cp new_data.csv data/raw/train.csv
# 2. Track changes
dvc add data/raw/train.csv
dvc push
# 3. Commit
git add data/raw/train.csv.dvc
git commit -m "Update training data v2 - add 500 new samples"
git tag data-v2.0
# 4. Rollback nếu cần
git checkout data-v1.0 -- data/raw/train.csv.dvc
dvc checkout # Pull data v1.0 từ remote
2.3 DVC 管道
# dvc.yaml — Define reproducible ML pipeline
stages:
prepare:
cmd: python src/data/prepare.py
deps:
- src/data/prepare.py
- data/raw/train.csv
outs:
- data/processed/train_clean.csv
- data/processed/val_clean.csv
featurize:
cmd: python src/features/build_features.py
deps:
- src/features/build_features.py
- data/processed/train_clean.csv
params:
- configs/features.yaml:
- feature_columns
- scaling_method
outs:
- data/features/train_features.parquet
- data/features/val_features.parquet
train:
cmd: python src/models/train.py
deps:
- src/models/train.py
- data/features/train_features.parquet
params:
- configs/training.yaml:
- model_type
- learning_rate
- n_estimators
outs:
- models/model.pkl
metrics:
- metrics/train_metrics.json:
cache: false
evaluate:
cmd: python src/models/evaluate.py
deps:
- src/models/evaluate.py
- models/model.pkl
- data/features/val_features.parquet
metrics:
- metrics/eval_metrics.json:
cache: false
plots:
- metrics/plots/:
cache: false
# Reproduce toàn bộ pipeline
dvc repro
# Chỉ chạy stage bị thay đổi
dvc repro train # Nếu chỉ đổi training config
# Xem pipeline DAG
dvc dag
# prepare → featurize → train → evaluate
# So sánh metrics giữa các version
dvc metrics diff
# Path Metric HEAD workspace
# metrics/eval.json accuracy 0.89 0.92
# metrics/eval.json f1_score 0.87 0.91
2.4 DVC + Git 分支 = 每個分支進行實驗
# Branch-based experiments
git checkout -b experiment/new-features
# Modify features
python src/features/build_features.py
dvc repro
# Compare với main
dvc metrics diff main
dvc plots diff main
# Nếu tốt hơn → merge
git checkout main
git merge experiment/new-features
dvc push
3. 資料品質驗證
"""Data validation — Great Expectations"""
# pip install great_expectations
import great_expectations as gx
# Setup
context = gx.get_context()
# Define expectations
validator = context.sources.pandas_default.read_csv(
"data/processed/train.csv"
)
# Expectations
validator.expect_column_values_to_not_be_null("user_id")
validator.expect_column_values_to_be_between("age", min_value=0, max_value=150)
validator.expect_column_values_to_be_in_set(
"status", ["active", "inactive", "churned"]
)
validator.expect_column_mean_to_be_between("revenue", min_value=10, max_value=1000)
validator.expect_table_row_count_to_be_between(min_value=1000, max_value=1000000)
# Validate
results = validator.validate()
if not results.success:
print("❌ Data validation FAILED!")
for result in results.results:
if not result.success:
print(f" FAIL: {result.expectation_config.expectation_type}")
"""Data validation đơn giản với Pandera"""
# pip install pandera
import pandera as pa
import pandas as pd
# Define schema
schema = pa.DataFrameSchema({
"user_id": pa.Column(int, pa.Check.gt(0), nullable=False),
"age": pa.Column(int, pa.Check.in_range(0, 150)),
"revenue": pa.Column(float, pa.Check.ge(0)),
"status": pa.Column(str, pa.Check.isin(["active", "inactive", "churned"])),
"signup_date": pa.Column(pd.Timestamp),
})
# Validate
df = pd.read_csv("data/train.csv")
try:
validated_df = schema.validate(df)
print("✅ Data is valid!")
except pa.errors.SchemaError as e:
print(f"❌ Validation failed: {e}")
4. 功能商店 — 盛宴
4.1 為什麼我們需要特徵儲存?
Vấn đề:
❌ Training dùng feature A tính theo logic X
❌ Serving dùng feature A tính theo logic Y
→ Training-Serving Skew → Model fail
❌ Team A & Team B cùng tính "user_age" nhưng khác logic
❌ Feature computation lặp lại giữa các model
❌ Realtime features rất khó implement
Feature Store giải quyết:
✅ Single source of truth cho features
✅ Offline (batch) & Online (realtime) serving
✅ Feature reuse across teams & models
✅ Point-in-time correct joins
✅ Feature versioning & discovery
4.2 盛宴設置
# Install
pip install feast
# Init project
feast init my_feature_store
cd my_feature_store
"""feature_store.yaml — Feast config"""
# feature_store.yaml
"""
project: my_project
registry: data/registry.db
provider: local
online_store:
type: sqlite
path: data/online_store.db
offline_store:
type: file
entity_key_serialization_version: 2
"""
4.3 定義特徵
"""features.py — Feature definitions"""
from datetime import timedelta
from feast import Entity, FeatureView, Field, FileSource
from feast.types import Float32, Int64, String
# Entity = đối tượng (user, product, ...)
user = Entity(
name="user_id",
join_keys=["user_id"],
description="User entity",
)
# Data source
user_stats_source = FileSource(
path="data/user_stats.parquet",
timestamp_field="event_timestamp",
created_timestamp_column="created_timestamp",
)
# Feature View
user_stats_fv = FeatureView(
name="user_stats",
entities=[user],
ttl=timedelta(days=1), # Feature freshness
schema=[
Field(name="total_purchases", dtype=Int64),
Field(name="avg_order_value", dtype=Float32),
Field(name="days_since_last_order", dtype=Int64),
Field(name="favorite_category", dtype=String),
Field(name="lifetime_value", dtype=Float32),
],
source=user_stats_source,
online=True, # Có serve online
)
4.4 使用特徵庫
"""Feature retrieval cho training & serving"""
from feast import FeatureStore
import pandas as pd
store = FeatureStore(repo_path=".")
# === Apply feature definitions ===
# feast apply (CLI)
# === OFFLINE: Get features cho Training ===
entity_df = pd.DataFrame({
"user_id": [1, 2, 3, 4, 5],
"event_timestamp": pd.to_datetime("2024-01-15"),
})
training_df = store.get_historical_features(
entity_df=entity_df,
features=[
"user_stats:total_purchases",
"user_stats:avg_order_value",
"user_stats:days_since_last_order",
"user_stats:lifetime_value",
],
).to_df()
print(training_df)
# user_id | total_purchases | avg_order_value | days_since_last | ltv
# 1 | 15 | 45.2 | 3 | 678.0
# 2 | 3 | 120.0 | 45 | 360.0
# ...
# === ONLINE: Get features cho Serving (real-time) ===
# Materialize (push offline → online store)
# feast materialize-incremental $(date +%Y-%m-%dT%H:%M:%S)
online_features = store.get_online_features(
features=[
"user_stats:total_purchases",
"user_stats:avg_order_value",
"user_stats:lifetime_value",
],
entity_rows=[
{"user_id": 1},
{"user_id": 2},
],
).to_dict()
print(online_features)
# {'user_id': [1, 2], 'total_purchases': [15, 3], ...}
4.5 特徵工程流程
"""Feature engineering pipeline kết hợp với Feast"""
import pandas as pd
from datetime import datetime
def compute_user_features(raw_events_df):
"""Tính features từ raw events"""
now = datetime.now()
features = raw_events_df.groupby("user_id").agg(
total_purchases=("order_id", "count"),
avg_order_value=("amount", "mean"),
total_revenue=("amount", "sum"),
last_order_date=("order_date", "max"),
).reset_index()
# Derived features
features["days_since_last_order"] = (
now - features["last_order_date"]
).dt.days
features["lifetime_value"] = features["total_revenue"]
# Categorical
category_mode = raw_events_df.groupby("user_id")["category"].agg(
lambda x: x.mode()[0] if len(x.mode()) > 0 else "unknown"
).reset_index()
category_mode.columns = ["user_id", "favorite_category"]
features = features.merge(category_mode, on="user_id")
# Add timestamps for Feast
features["event_timestamp"] = now
features["created_timestamp"] = now
return features
# Chạy pipeline
raw_data = pd.read_parquet("data/raw_events.parquet")
user_features = compute_user_features(raw_data)
user_features.to_parquet("data/user_stats.parquet", index=False)
# Apply & materialize
# feast apply
# feast materialize-incremental $(date +%Y-%m-%dT%H:%M:%S)
5. 最佳實踐
Data Versioning:
✅ NEVER modify raw data — luôn tạo bản processed
✅ Tag data versions (data-v1.0, data-v2.0)
✅ Link data version ↔ model version
✅ Validate data trước khi train (schema, stats)
✅ Automate data pipeline (DVC repro)
Feature Store:
✅ Centralize feature definitions
✅ Same code cho offline & online
✅ Monitor feature freshness
✅ Document features (what, why, how computed)
✅ Point-in-time correct joins (avoid data leakage)
總結
| 概念 | 記住 |
|---|---|
| DVC | Git 用於資料、版本+管道+複製 |
| dvc.yaml | 定義可重複的機器學習管道 |
| 資料驗證 | 遠大前程 / Pandera — 訓練前驗證 |
| 功能商店 | 盛宴-功能的單一事實來源 |
| 線下商店 | 培訓的歷史特徵 |
| 網上商店 | 即時服務功能 |
| 訓練-服務偏誤 | 特徵儲存解決 |
練習
- DVC 設定: 為 ML 專案初始化 DVC。 Track 1 資料集,推送到遠端。回滾版本。
- 管道: 建立具有 3 個階段的 DVC 管道:準備 → 訓練 → 評估。使用
dvc repro。 - 資料驗證: 為資料集編寫 Pandera 架構。使用錯誤格式的資料進行測試。
- Feast: 在本地設定Feast,定義5個特徵,具體化,取得線上特徵。
下一篇文章: 模型註冊表、版本控制和打包。