Machine Learning: 项目部署 — SalesPredict上线与运维指南

部署是终点也是起点——上线意味着运维开始,持续可靠才是真正的交付。

1. 你将学到


2. 一个DevOps工程师的真实故事

(1) 痛点:模型训练完了但部署方案不完整

Bob训练好了MAPE 8%的LightGBM模型,但部署方案只有FastAPI + Docker——没有监控、没有缓存、没有限流、没有重训练机制。上线第一天就因为流量突增导致API超时,第二天因为特征格式变更导致预测全错。没有运维的部署就像没有刹车的汽车——迟早出事。

(2) 完整生产部署的解法

生产级部署=服务(API) + 编排(Docker) + 监控(Grafana) + 告警(Prometheus) + 重训练(Airflow)。

YAML
# Production-grade deployment stack
services:
  api: FastAPI + Uvicorn (model serving)
  redis: Cache layer (hot predictions)
  nginx: Load balancer + rate limiting
  prometheus: Metrics collection
  grafana: Monitoring dashboard
  airflow: Retraining scheduler

(3) 收益:上线3个月零故障,MAPE稳定8%

Bob完成全链路部署后,3个月零故障运行,模型MAPE稳定在8%,一次双十一漂移在3天内自动检测并重训练修复。


3. FastAPI推理服务

(1) 生产级API实现

▶ 示例:完整FastAPI服务

PYTHON
# File: app/main.py
from fastapi import FastAPI, HTTPException, Request
from fastapi.responses import JSONResponse
from pydantic import BaseModel, Field
from typing import Optional
import joblib
import numpy as np
import time
import logging

logger = logging.getLogger("salespredict")
app = FastAPI(title="SalesPredict API", version="1.0.0")

# Global model (loaded at startup)
model = None

@app.on_event("startup")
async def load_model():
    global model
    model = joblib.load("model/salespredict_lgbm.joblib")
    logger.info("Model loaded successfully")

class PredictionInput(BaseModel):
    ad_spend_k_usd: float = Field(..., ge=0, le=1000)
    traffic_k: float = Field(..., ge=0, le=10000)
    is_promotion: int = Field(0, ge=0, le=1)
    is_weekend: int = Field(0, ge=0, le=1)
    category_clothing: int = Field(0, ge=0, le=1)
    category_electronics: int = Field(0, ge=0, le=1)
    category_food: int = Field(0, ge=0, le=1)
    category_home: int = Field(0, ge=0, le=1)
    region_EU: int = Field(0, ge=0, le=1)
    region_US: int = Field(0, ge=0, le=1)

    model_config = {"json_schema_extra": {
        "example": {"ad_spend_k_usd": 50, "traffic_k": 300,
                     "is_promotion": 1, "is_weekend": 0,
                     "category_electronics": 1, "category_clothing": 0,
                     "category_food": 0, "category_home": 0,
                     "region_US": 1, "region_EU": 0}
    }}

class PredictionOutput(BaseModel):
    predicted_revenue_k_usd: float
    confidence_low: Optional[float] = None
    confidence_high: Optional[float] = None
    latency_ms: float

@app.get("/health")
def health():
    return {"status": "healthy", "model_loaded": model is not None}

@app.post("/predict", response_model=PredictionOutput)
def predict(input_data: PredictionInput):
    start = time.time()
    try:
        features = np.array([[input_data.ad_spend_k_usd, input_data.traffic_k,
                               input_data.is_promotion, input_data.is_weekend,
                               input_data.category_clothing, input_data.category_electronics,
                               input_data.category_food, input_data.category_home,
                               input_data.region_EU, input_data.region_US]])
        prediction = float(model.predict(features)[0])
        latency = (time.time() - start) * 1000

        # Simple confidence interval (±20%)
        return PredictionOutput(
            predicted_revenue_k_usd=round(prediction, 2),
            confidence_low=round(prediction * 0.8, 2),
            confidence_high=round(prediction * 1.2, 2),
            latency_ms=round(latency, 2),
        )
    except Exception as e:
        logger.error(f"Prediction error: {e}")
        raise HTTPException(status_code=500, detail=str(e))

@app.post("/predict_batch")
def predict_batch(inputs: list[PredictionInput], max_batch: int = 100):
    if len(inputs) > max_batch:
        raise HTTPException(status_code=400, detail=f"Batch size exceeds {max_batch}")
    start = time.time()
    features = np.array([[d.ad_spend_k_usd, d.traffic_k, d.is_promotion, d.is_weekend,
                           d.category_clothing, d.category_electronics, d.category_food,
                           d.category_home, d.region_EU, d.region_US] for d in inputs])
    predictions = model.predict(features)
    latency = (time.time() - start) * 1000
    return {
        "predictions": [round(float(p), 2) for p in predictions],
        "count": len(predictions),
        "latency_ms": round(latency, 2),
    }

