Python版本:Python 3.12+
开发工具:PyCharm 或 VS Code
操作系统:Windows / macOS / Linux (通用)
核心依赖:FastAPI 0.115+、MLflow 2.20+、Evidently 0.6+、Docker 25+


摘要:本章将构建一个完整的MLOps流水线,涵盖特征存储、模型训练、版本管理、API部署、容器化和监控告警。通过实战掌握2025年最前沿的机器学习工程化方法。


学习目标

完成本章学习后,你将能够:

能力维度 具体技能
架构设计 设计端到端MLOps架构,理解各组件职责边界
特征工程 使用 Feast 构建特征存储,实现训练与推理特征一致性
模型管理 使用 MLflow 进行实验追踪、模型注册和版本控制
服务部署 使用 FastAPI 构建高性能推理服务,集成异步处理
容器化 编写生产级 Dockerfile,实现多阶段构建和镜像优化
监控运维 使用 Evidently 实现数据漂移检测和模型性能监控
自动化 构建 CI/CD 流水线,实现模型自动训练和部署

本章内容概览

本章内容涵盖MLOps全流程,分为以下模块:

  • 特征存储(Feast)
  • 模型管理(MLflow)
  • API部署(FastAPI)
  • 容器化(Docker)
  • 监控告警(Evidently)

建议按顺序学习,每个模块可独立实践。


1. MLOps 架构全景

| :— | :— |
| 架构设计 | 设计端到端MLOps架构,理解各组件职责边界 |
| 特征工程 | 使用 Feast 构建特征存储,实现训练与推理特征一致性 |
| 模型管理 | 使用 MLflow 进行实验追踪、模型注册和版本控制 |
| 服务部署 | 使用 FastAPI 构建高性能推理服务,集成异步处理 |
| 容器化 | 编写生产级 Dockerfile,实现多阶段构建和镜像优化 |
| 监控运维 | 使用 Evidently 实现数据漂移检测和模型性能监控 |
| 自动化 | 构建 CI/CD 流水线,实现模型自动训练和部署 |


1. MLOps 架构全景

1.1 什么是 MLOps

MLOps(Machine Learning Operations)是机器学习开发与运维的工程化实践,旨在通过标准化流程和工具链实现模型的高效开发、部署、监控与迭代。

用一个比喻理解:传统机器学习项目像手工作坊,依赖个人经验;MLOps 则是现代化工厂,通过流水线、质检系统和自动化设备确保产品质量稳定可控。

1.2 2025年 MLOps 技术栈

┌─────────────────────────────────────────────────────────────────┐
│                        MLOps 技术架构                            │
├─────────────────────────────────────────────────────────────────┤
│  应用层  │  Streamlit Dashboard  │  REST API  │  Batch Inference │
├─────────┼─────────────────────────┼────────────┼─────────────────┤
│  服务层  │  FastAPI Service  │  Feature Store  │  Model Registry │
├─────────┼─────────────────────────────────────────────────────────┤
│  监控层  │  Evidently (漂移检测)  │  Prometheus  │  Grafana       │
├─────────┼─────────────────────────────────────────────────────────┤
│  编排层  │  MLflow (实验追踪)  │  Airflow/Prefect (工作流)       │
├─────────┼─────────────────────────────────────────────────────────┤
│  基础设施 │  Docker  │  Kubernetes  │  S3/MinIO (对象存储)          │
└─────────────────────────────────────────────────────────────────┘

1.3 核心组件选型对比

组件类型 开源方案 商业方案 适用场景
实验追踪 MLflow, WandB Weights & Biases 中小团队首选 MLflow
特征存储 Feast, Hopsworks Tecton, SageMaker 实时场景用 Feast
模型监控 Evidently, WhyLabs Arize, Fiddler 快速启动用 Evidently
工作流编排 Airflow, Prefect AWS Step Functions 复杂依赖用 Airflow
模型服务 FastAPI, BentoML SageMaker, Vertex AI 高性能用 FastAPI

2. 项目架构设计

2.1 业务场景:电商用户购买预测

项目背景:某电商平台需要实时预测用户的购买意向,以便进行精准营销和个性化推荐。

核心挑战

  • 特征工程复杂:需整合用户行为、商品属性、上下文信息
  • 实时性要求:预测延迟需控制在 100ms 以内
  • 数据漂移:用户行为随季节、促销活动变化

技术目标

  • 模型 AUC >= 0.85
  • P99 延迟 < 100ms
  • 支持每日自动重训练

