Chuyển đến nội dung chính

第 3 課:資料版本控制與特徵存儲

DVC(資料版本控制)設定和工作流程。數據管道。用於資料湖版本控制的 LakeFS。特徵儲存:盛宴設定、特徵工程管道、線上/離線服務。數據品質驗證。

🧠 人工智慧與機器學習 — 第 2 課 第 3 課:資料版本控制與特徵存儲

MLOps 和 LLMOps:將 AI 引入生產

第 1 部分:MLOps 基礎

亞洲開發網

簡介

您可以使用 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)

總結

概念記住
DVCGit 用於資料、版本+管道+複製
dvc.yaml定義可重複的機器學習管道
資料驗證遠大前程 / Pandera — 訓練前驗證
功能商店盛宴-功能的單一事實來源
線下商店培訓的歷史特徵
網上商店即時服務功能
訓練-服務偏誤特徵儲存解決

練習

  1. DVC 設定: 為 ML 專案初始化 DVC。 Track 1 資料集,推送到遠端。回滾版本。
  2. 管道: 建立具有 3 個階段的 DVC 管道:準備 → 訓練 → 評估。使用 dvc repro。
  3. 資料驗證: 為資料集編寫 Pandera 架構。使用錯誤格式的資料進行測試。
  4. Feast: 在本地設定Feast,定義5個特徵,具體化,取得線上特徵。

下一篇文章: 模型註冊表、版本控制和打包。