@app.middleware("http")
async def log_requests(request: Request, call_next):
    start = time.time()
    response = await call_next(request)
    latency = (time.time() - start) * 1000
    logger.info(f"{request.method} {request.url.path} - {response.status_code} - {latency:.1f}ms")
    return response

输出:

TEXT 📖 仅展示
INFO:     Application startup complete.
INFO:     Model loaded successfully
INFO:     Uvicorn running on http://0.0.0.0:8000

(2) API性能指标

端点 单次延迟 批量延迟(100) QPS
/predict < 10ms 100+
/predict_batch < 50ms 50+
/health < 1ms 1000+

4. Docker部署

(1) 多阶段Dockerfile

▶ 示例:生产级Dockerfile

DOCKERFILE
# Stage 1: Build
FROM python:3.11-slim AS builder
WORKDIR /build
COPY requirements.txt .
RUN pip install --no-cache-dir --prefix=/install -r requirements.txt

# Stage 2: Runtime
FROM python:3.11-slim
WORKDIR /app

COPY --from=builder /install /usr/local
COPY app/ ./app/
COPY model/ ./model/

RUN useradd -m -r appuser && chown -R appuser:appuser /app
USER appuser

EXPOSE 8000

HEALTHCHECK --interval=30s --timeout=5s --retries=3 \
  CMD python -c "import urllib.request; urllib.request.urlopen('http://localhost:8000/health')"

CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "4"]

输出:

TEXT 📖 仅展示
Successfully built 3a7f2b1c9d4e
Successfully tagged salespredict-api:latest

(2) docker-compose编排

▶ 示例:完整服务栈

YAML
# docker-compose.yml
version: "3.8"

services:
  api:
    build: .
    environment:
      - REDIS_URL=redis://redis:6379/0
      - MLFLOW_TRACKING_URI=http://mlflow:5000
    depends_on:
      - redis
    deploy:
      replicas: 2
      resources:
        limits:
          memory: 2G
    restart: unless-stopped
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
      interval: 30s
      timeout: 5s

  redis:
    image: redis:7-alpine
    command: redis-server --maxmemory 256mb --maxmemory-policy allkeys-lru
    volumes:
      - redis_data:/data
    restart: unless-stopped

  nginx:
    image: nginx:alpine
    ports:
      - "80:80"
    volumes:
      - ./nginx.conf:/etc/nginx/conf.d/default.conf:ro
    depends_on:
      - api
    restart: unless-stopped

  prometheus:
    image: prom/prometheus:latest
    ports:
      - "9090:9090"
    volumes:
      - ./prometheus.yml:/etc/prometheus/prometheus.yml:ro
      - prometheus_data:/prometheus
    restart: unless-stopped

  grafana:
    image: grafana/grafana:latest
    ports:
      - "3000:3000"
    environment:
      - GF_SECURITY_ADMIN_PASSWORD=admin
    volumes:
      - grafana_data:/var/lib/grafana
    depends_on:
      - prometheus
    restart: unless-stopped

volumes:
  redis_data:
  prometheus_data:
  grafana_data:
TEXT 📖 仅展示
# nginx.conf
upstream api_backend {
    least_conn;
    server api:8000;
}

server {
    listen 80;
    location / {
        proxy_pass http://api_backend;
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
    }
    location /health {
        proxy_pass http://api_backend/health;
    }
}

5. 监控仪表盘

(1) Prometheus指标采集

▶ 示例:API指标暴露

PYTHON
# Add to app/main.py
from prometheus_client import Counter, Histogram, generate_latest
from fastapi import Response

PREDICTIONS_COUNT = Counter("predictions_total", "Total predictions made")
PREDICTION_LATENCY = Histogram("prediction_latency_seconds", "Prediction latency",
                                buckets=[0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0])

@app.get("/metrics")
def metrics():
    return Response(content=generate_latest(), media_type="text/plain")

# Instrument the predict endpoint
@app.post("/predict", response_model=PredictionOutput)
def predict(input_data: PredictionInput):
    start = time.time()
    PREDICTIONS_COUNT.inc()
    # ... (prediction logic) ...
    PREDICTION_LATENCY.observe(time.time() - start)
    # ... (return result) ...

输出:

TEXT 📖 仅展示
INFO:     Metrics endpoint registered at /metrics
INFO:     Prediction counter initialized