2.2 系统架构图

                    ┌─────────────────┐
                    │   用户请求       │
                    └────────┬────────┘
                             │
                    ┌────────▼────────┐
                    │  FastAPI 服务    │
                    │  (模型推理)      │
                    └────────┬────────┘
                             │
           ┌─────────────────┼─────────────────┐
           │                 │                 │
    ┌──────▼──────┐  ┌──────▼──────┐  ┌──────▼──────┐
    │ Feature Store│  │ Model Registry│  │ 监控服务    │
    │   (Feast)   │  │   (MLflow)   │  │ (Evidently)│
    └─────────────┘  └─────────────┘  └─────────────┘
           │                 │                 │
           └─────────────────┼─────────────────┘
                             │
                    ┌────────▼────────┐
                    │   对象存储       │
                    │  (MinIO/S3)     │
                    └─────────────────┘

2.3 项目目录结构

mlops_project/
├── app/                          # FastAPI 应用
│   ├── __init__.py
│   ├── main.py                   # 服务入口
│   ├── api/
│   │   ├── __init__.py
│   │   ├── predict.py            # 预测接口
│   │   └── health.py             # 健康检查
│   ├── core/
│   │   ├── __init__.py
│   │   ├── config.py             # 配置管理
│   │   └── logging.py            # 日志配置
│   └── services/
│       ├── __init__.py
│       ├── model_service.py      # 模型加载与推理
│       └── feature_service.py    # 特征获取
├── feature_repo/                 # Feast 特征仓库
│   ├── feature_store.yaml
│   ├── features/
│   │   ├── user_features.py
│   │   └── item_features.py
│   └── data/
│       └── sample_data.parquet
├── training/                     # 模型训练
│   ├── train.py                  # 训练脚本
│   ├── evaluate.py               # 评估脚本
│   └── pipeline.py               # 训练流水线
├── monitoring/                   # 监控配置
│   ├── drift_detection.py
│   └── dashboards/
├── deployment/                   # 部署配置
│   ├── Dockerfile
│   ├── docker-compose.yml
│   └── k8s/
├── tests/                        # 测试代码
│   ├── unit/
│   └── integration/
├── requirements.txt
└── README.md

3. 特征存储:Feast 实战

3.1 为什么需要特征存储

传统方式的痛点

问题 描述 后果
训练-推理不一致 训练时离线计算特征,推理时实时计算逻辑不同 模型性能下降
特征重复开发 不同团队重复开发相同特征 资源浪费
特征版本混乱 特征定义变更无版本控制 结果不可复现

特征存储的价值

  • 特征定义统一:一次定义,多处使用
  • 在线/离线一致性:同一特征定义服务训练和推理
  • 特征版本管理:支持特征回滚和 A/B 测试

3.2 Feast 环境配置

# 安装 Feast
pip install feast==0.40.0

# 初始化特征仓库
feast init feature_repo
cd feature_repo

3.3 特征定义

# feature_repo/features/user_features.py
from feast import Entity, Feature, FeatureView, ValueType
from feast.types import Float32, Int64, String
from datetime import timedelta

# 定义实体
user = Entity(
    name="user_id",
    value_type=ValueType.INT64,
    description="用户唯一标识"
)

# 定义特征视图
user_feature_view = FeatureView(
    name="user_features",
    entities=["user_id"],
    ttl=timedelta(days=1),
    schema=[
        Field(name="age", dtype=Int64),
        Field(name="gender", dtype=String),
        Field(name="register_days", dtype=Int64),
        Field(name="avg_order_value", dtype=Float32),
        Field(name="total_orders", dtype=Int64),
        Field(name="last_login_days", dtype=Int64),
        Field(name="favorite_category", dtype=String),
    ],
    online=True,
    source=user_stats_source,
    tags={"team": "recommendation"},
)
# feature_repo/features/item_features.py
from feast import Entity, FeatureView, Field
from feast.types import Float32, Int64, String
from datetime import timedelta

item = Entity(
    name="item_id",
    value_type=ValueType.INT64,
    description="商品唯一标识"
)

item_feature_view = FeatureView(
    name="item_features",
    entities=["item_id"],
    ttl=timedelta(hours=6),
    schema=[
        Field(name="category_id", dtype=Int64),
        Field(name="brand_id", dtype=Int64),
        Field(name="price", dtype=Float32),
        Field(name="avg_rating", dtype=Float32),
        Field(name="review_count", dtype=Int64),
        Field(name="stock_quantity", dtype=Int64),
        Field(name="is_new_arrival", dtype=Int64),
    ],
    online=True,
    source=item_stats_source,
)

3.4 特征存储配置

