Machine Learning: 项目部署 — SalesPredict上线与运维指南
部署是终点也是起点——上线意味着运维开始,持续可靠才是真正的交付。
1. 你将学到
- FastAPI推理服务:请求校验、批量预测、缓存策略(Redis)、限流与降级
- Docker部署:多阶段Dockerfile、docker-compose编排、健康检查配置
- 监控仪表盘:Grafana + Prometheus,预测量/延迟/准确率/Revenue实时监控
- 漂移检测与自动重训练:PSI周检测,超阈值触发Airflow重训练流水线
- 项目复盘:Bob的SalesPredict从0到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) 项目复盘
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/开源贡献)。基础已扎实,选择兴趣方向深入。
📖 小节
- FastAPI生产服务:Pydantic校验、批量预测、中间件日志、Prometheus指标
- Docker部署:多阶段构建减镜像、docker-compose编排全栈、HEALTHCHECK健康检查
- 监控仪表盘:Prometheus采集+Grafana展示,三层监控(基础设施/ML/业务)
- 漂移检测:PSI周检测,超0.2触发告警,Airflow DAG编排重训练流水线
- 自动重训练:检测→数据回捞→训练→A/B验证→灰度上线,全程自动化
- SalesPredict全链路:Excel 25% MAPE → LightGBM 8% MAPE,年节省400k USD库存成本
📝 作业
- 基础题(难度⭐):将训练好的模型保存为joblib文件,编写最小FastAPI /predict端点,用uvicorn运行并测试。提示:参考第3节API代码。
- 进阶题(难度⭐⭐):编写Dockerfile + docker-compose.yml(API+Redis+Nginx),构建镜像并运行完整服务栈,验证/health端点和负载均衡。提示:参考第4节Docker配置。
- 挑战题(难度⭐⭐⭐):实现完整的生产部署——FastAPI+Prometheus指标+Grafana仪表盘+漂移监控脚本,运行1周模拟数据,检测注入的漂移并触发告警。提示:综合第3-6节所有代码。