(2) Grafana仪表盘配置

面板 指标 告警阈值
预测QPS rate(predictions_total[5m]) < 10 (异常低)
延迟p95 histogram_quantile(0.95, prediction_latency_seconds) > 100ms
错误率 rate(http_requests_total{status=~"5xx"}[5m]) > 1%
预测均值 avg(predicted_revenue_k_usd) 偏离基线20%

▶ 示例:漂移监控脚本

PYTHON
# File: monitoring/drift_monitor.py
import numpy as np
import requests
import json
from datetime import datetime

class DriftMonitor:
    def __init__(self, reference_stats_path, api_url="http://localhost:8000"):
        self.reference = np.load(reference_stats_path, allow_pickle=True).item()
        self.api_url = api_url

    def calculate_psi(self, reference, current, n_bins=10):
        breakpoints = np.percentile(reference, np.linspace(0, 100, n_bins + 1))
        breakpoints[0], breakpoints[-1] = -np.inf, np.inf
        ref_counts = np.histogram(reference, bins=breakpoints)[0]
        cur_counts = np.histogram(current, bins=breakpoints)[0]
        ref_pct = np.clip(ref_counts / len(reference), 1e-6, None)
        cur_pct = np.clip(cur_counts / len(current), 1e-6, None)
        return np.sum((cur_pct - ref_pct) * np.log(cur_pct / ref_pct))

    def check_weekly_drift(self, current_features):
        """Run weekly drift check on all features."""
        results = {}
        for feature_name, current_data in current_features.items():
            if feature_name in self.reference:
                psi = self.calculate_psi(self.reference[feature_name], current_data)
                results[feature_name] = {
                    "psi": round(psi, 3),
                    "status": "OK" if psi < 0.1 else ("WARNING" if psi < 0.2 else "DRIFT"),
                }

        # Log results
        drift_detected = any(r["status"] == "DRIFT" for r in results.values())
        if drift_detected:
            self._send_alert(results)

        return results, drift_detected

    def _send_alert(self, results):
        alert_msg = f"DRIFT ALERT at {datetime.now()}\n"
        for feat, info in results.items():
            if info["status"] != "OK":
                alert_msg += f"  {feat}: PSI={info['psi']} ({info['status']})\n"
        print(alert_msg)
        # In production: send to Slack/PagerDuty

# Weekly monitoring job
monitor = DriftMonitor("model/reference_stats.npy")
# features = load_current_week_features()
# results, drift = monitor.check_weekly_drift(features)

输出:

TEXT 📖 仅展示
DriftMonitor initialized with reference stats
Weekly check scheduled: PSI threshold=0.2

6. 自动重训练与项目复盘

(1) 自动重训练流水线

▶ 示例:Airflow DAG概念

PYTHON
# File: dags/retrain_pipeline.py (Airflow DAG concept)
# from airflow import DAG
# from airflow.operators.python import PythonOperator

def retrain_pipeline():
    """Automated retraining pipeline triggered by drift detection."""
    # Step 1: Check drift
    # drift_detected = check_weekly_drift()

    # Step 2: Fetch recent data
    # data = fetch_recent_data(months=3)

    # Step 3: Retrain model
    # new_model, metrics = train_lightgbm(data)

    # Step 4: Compare with production
    # if metrics["mape"] < production_mape:
    #     register_model(new_model, stage="Staging")

    # Step 5: A/B test (2 weeks)
    # run_ab_test(new_model, duration_weeks=2)

    # Step 6: Promote if A/B shows improvement
    # if ab_test_significant:
    #     promote_to_production(new_model)

    # Step 7: Update reference stats
    # update_reference_stats(new_model)
    pass

# DAG definition
# with DAG("salespredict_retrain", schedule_interval="0 6 * * 1") as dag:
#     check_drift = PythonOperator(task_id="check_drift", python_callable=check_drift)
#     retrain = PythonOperator(task_id="retrain", python_callable=retrain_model)
#     ab_test = PythonOperator(task_id="ab_test", python_callable=run_ab_test)
#     promote = PythonOperator(task_id="promote", python_callable=promote_model)
#     check_drift >> retrain >> ab_test >> promote

输出:

TEXT 📖 仅展示
DAG 'salespredict_retrain' registered
Schedule: every Monday 06:00 UTC
Tasks: check_drift >> retrain >> ab_test >> promote

(2) 项目复盘