# feature_repo/feature_store.yaml
project: ecommerce_prediction
provider: local
registry: data/registry.db
online_store:
    type: sqlite
    path: data/online_store.db
offline_store:
    type: file
entity_key_serialization_version: 2

3.5 特征获取实战

# 训练时获取历史特征
from feast import FeatureStore
import pandas as pd

store = FeatureStore(repo_path="feature_repo")

# 构建训练样本
train_df = pd.read_parquet("data/train_events.parquet")

# 获取特征
feature_refs = [
    "user_features:age",
    "user_features:avg_order_value",
    "user_features:total_orders",
    "item_features:price",
    "item_features:avg_rating",
]

training_df = store.get_historical_features(
    entity_df=train_df[["user_id", "item_id", "event_timestamp"]],
    features=feature_refs,
).to_df()

print(f"训练数据形状: {training_df.shape}")
print(training_df.head())
# 推理时获取实时特征
from feast import FeatureStore

store = FeatureStore(repo_path="feature_repo")

# 在线获取特征
features = store.get_online_features(
    features=[
        "user_features:age",
        "user_features:avg_order_value",
        "item_features:price",
    ],
    entity_rows=[
        {"user_id": 12345, "item_id": 67890}
    ],
).to_dict()

print(features)

4. 模型训练与版本管理

4.1 MLflow 实验追踪

# training/train.py
import mlflow
import mlflow.sklearn
from sklearn.ensemble import GradientBoostingClassifier
from sklearn.model_selection import train_test_split
from sklearn.metrics import roc_auc_score, precision_recall_curve
import pandas as pd
import joblib

def train_model(data_path: str):
    """训练模型并记录到 MLflow"""
  
    # 设置 MLflow 跟踪服务器
    mlflow.set_tracking_uri("http://localhost:5000")
    mlflow.set_experiment("ecommerce_purchase_prediction")
  
    # 加载数据
    df = pd.read_parquet(data_path)
    X = df.drop(["user_id", "item_id", "event_timestamp", "label"], axis=1)
    y = df["label"]
  
    X_train, X_test, y_train, y_test = train_test_split(
        X, y, test_size=0.2, random_state=42, stratify=y
    )
  
    # 开始实验
    with mlflow.start_run(run_name="gbdt_v1"):
        # 记录参数
        params = {
            "n_estimators": 200,
            "max_depth": 6,
            "learning_rate": 0.1,
            "subsample": 0.8,
            "random_state": 42
        }
        mlflow.log_params(params)
      
        # 训练模型
        model = GradientBoostingClassifier(**params)
        model.fit(X_train, y_train)
      
        # 评估
        y_pred = model.predict_proba(X_test)[:, 1]
        auc = roc_auc_score(y_test, y_pred)
      
        # 记录指标
        mlflow.log_metric("auc", auc)
        mlflow.log_metric("train_samples", len(X_train))
        mlflow.log_metric("test_samples", len(X_test))
      
        # 记录特征重要性
        feature_importance = pd.DataFrame({
            "feature": X.columns,
            "importance": model.feature_importances_
        }).sort_values("importance", ascending=False)
      
        mlflow.log_table("feature_importance", feature_importance)
      
        # 保存模型
        mlflow.sklearn.log_model(
            model,
            artifact_path="model",
            registered_model_name="purchase_prediction_model"
        )
      
        print(f"模型训练完成,AUC: {auc:.4f}")
        return model, auc

if __name__ == "__main__":
    train_model("data/training_data.parquet")

4.2 模型注册与版本管理

# training/model_registry.py
from mlflow.tracking import MlflowClient

def promote_model_to_production(model_name: str, version: int):
    """将模型提升到生产环境"""
    client = MlflowClient()
  
    # 将指定版本提升到 Production
    client.transition_model_version_stage(
        name=model_name,
        version=version,
        stage="Production",
        archive_existing_versions=True
    )
  
    print(f"模型 {model_name} 版本 {version} 已部署到生产环境")

def get_production_model_uri(model_name: str) -> str:
    """获取生产环境模型路径"""
    client = MlflowClient()
  
    # 获取 Production 阶段的最新版本
    versions = client.get_latest_versions(model_name, stages=["Production"])
  
    if not versions:
        raise ValueError(f"没有找到生产环境的模型: {model_name}")
  
    latest_version = versions[0]
    return f"models:/{model_name}/{latest_version.version}"

# 使用示例
if __name__ == "__main__":
    promote_model_to_production("purchase_prediction_model", version=3)
    model_uri = get_production_model_uri("purchase_prediction_model")
    print(f"生产模型 URI: {model_uri}")

4.3 自动化训练流水线

# training/pipeline.py
import prefect
from prefect import flow, task
from prefect.task_runners import ConcurrentTaskRunner

@task
def extract_features():
    """特征提取"""
    print("从特征存储提取特征...")
    # 实际实现...
    return "data/training_data.parquet"

@task
def train_model_task(data_path: str):
    """训练模型"""
    from train import train_model
    model, auc = train_model(data_path)
    return auc

@task
def evaluate_model(auc: float) -> bool:
    """模型评估"""
    threshold = 0.85
    if auc >= threshold:
        print(f"模型通过评估,AUC: {auc:.4f}")
        return True
    else:
        print(f"模型未通过评估,AUC: {auc:.4f} < {threshold}")
        return False

@task
def deploy_model():
    """部署模型"""
    print("部署模型到生产环境...")
    # 实际部署逻辑

@flow(name="ml_training_pipeline", task_runner=ConcurrentTaskRunner())
def training_pipeline():
    """完整的训练流水线"""
  
    # 提取特征
    data_path = extract_features()
  
    # 训练模型
    auc = train_model_task(data_path)
  
    # 评估模型
    passed = evaluate_model(auc)
  
    # 条件部署
    if passed:
        deploy_model()
    else:
        print("模型未通过评估,跳过部署")

if __name__ == "__main__":
    training_pipeline()

5. FastAPI 高性能推理服务

5.1 核心服务架构

# app/main.py
from fastapi import FastAPI, HTTPException
from fastapi.middleware.cors import CORSMiddleware
from contextlib import asynccontextmanager
import uvicorn

from app.api import predict, health
from app.core.config import settings
from app.services.model_service import ModelService

@asynccontextmanager
async def lifespan(app: FastAPI):
    """应用生命周期管理"""
    # 启动时加载模型
    app.state.model_service = ModelService()
    await app.state.model_service.load_model()
    print("模型加载完成")
  
    yield
  
    # 关闭时清理资源
    await app.state.model_service.cleanup()
    print("资源清理完成")

app = FastAPI(
    title="电商购买预测服务",
    description="基于 MLOps 的实时购买意向预测 API",
    version="1.0.0",
    lifespan=lifespan
)

# CORS 配置
app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

# 注册路由
app.include_router(predict.router, prefix="/api/v1", tags=["prediction"])
app.include_router(health.router, prefix="/api/v1", tags=["health"])

if __name__ == "__main__":
    uvicorn.run(
        "app.main:app",
        host="0.0.0.0",
        port=8000,
        reload=False,
        workers=1
    )

5.2 预测接口实现

# app/api/predict.py
from fastapi import APIRouter, HTTPException, Depends, BackgroundTasks
from pydantic import BaseModel, Field
from typing import List, Optional
import numpy as np
import time

from app.services.model_service import ModelService
from app.services.feature_service import FeatureService
from app.core.logging import logger

router = APIRouter()

class PredictionRequest(BaseModel):
    """预测请求模型"""
    user_id: int = Field(..., description="用户ID")
    item_id: int = Field(..., description="商品ID")
    context: Optional[dict] = Field(None, description="上下文信息")

class PredictionResponse(BaseModel):
    """预测响应模型"""
    user_id: int
    item_id: int
    purchase_probability: float = Field(..., ge=0, le=1)
    risk_level: str
    confidence: str
    latency_ms: float
    model_version: str

class BatchPredictionRequest(BaseModel):
    """批量预测请求"""
    items: List[PredictionRequest] = Field(..., max_length=100)

def get_model_service():
    from app.main import app
    return app.state.model_service

def get_feature_service():
    return FeatureService()

@router.post("/predict", response_model=PredictionResponse)
async def predict(
    request: PredictionRequest,
    background_tasks: BackgroundTasks,
    model_service: ModelService = Depends(get_model_service),
    feature_service: FeatureService = Depends(get_feature_service)
):
    """
    单条预测接口
  
    预测用户购买某商品的意向概率
    """
    start_time = time.time()
  
    try:
        # 获取特征
        features = await feature_service.get_features(
            request.user_id, 
            request.item_id
        )
      
        # 模型推理
        probability = await model_service.predict(features)
      
        # 确定风险等级
        if probability < 0.3:
            risk_level = "低意向"
            confidence = "高"
        elif probability < 0.7:
            risk_level = "中意向"
            confidence = "中"
        else:
            risk_level = "高意向"
            confidence = "高"
      
        latency = (time.time() - start_time) * 1000
      
        # 异步记录预测日志
        background_tasks.add_task(
            log_prediction,
            request.user_id,
            request.item_id,
            probability,
            latency
        )
      
        return PredictionResponse(
            user_id=request.user_id,
            item_id=request.item_id,
            purchase_probability=probability,
            risk_level=risk_level,
            confidence=confidence,
            latency_ms=round(latency, 2),
            model_version=model_service.current_version
        )
      
    except Exception as e:
        logger.error(f"预测失败: {str(e)}")
        raise HTTPException(status_code=500, detail=f"预测失败: {str(e)}")