100%
graph TB
    START[Week 1: Project Kickoff] --> DATA[Week 2-3: Data Pipeline]
    DATA --> BASELINE[Week 3: Baseline LR<br/>MAPE 15%]
    BASELINE --> XGB[Week 4-5: XGBoost<br/>MAPE 9%]
    XGB --> LGBM[Week 5-6: LightGBM + Optuna<br/>MAPE 8%]
    LGBM --> DEPLOY[Week 7: FastAPI + Docker]
    DEPLOY --> MONITOR[Week 8: Monitoring + Drift]
    MONITOR --> LIVE[Production Live<br/>MAPE 8% stable]
维度 起始值 最终值 提升
预测MAPE 25%(Excel) 8%(LightGBM) 68%↓
库存成本 500k USD/年 100k USD/年 400k节省
预测延迟 人工2天 API 10ms 1.7亿倍↓
实验效率 3次/周 80+次/天 180倍↑

▶ 示例:全链路回顾

PYTHON
project_retrospective = {
    "what_went_well": [
        "MLflow experiment tracking prevented confusion over 50+ runs",
        "TimeSeriesSplit caught future leakage before production",
        "Optuna found better params than manual search in 1/10 time",
        "Docker deployment enabled zero-downtime model updates",
    ],
    "what_could_improve": [
        "Data pipeline should have been built before model training",
        "Feature store would reduce duplication between training and serving",
        "A/B test should have run longer (3 weeks instead of 2)",
        "Monitoring dashboard should have been set up from day 1",
    ],
    "key_learnings": [
        "Data quality > model complexity (clean data + simple model beats dirty data + complex model)",
        "Design first, code second (1 week design saved 4 weeks of rework)",
        "Deploy early, iterate fast (production feedback > lab experiments)",
        "Monitor everything (drift detection caught Double-11 issue in 3 days)",
    ],
    "next_steps": [
        "Add LSTM model for temporal pattern capture",
        "Implement real-time feature serving with Feast",
        "Build automated retraining pipeline with Airflow",
        "Expand to 3 more markets (Japan, India, Brazil)",
    ],
}

for category, items in project_retrospective.items():
    print(f"\n{category.replace('_', ' ').title()}:")
    for item in items:
        print(f"  - {item}")

输出:

TEXT 📖 仅展示
What Went Well:
  - MLflow experiment tracking prevented confusion over 50+ runs
  - TimeSeriesSplit caught future leakage before production
What Could Improve:
  - Data pipeline should have been built before model training
Key Learnings:
  - Data quality > model complexity
Next Steps:
  - Add LSTM model for temporal pattern capture

❓ 常见问题

Q API服务需要多少资源?
A LightGBM推理CPU密集型——2核4GB可处理100+ QPS。PyTorch MLP建议4核8GB+GPU(如果模型大)。Redis缓存256MB够用。
Q 模型更新怎么不停服?
A 两种方案——1) 蓝绿部署(新旧版本切换,切流量);2) 滚动更新(docker-compose rolling update,逐个替换)。配合MLflow版本管理。
Q 监控仪表盘需要监控哪些指标?
A 三层——基础设施(延迟/可用性/内存)、ML指标(预测分布/特征统计/模型版本)、业务指标(MAPE/Revenue偏差/转化率)。任何层异常都需要告警。
Q 自动重训练多久一次?
A 默认每周检测漂移、每月重训练(无漂移也定期更新)。关键事件(促销/新品上线)后立即检测。漂移触发即时重训练。
Q 部署后模型退化怎么办?
A 四步响应——1) 检查数据漂移(PSI);2) 检查概念漂移(MAPE趋势);3) 重训练+验证;4) A/B测试后灰度上线。全程MLflow记录。
Q 25课学完了,下一步学什么?
A 三个方向——1) 深度学习进阶(Transformer/NLP/计算机视觉);2) MLOps进阶(Kubeflow/特征存储/数据版本管理);3) 实战更多项目(Kaggle/开源贡献)。基础已扎实,选择兴趣方向深入。

📖 小节


📝 作业

  1. 基础题(难度⭐):将训练好的模型保存为joblib文件,编写最小FastAPI /predict端点,用uvicorn运行并测试。提示:参考第3节API代码。
  2. 进阶题(难度⭐⭐):编写Dockerfile + docker-compose.yml(API+Redis+Nginx),构建镜像并运行完整服务栈,验证/health端点和负载均衡。提示:参考第4节Docker配置。
  3. 挑战题(难度⭐⭐⭐):实现完整的生产部署——FastAPI+Prometheus指标+Grafana仪表盘+漂移监控脚本,运行1周模拟数据,检测注入的漂移并触发告警。提示:综合第3-6节所有代码。

← 上一课:项目开发 | 教程结束 →

Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