@router.post("/predict/batch")
async def predict_batch(
    request: BatchPredictionRequest,
    model_service: ModelService = Depends(get_model_service),
    feature_service: FeatureService = Depends(get_feature_service)
):
    """
    批量预测接口
  
    支持一次预测最多100条记录
    """
    results = []
  
    for item in request.items:
        features = await feature_service.get_features(
            item.user_id, 
            item.item_id
        )
        probability = await model_service.predict(features)
      
        results.append({
            "user_id": item.user_id,
            "item_id": item.item_id,
            "purchase_probability": probability
        })
  
    return {"predictions": results, "count": len(results)}

def log_prediction(user_id: int, item_id: int, probability: float, latency: float):
    """记录预测日志"""
    logger.info(
        f"prediction",
        extra={
            "user_id": user_id,
            "item_id": item_id,
            "probability": probability,
            "latency_ms": latency
        }
    )

5.3 模型服务封装

# app/services/model_service.py
import mlflow
import numpy as np
import joblib
from typing import Any, Dict
import asyncio
from cachetools import TTLCache

class ModelService:
    """模型服务封装"""
  
    def __init__(self):
        self.model = None
        self.current_version = None
        self.feature_order = None
        self.cache = TTLCache(maxsize=10000, ttl=300)  # 5分钟缓存
      
    async def load_model(self, model_uri: str = None):
        """加载模型"""
        if model_uri is None:
            # 从 MLflow 获取生产环境模型
            model_uri = "models:/purchase_prediction_model/Production"
      
        # 在异步线程中加载模型
        loop = asyncio.get_event_loop()
        self.model = await loop.run_in_executor(
            None, 
            lambda: mlflow.sklearn.load_model(model_uri)
        )
      
        # 加载特征顺序
        self.feature_order = joblib.load("models/feature_order.pkl")
      
        # 解析版本
        self.current_version = model_uri.split("/")[-1]
      
        print(f"模型加载成功: {model_uri}")
  
    async def predict(self, features: Dict[str, Any]) -> float:
        """异步预测"""
        # 检查缓存
        cache_key = f"{features.get('user_id')}_{features.get('item_id')}"
        if cache_key in self.cache:
            return self.cache[cache_key]
      
        # 特征排序
        feature_vector = [features.get(f, 0) for f in self.feature_order]
        X = np.array([feature_vector])
      
        # 异步推理
        loop = asyncio.get_event_loop()
        probability = await loop.run_in_executor(
            None,
            lambda: self.model.predict_proba(X)[0, 1]
        )
      
        # 更新缓存
        self.cache[cache_key] = probability
      
        return float(probability)
  
    async def cleanup(self):
        """清理资源"""
        self.model = None
        self.cache.clear()

5.4 特征服务封装

# app/services/feature_service.py
from feast import FeatureStore
import pandas as pd
from typing import Dict, Any

class FeatureService:
    """特征服务封装"""
  
    def __init__(self):
        self.store = FeatureStore(repo_path="feature_repo")
      
    async def get_features(self, user_id: int, item_id: int) -> Dict[str, Any]:
        """获取在线特征"""
        features = self.store.get_online_features(
            features=[
                "user_features:age",
                "user_features:avg_order_value",
                "user_features:total_orders",
                "user_features:last_login_days",
                "item_features:price",
                "item_features:avg_rating",
                "item_features:category_id",
            ],
            entity_rows=[{"user_id": user_id, "item_id": item_id}],
        ).to_dict()
      
        # 转换为字典
        result = {"user_id": user_id, "item_id": item_id}
        for key, values in features.items():
            if values:
                result[key] = values[0]
      
        return result

6. Docker 容器化部署

6.1 生产级 Dockerfile

# deployment/Dockerfile
# 多阶段构建优化镜像体积

# 阶段1:构建依赖
FROM python:3.12-slim as builder

WORKDIR /build

# 安装编译依赖
RUN apt-get update && apt-get install -y --no-install-recommends \
    gcc \
    g++ \
    && rm -rf /var/lib/apt/lists/*

# 复制依赖文件
COPY requirements.txt .

# 安装依赖到虚拟环境
RUN python -m venv /opt/venv
ENV PATH="/opt/venv/bin:$PATH"
RUN pip install --no-cache-dir --upgrade pip && \
    pip install --no-cache-dir -r requirements.txt

# 阶段2:运行环境
FROM python:3.12-slim

WORKDIR /app

# 创建非 root 用户
RUN groupadd -r appuser && useradd -r -g appuser appuser

# 从构建阶段复制虚拟环境
COPY --from=builder /opt/venv /opt/venv
ENV PATH="/opt/venv/bin:$PATH"

# 复制应用代码
COPY app/ ./app/
COPY feature_repo/ ./feature_repo/
COPY models/ ./models/

# 设置权限
RUN chown -R appuser:appuser /app
USER appuser

# 健康检查
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
    CMD python -c "import requests; requests.get('http://localhost:8000/api/v1/health')" || exit 1

# 暴露端口
EXPOSE 8000

# 启动命令
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "2"]

6.2 Docker Compose 编排

# deployment/docker-compose.yml
version: '3.8'

services:
  # MLflow 跟踪服务器
  mlflow:
    image: python:3.12-slim
    command: >
      bash -c "pip install mlflow==2.20.0 && 
               mlflow server 
               --host 0.0.0.0 
               --port 5000 
               --backend-store-uri sqlite:///mlflow/mlflow.db 
               --default-artifact-root /mlflow/artifacts"
    ports:
      - "5000:5000"
    volumes:
      - mlflow_data:/mlflow
    networks:
      - mlops_network

  # MinIO 对象存储
  minio:
    image: minio/minio:latest
    command: server /data --console-address ":9001"
    ports:
      - "9000:9000"
      - "9001:9001"
    environment:
      MINIO_ROOT_USER: minioadmin
      MINIO_ROOT_PASSWORD: minioadmin
    volumes:
      - minio_data:/data
    networks:
      - mlops_network

  # 模型推理服务
  prediction-api:
    build:
      context: ..
      dockerfile: deployment/Dockerfile
    ports:
      - "8000:8000"
    environment:
      - MLFLOW_TRACKING_URI=http://mlflow:5000
      - FEATURE_STORE_PATH=./feature_repo
      - MODEL_NAME=purchase_prediction_model
    depends_on:
      - mlflow
      - minio
    networks:
      - mlops_network
    restart: unless-stopped

  # 监控服务
  monitoring:
    build:
      context: ..
      dockerfile: deployment/Dockerfile.monitoring
    ports:
      - "8001:8001"
    environment:
      - MLFLOW_TRACKING_URI=http://mlflow:5000
      - PREDICTION_API_URL=http://prediction-api:8000
    depends_on:
      - prediction-api
    networks:
      - mlops_network

volumes:
  mlflow_data:
  minio_data:

networks:
  mlops_network:
    driver: bridge

6.3 构建与运行

# 构建镜像
docker-compose -f deployment/docker-compose.yml build

# 启动服务
docker-compose -f deployment/docker-compose.yml up -d

# 查看日志
docker-compose -f deployment/docker-compose.yml logs -f prediction-api

# 停止服务
docker-compose -f deployment/docker-compose.yml down

7. 模型监控与漂移检测

7.1 监控指标体系

监控维度 指标 告警阈值 处理方式
数据质量 缺失值比例 > 5% 暂停服务,人工介入
数据漂移 PSI 指数 > 0.2 触发重训练
概念漂移 AUC 下降 > 0.05 触发重训练
服务性能 P99 延迟 > 100ms 扩容或优化
服务可用性 错误率 > 1% 回滚到上一版本

7.2 Evidently 漂移检测

# monitoring/drift_detection.py
import evidently
from evidently.report import Report
from evidently.metric_preset import DataDriftPreset, TargetDriftPreset
from evidently.metrics import ColumnDriftMetric, DatasetDriftMetric
import pandas as pd
import requests
from datetime import datetime, timedelta

class DriftDetector:
    """漂移检测器"""
  
    def __init__(self, reference_data_path: str):
        self.reference_data = pd.read_parquet(reference_data_path)
        self.drift_report = None
      
    def generate_report(self, current_data: pd.DataFrame) -> dict:
        """生成漂移检测报告"""
      
        report = Report(metrics=[
            DatasetDriftMetric(),
            ColumnDriftMetric(column_name="age"),
            ColumnDriftMetric(column_name="avg_order_value"),
            ColumnDriftMetric(column_name="price"),
        ])
      
        report.run(
            reference_data=self.reference_data,
            current_data=current_data
        )
      
        self.drift_report = report
      
        # 提取关键指标
        results = report.as_dict()
        dataset_drift = results["metrics"][0]["result"]
      
        return {
            "dataset_drift": dataset_drift["dataset_drift"],
            "drift_share": dataset_drift["share_of_drifted_columns"],
            "number_of_columns": dataset_drift["number_of_columns"],
            "number_of_drifted_columns": dataset_drift["number_of_drifted_columns"],
            "timestamp": datetime.now().isoformat()
        }
  
    def check_drift_threshold(self, report: dict, threshold: float = 0.2) -> bool:
        """检查是否超过漂移阈值"""
        return report["drift_share"] > threshold
  
    def save_report(self, output_path: str):
        """保存 HTML 报告"""
        if self.drift_report:
            self.drift_report.save_html(output_path)

# 实时监控脚本
import schedule
import time

def monitor_job():
    """定时监控任务"""
    detector = DriftDetector("data/reference_data.parquet")
  
    # 获取最近1小时的数据
    current_data = pd.read_parquet("data/recent_predictions.parquet")
  
    # 生成报告
    report = detector.generate_report(current_data)
  
    print(f"漂移检测完成: {report}")
  
    # 检查是否需要告警
    if detector.check_drift_threshold(report):
        print("警告: 检测到数据漂移!")
        # 发送告警通知
        send_alert(report)
        # 触发重训练
        trigger_retraining()

def send_alert(report: dict):
    """发送告警通知"""
    # 集成企业微信/钉钉/邮件等
    webhook_url = "https://oapi.dingtalk.com/robot/send?access_token=xxx"
    message = {
        "msgtype": "text",
        "text": {
            "content": f"模型漂移告警: 漂移比例 {report['drift_share']:.2%}"
        }
    }
    requests.post(webhook_url, json=message)

def trigger_retraining():
    """触发模型重训练"""
    # 调用训练流水线
    requests.post("http://training-service:8000/trigger")

# 每小时执行一次
schedule.every().hour.do(monitor_job)

if __name__ == "__main__":
    while True:
        schedule.run_pending()
        time.sleep(60)

7.3 性能监控面板

# monitoring/metrics_exporter.py
from prometheus_client import Counter, Histogram, Gauge, start_http_server
import time

# 定义指标
PREDICTION_COUNTER = Counter(
    'model_predictions_total',
    'Total number of predictions',
    ['model_version', 'status']
)

PREDICTION_LATENCY = Histogram(
    'model_prediction_latency_seconds',
    'Prediction latency in seconds',
    buckets=[0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0]
)

MODEL_AUC = Gauge(
    'model_auc_score',
    'Current model AUC score'
)

DRIFT_SCORE = Gauge(
    'data_drift_score',
    'Data drift detection score'
)

class MetricsCollector:
    """指标收集器"""
  
    @staticmethod
    def record_prediction(latency: float, success: bool, version: str):
        """记录预测指标"""
        status = "success" if success else "failure"
        PREDICTION_COUNTER.labels(model_version=version, status=status).inc()
        PREDICTION_LATENCY.observe(latency)
  
    @staticmethod
    def update_model_auc(auc: float):
        """更新模型 AUC"""
        MODEL_AUC.set(auc)
  
    @staticmethod
    def update_drift_score(score: float):
        """更新漂移分数"""
        DRIFT_SCORE.set(score)

# 启动指标服务器
if __name__ == "__main__":
    start_http_server(8001)
    print("Metrics server started on port 8001")

8. CI/CD 自动化流水线

8.1 GitHub Actions 工作流

# .github/workflows/mlops-pipeline.yml
name: MLOps Pipeline

on:
  push:
    branches: [ main ]
  schedule:
    # 每天凌晨2点执行
    - cron: '0 2 * * *'
  workflow_dispatch:

jobs:
  test:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
    
      - name: Set up Python
        uses: actions/setup-python@v5
        with:
          python-version: '3.12'
    
      - name: Install dependencies
        run: |
          pip install -r requirements.txt
          pip install pytest pytest-cov
    
      - name: Run tests
        run: pytest tests/ --cov=app --cov-report=xml
    
      - name: Upload coverage
        uses: codecov/codecov-action@v3
        with:
          file: ./coverage.xml

  train:
    needs: test
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
    
      - name: Set up Python
        uses: actions/setup-python@v5
        with:
          python-version: '3.12'
    
      - name: Install dependencies
        run: pip install -r requirements.txt
    
      - name: Start MLflow
        run: |
          mlflow server --host 0.0.0.0 --port 5000 &
          sleep 5
    
      - name: Train model
        run: python training/train.py
        env:
          MLFLOW_TRACKING_URI: http://localhost:5000
    
      - name: Evaluate model
        run: python training/evaluate.py

  build:
    needs: train
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
    
      - name: Set up Docker Buildx
        uses: docker/setup-buildx-action@v3
    
      - name: Login to Docker Hub
        uses: docker/login-action@v3
        with:
          username: ${{ secrets.DOCKER_USERNAME }}
          password: ${{ secrets.DOCKER_PASSWORD }}
    
      - name: Build and push
        uses: docker/build-push-action@v5
        with:
          context: .
          file: ./deployment/Dockerfile
          push: true
          tags: |
            ${{ secrets.DOCKER_USERNAME }}/mlops-prediction-api:${{ github.sha }}
            ${{ secrets.DOCKER_USERNAME }}/mlops-prediction-api:latest

  deploy:
    needs: build
    runs-on: ubuntu-latest
    if: github.ref == 'refs/heads/main'
    steps:
      - name: Deploy to production
        run: |
          echo "部署到生产环境"
          # 实际部署命令

9. 避坑小贴士

9.1 特征存储陷阱

陷阱 描述 解决方案
在线/离线不一致 训练时批处理计算特征,推理时实时计算逻辑不同 统一使用 Feast 特征视图定义
TTL 设置不当 特征过期时间过短导致获取失败,过长导致数据陈旧 根据业务场景设置合理 TTL
实体缺失 新用户/新商品无历史特征 设计默认特征值或冷启动策略

9.2 模型服务陷阱

陷阱 描述 解决方案
模型热加载 更新模型时服务中断 使用蓝绿部署或金丝雀发布
内存泄漏 频繁加载模型导致内存持续增长 限制缓存大小,定期清理
并发瓶颈 同步推理导致请求排队 使用异步处理或模型批量化推理

9.3 监控陷阱

陷阱 描述 解决方案
延迟标签 真实标签获取延迟导致无法及时评估 使用代理指标或无监督漂移检测
阈值僵化 固定阈值无法适应业务变化 使用动态阈值或基于历史分布的阈值
告警疲劳 过多无效告警导致忽视真正问题 分级告警,合并相似告警

10. 工程化思维总结

10.1 MLOps 成熟度模型

级别 特征 关键能力
L1 手动 脚本化训练,手动部署 版本控制、基础测试
L2 自动化 CI/CD 流水线,自动部署 自动化测试、容器化
L3 监控化 实时监控,自动告警 漂移检测、性能监控
L4 智能化 自动重训练,A/B 测试 实验平台、流量切分
L5 全自动化 端到端自动化,自愈系统 全自动闭环、智能决策

10.2 关键设计原则

  1. 可复现性:所有实验必须可复现,记录完整的参数、代码和数据版本
  2. 可观测性:系统状态必须可见,包括模型性能、数据分布、服务指标
  3. 可回滚:任何变更都能快速回滚,保留历史版本
  4. 松耦合:组件间通过标准接口通信,便于独立迭代
  5. 自动化:重复性工作必须自动化,减少人为错误

本章小结

本章构建了一个完整的 MLOps 实战项目,涵盖以下核心内容:

模块 技术要点
特征存储 Feast 实现训练/推理特征一致性
模型管理 MLflow 实验追踪与版本控制
服务部署 FastAPI 高性能异步推理服务
容器化 Docker 多阶段构建与 Compose 编排
监控告警 Evidently 漂移检测 + Prometheus 指标采集
自动化 GitHub Actions CI/CD 流水线

核心收获

  • MLOps 不是工具堆砌,而是工程化思维的实践
  • 特征存储是解决训练-推理不一致的关键
  • 监控是模型可持续运行的保障
  • 自动化是提升效率的必由之路

上一篇/下一篇导航

  • 上一篇:[第15章 PySpark大数据处理](file:///d:/python/Python全栈开发/数据科学实战之路/第15章-PySpark大数据处理/第15章-PySpark大数据处理.md)
  • 下一篇:[第17章 Dask并行计算](file:///d:/python/Python全栈开发/数据科学实战之路/第17章-Dask并行计算/第17章-Dask并行计算.md)

本章内容到此结束。MLOps 是一个快速发展的领域,建议持续关注最新工具和实践,不断优化你的机器学习工程化能力。

更多推荐