graph TB
subgraph L4[模型质量层]
M1[预测准确率]
M2[数据漂移 PSI]
M3[特征分布偏移]
M4[公平性指标]
end
subgraph L3[ML 系统层]
S1[推理延迟分布]
S2[吞吐量 QPS]
S3[GPU 利用率]
S4[批处理队列长度]
end
subgraph L2[应用层]
A1[请求成功率]
A2[错误率]
A3[请求日志]
end
subgraph L1[基础设施层]
I1[节点健康度]
I2[网络 I/O]
I3[磁盘使用率]
I4[内存使用率]
end
L1 --> L2 --> L3 --> L4
第17章 可观测性与运维
第17章 可观测性与运维
“你没有监控的系统,就是在黑暗中飞行——而且你还蒙着眼睛。”
17.1 为什么 ML 可观测性不同于传统系统
传统软件系统的可观测性关注三件事:日志(Logs)、指标(Metrics)、链路追踪(Traces)。这三者对 ML 系统同样重要,但还不够。
ML 系统多了一个维度:模型行为本身。一个推理服务可能返回 200 OK,延迟也在正常范围内,但模型输出的质量可能已经因为数据漂移而严重退化。传统的 SLO 完全检测不到这种问题。
ML 系统的可观测性分层
这个分层模型的核心洞察是:越底层的指标越容易采集,但越不能反映真实业务影响。基础设施正常不代表应用正常,应用正常不代表模型正常,模型正常不代表预测有效。
17.2 训练与推理统一监控体系
训练阶段监控
训练阶段的监控重点是资源利用率和训练过程的健康度。
# train_monitor.py - 训练过程监控
import time
import torch
import psutil
from prometheus_client import CollectorRegistry, push_to_gateway
class TrainingMonitor:
def __init__(self, pushgateway_url="localhost:9091"):
self.registry = CollectorRegistry()
self.pushgateway = pushgateway_url
self.start_time = time.time()
def log_step_metrics(self, step, loss, lr, model, optimizer):
"""每个训练 step 记录的指标"""
metrics = {
"train_loss": loss,
"learning_rate": lr,
"step": step,
"elapsed_hours": (time.time() - self.start_time) / 3600,
}
# GPU 指标
for i in range(torch.cuda.device_count()):
metrics[f"gpu_{i}_memory_used_gb"] = (
torch.cuda.memory_allocated(i) / 1024**3
)
metrics[f"gpu_{i}_memory_cached_gb"] = (
torch.cuda.memory_reserved(i) / 1024**3
)
# CPU 和内存
metrics["cpu_percent"] = psutil.cpu_percent()
metrics["ram_used_gb"] = psutil.virtual_memory().used / 1024**3
# 关键:检查梯度异常
total_grad_norm = 0.0
for p in model.parameters():
if p.grad is not None:
total_grad_norm += p.grad.data.norm().item() ** 2
total_grad_norm = total_grad_norm ** 0.5
metrics["grad_norm"] = total_grad_norm
# 推送到 Prometheus Pushgateway
push_to_gateway(
self.pushgateway,
job="training_monitor",
registry=self.registry,
grouping_key={"run_id": self.run_id}
)
# 异常检测
self._check_anomalies(metrics)
return metrics
def _check_anomalies(self, metrics):
"""训练异常检测"""
issues = []
# 梯度爆炸
if metrics.get("grad_norm", 0) > 1000:
issues.append(f"⚠️ 梯度爆炸: grad_norm={metrics['grad_norm']:.2f}")
# Loss 为 NaN
if metrics.get("train_loss", 0) != metrics.get("train_loss", 0): # NaN check
issues.append("⚠️ Loss 为 NaN,训练已崩溃")
# GPU 内存接近上限
for i in range(torch.cuda.device_count()):
mem_used = metrics.get(f"gpu_{i}_memory_used_gb", 0)
if mem_used > 75: # 80GB 卡,75GB 阈值
issues.append(f"⚠️ GPU {i} 内存即将溢出: {mem_used:.1f}GB")
if issues:
for issue in issues:
print(issue)
# 发送告警
self._send_alert(issues)
def _send_alert(self, issues):
"""发送告警到飞书/钉钉/Slack"""
# 实际实现中对接告警平台
pass推理阶段监控
推理阶段除了常规的系统指标,还需要监控模型特有的指标:
# inference_monitor.py - 推理服务监控
from dataclasses import dataclass, field
from collections import deque
import numpy as np
import time
@dataclass
class PredictionLog:
"""单次预测的完整日志"""
request_id: str
timestamp: float
input_features: dict
prediction: any
confidence: float
latency_ms: float
model_version: str
class InferenceMonitor:
def __init__(self, window_size=10000):
self.prediction_buffer = deque(maxlen=window_size)
self.latency_buffer = deque(maxlen=window_size)
self.confidence_buffer = deque(maxlen=window_size)
def log_prediction(self, log: PredictionLog):
self.prediction_buffer.append(log)
self.latency_buffer.append(log.latency_ms)
self.confidence_buffer.append(log.confidence)
def compute_serving_metrics(self):
"""计算推理服务指标"""
latencies = np.array(self.latency_buffer)
confidences = np.array(self.confidence_buffer)
return {
# 延迟分位数
"latency_p50_ms": np.percentile(latencies, 50),
"latency_p95_ms": np.percentile(latencies, 95),
"latency_p99_ms": np.percentile(latencies, 99),
# 吞吐量
"qps": len(self.latency_buffer) / max(
time.time() - self.prediction_buffer[0].timestamp, 1
),
# 置信度分布(异常检测的关键信号)
"confidence_mean": np.mean(confidences),
"confidence_std": np.std(confidences),
"low_confidence_ratio": np.mean(confidences < 0.5), # 低置信度比例
# 错误率
"error_rate": np.mean(np.array(latencies) < 0), # 负延迟表示错误
}
def detect_drift(self, baseline_confidence):
"""
基于置信度分布漂移检测
当平均置信度显著下降时,可能存在数据漂移
"""
current_confidence = np.mean(self.confidence_buffer)
drift_ratio = (baseline_confidence - current_confidence) / baseline_confidence
if drift_ratio > 0.1: # 置信度下降超过 10%
return {
"drift_detected": True,
"severity": "high" if drift_ratio > 0.2 else "medium",
"message": f"置信度下降 {drift_ratio:.1%},疑似数据漂移"
}
return {"drift_detected": False}17.3 日志、指标与链路追踪
日志策略
ML 系统的日志需要覆盖三种信息:
# structured_logging.py
import structlog
import json
logger = structlog.get_logger()
def configure_logging():
structlog.configure(
processors=[
structlog.processors.TimeStamper(fmt="iso"),
structlog.processors.add_log_level,
structlog.processors.JSONRenderer(),
],
)
def log_inference_request(request_id, model_input, model_output, metadata):
"""结构化日志:推理请求"""
logger.info(
"inference_request",
request_id=request_id,
model_version=metadata["model_version"],
input_shape=list(model_input.shape) if hasattr(model_input, 'shape') else None,
output_class=str(model_output),
latency_ms=metadata["latency_ms"],
gpu_id=metadata.get("gpu_id"),
# 关键:记录输入特征的关键统计量(不记录原始数据)
input_stats={
"mean": float(model_input.float().mean()),
"std": float(model_input.float().std()),
"min": float(model_input.float().min()),
"max": float(model_input.float().max()),
},
)
def log_training_step(run_id, step, loss, metrics):
"""结构化日志:训练步骤"""
logger.info(
"training_step",
run_id=run_id,
step=step,
loss=loss,
**metrics,
)日志成本管理: ML 系统的日志量可能非常大(每个推理请求都记录)。设置合理的采样率:正常请求 1% 采样,异常请求 100% 记录。使用像 OpenTelemetry 这样的标准来确保日志可以在不同系统间关联。
Prometheus 指标采集
# prometheus-config.yaml - ML 系统的 Prometheus 配置
scrape_configs:
# 推理服务指标
- job_name: 'inference-server'
scrape_interval: 15s
metrics_path: /metrics
static_configs:
- targets: ['inference-service:8080']
# 训练任务指标(通过 Pushgateway)
- job_name: 'training-jobs'
scrape_interval: 30s
honor_labels: true
static_configs:
- targets: ['pushgateway:9091']
# GPU 指标
- job_name: 'gpu-metrics'
scrape_interval: 15s
static_configs:
- targets: ['dcgm-exporter:9400']
# 特征存储指标
- job_name: 'feature-store'
scrape_interval: 30s
static_configs:
- targets: ['feast-server:6566']
# 告警规则
rule_files:
- 'alerts/*.yml'关键告警规则
# alerts/ml-alerts.yml
groups:
- name: ml_serving_alerts
rules:
# 推理延迟告警
- alert: HighInferenceLatency
expr: |
histogram_quantile(0.99,
rate(inference_latency_seconds_bucket[5m])
) > 0.5
for: 5m
labels:
severity: warning
annotations:
summary: "推理 P99 延迟超过 500ms"
description: "{{ $labels.model }} 的 P99 延迟为 {{ $value }}s"
# GPU 内存使用告警
- alert: GPUMemoryHigh
expr: |
DCGM_FI_DEV_FB_USED / DCGM_FI_DEV_FB_TOTAL > 0.9
for: 10m
labels:
severity: critical
annotations:
summary: "GPU {{ $labels.gpu }} 内存使用率超过 90%"
# 模型错误率告警
- alert: HighModelErrorRate
expr: |
rate(inference_errors_total[5m]) / rate(inference_requests_total[5m]) > 0.01
for: 2m
labels:
severity: critical
annotations:
summary: "模型 {{ $labels.model_version }} 错误率超过 1%"
# 数据漂移告警
- alert: DataDriftDetected
expr: model_data_drift_psi > 0.2
for: 15m
labels:
severity: warning
annotations:
summary: "检测到数据漂移 (PSI > 0.2)"
description: "特征 {{ $labels.feature }} 的 PSI = {{ $value }}"分布式链路追踪
在微服务架构下,一个推理请求可能经过多个服务。链路追踪帮助你理解请求在不同服务间的流转:
# tracing.py - 使用 OpenTelemetry
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
# 配置 tracing
trace.set_tracer_provider(TracerProvider())
trace.get_tracer_provider().add_span_processor(
BatchSpanProcessor(OTLPSpanExporter(endpoint="otel-collector:4317"))
)
tracer = trace.get_tracer(__name__)
async def handle_inference_request(self, request):
"""推理请求的完整链路"""
with tracer.start_as_current_span("inference_request") as span:
span.set_attribute("request.id", request.id)
span.set_attribute("model.version", self.model_version)
# 子 span: 特征查询
with tracer.start_as_current_span("feature_lookup"):
features = await self.feature_store.get_online_features(
request.entity_id
)
span.set_attribute("feature.count", len(features))
# 子 span: 模型推理
with tracer.start_as_current_span("model_inference"):
output = self.model.predict(features)
span.set_attribute("prediction.confidence", output.confidence)
# 子 span: 后处理
with tracer.start_as_current_span("post_process"):
result = self.post_processor(output)
return result17.4 告警策略与事件响应
告警分级
不是所有告警都需要半夜把人叫醒。合理的告警分级至关重要:
| 级别 | 定义 | 响应时间 | 通知方式 |
|---|---|---|---|
| P0 - Critical | 服务完全不可用 | 5 分钟 | 电话 + 短信 + IM |
| P1 - High | 核心功能受损 | 15 分钟 | 短信 + IM |
| P2 - Medium | 性能退化但不影响核心功能 | 1 小时 | IM |
| P3 - Low | 需要关注但不紧急 | 1 工作日 | 邮件 / 仪表板 |
ML 特有的告警场景
graph TD
A[ML 告警触发] --> B{告警类型}
B -->|系统告警| C[基础设施/应用层]
B -->|模型告警| D[模型行为层]
B -->|数据告警| E[数据质量层]
C --> C1[推理服务不可用 - P0]
C --> C2[GPU OOM - P1]
C --> C3[延迟飙升 - P1/P2]
D --> D1[错误率飙升 - P0/P1]
D --> D2[置信度异常下降 - P2]
D --> D3[预测分布突变 - P2]
E --> E1[特征缺失率过高 - P1]
E --> E2[数据漂移 PSI > 0.2 - P2]
E --> E3[数据新鲜度过期 - P1]
事件响应 Runbook
# Runbook: 推理服务延迟飙升
## 症状
- Prometheus 告警: HighInferenceLatency
- P99 延迟 > 500ms 持续 5 分钟以上
## 排查步骤
### Step 1: 确认影响范围 (2 分钟)
```bash
# 查看哪些 Pod 受影响
kubectl get pods -l app=inference-server -o wide
# 查看 Grafana 仪表板
echo "打开 http://grafana.internal/d/inference-overview"Step 2: 检查 GPU 状态 (3 分钟)
# 检查 GPU 利用率和内存
nvidia-smi --query-gpu=index,utilization.gpu,memory.used,memory.total --format=csv
# 检查是否有异常进程
nvidia-smiStep 3: 检查请求模式 (3 分钟)
# 是否有异常流量?
kubectl logs -l app=inference-server --tail=1000 | jq '.request_size' | sort -n | tail -20
# 检查 batch size 是否异常
kubectl logs deployment/inference-server | grep "batch_size" | tail -100Step 4: 紧急处置 (5 分钟)
# 方案 A: 水平扩容
kubectl scale deployment inference-server --replicas=8
# 方案 B: 开启降级模式(跳过大模型,使用小模型)
kubectl set env deployment/inference-server MODEL_STRATEGY=fallback
# 方案 C: 限制并发
kubectl annotate deployment inference-server \
autoscaling.knative.dev/target: "5"Step 5: 事后复盘
- 记录事件时间线
- 更新告警阈值(如果误报)
- 添加自动化处理(如果可重复)
::: callout-warning
**告警疲劳:** 这是 ML 运维中最常见的问题。如果每天收到 100 条告警,人们会忽略所有告警。定期审查告警规则:过去 30 天内从未触发过的告警考虑删除;误报率超过 30% 的告警需要重新调整阈值。
:::
## 17.5 资源利用率与成本看板
### GPU 利用率监控
GPU 是 ML 基础设施中最昂贵的资源。监控 GPU 利用率不仅是技术问题,更是财务问题。
```python
# gpu_monitor.py
import subprocess
import json
from dataclasses import dataclass
from typing import List
@dataclass
class GPUMetric:
index: int
gpu_util: float # GPU 计算利用率 %
mem_util: float # 显存利用率 %
mem_used_gb: float
mem_total_gb: float
power_draw_w: float # 功耗 W
temperature_c: float # 温度 °C
running_processes: List[dict]
def collect_gpu_metrics() -> List[GPUMetric]:
"""使用 nvidia-smi 采集 GPU 指标"""
result = subprocess.run([
"nvidia-smi",
"--query-gpu=index,utilization.gpu,utilization.memory,"
"memory.used,memory.total,power.draw,temperature.gpu",
"--format=csv,noheader,nounits"
], capture_output=True, text=True)
metrics = []
for line in result.stdout.strip().split('\n'):
parts = [float(x.strip()) for x in line.split(',')]
metrics.append(GPUMetric(
index=int(parts[0]),
gpu_util=parts[1],
mem_util=parts[2],
mem_used_gb=parts[3] / 1024,
mem_total_gb=parts[4] / 1024,
power_draw_w=parts[5],
temperature_c=parts[6],
running_processes=[]
))
return metrics
def compute_gpu_efficiency(metrics: List[GPUMetric]):
"""
计算 GPU 利用效率
- SM 活跃率: GPU 计算是否在运行
- 内存带宽利用率: 数据传输是否充分
- 功耗效率: 每瓦特能做多少推理
"""
total_gpus = len(metrics)
active_gpus = sum(1 for m in metrics if m.gpu_util > 10)
idle_gpus = total_gpus - active_gpus
avg_util = sum(m.gpu_util for m in metrics) / total_gpus
avg_mem = sum(m.mem_used_gb / m.mem_total_gb * 100 for m in metrics) / total_gpus
total_power = sum(m.power_draw_w for m in metrics)
return {
"total_gpus": total_gpus,
"active_gpus": active_gpus,
"idle_gpus": idle_gpus,
"avg_gpu_utilization": avg_util,
"avg_memory_utilization": avg_mem,
"total_power_kw": total_power / 1000,
"estimated_hourly_cost_usd": total_gpus * 2.48, # A100 按需约 $2.48/h
}
成本看板设计
一个好的成本看板应该能回答以下问题:
# cost-dashboard.yaml - 成本看板指标定义
dashboards:
- name: "ML 平台成本总览"
panels:
# 面板 1: 每日 GPU 成本趋势
- title: "每日 GPU 成本"
query: |
sum(rate(gpu_hourly_cost[1d])) by (team)
# 面板 2: GPU 利用率分布
- title: "GPU 利用率分布"
query: |
histogram_quantile(0.50, rate(gpu_utilization_bucket[1h]))
histogram_quantile(0.95, rate(gpu_utilization_bucket[1h]))
# 面板 3: 各团队资源消耗
- title: "团队资源消耗"
query: |
sum(node_gpu_hourly_cost) by (team_label)
# 面板 4: 每千次推理成本
- title: "每千次推理成本"
query: |
daily_gpu_cost / (total_inference_count / 1000)
# 面板 5: 空闲 GPU 浪费
- title: "空闲 GPU 成本浪费"
query: |
count(gpu_utilization < 5) * gpu_hourly_cost
# 面板 6: 训练 vs 推理成本占比
- title: "训练 vs 推理成本"
query: |
sum(gpu_hourly_cost{workload_type="training"})
sum(gpu_hourly_cost{workload_type="inference"})关键成本指标
| 指标 | 计算方式 | 目标值 |
|---|---|---|
| GPU 平均利用率 | avg(DCGM_FI_DEV_GPU_UTIL) |
> 70% |
| GPU 闲置率 | count(gpu_util < 5%) / total_gpus |
< 10% |
| 每 1K 推理成本 | daily_cost / (inferences / 1000) |
逐月下降 |
| 训练 ROI | (模型上线后收益 - 训练成本) / 训练成本 |
> 3x |
| 存储成本占比 | storage_cost / total_infra_cost |
< 15% |
17.6 容量规划与集群健康度
容量规划模型
容量规划的核心是预测未来需求并提前准备资源:
# capacity_planning.py
import numpy as np
from scipy import stats
from datetime import datetime, timedelta
class CapacityPlanner:
def __init__(self, historical_data):
"""
historical_data: list of {date, qps, gpu_count, latency_p99}
"""
self.data = historical_data
def forecast_qps(self, days_ahead=30):
"""预测未来 QPS"""
dates = [d["date"] for d in self.data]
qps_values = [d["qps"] for d in self.data]
# 线性回归
x = np.arange(len(qps_values))
slope, intercept, r_value, _, _ = stats.linregress(x, qps_values)
forecast = []
for i in range(len(qps_values), len(qps_values) + days_ahead):
forecast.append({
"date": dates[-1] + timedelta(days=i - len(qps_values) + 1),
"predicted_qps": slope * i + intercept,
"confidence": r_value ** 2 # R²
})
return forecast
def recommend_gpu_allocation(self, target_latency_p99_ms=200):
"""
基于预测 QPS 推荐 GPU 数量
"""
forecast = self.forecast_qps()
# 每个 GPU 能承载的最大 QPS(通过压测获得)
max_qps_per_gpu = 50 # 示例值
recommendations = []
for f in forecast:
required_qps = f["predicted_qps"] * 1.3 # 30% buffer
gpu_needed = int(np.ceil(required_qps / max_qps_per_gpu))
recommendations.append({
"date": f["date"],
"predicted_qps": f["predicted_qps"],
"recommended_gpus": gpu_needed,
"current_gpus": self.data[-1]["gpu_count"],
"action": "scale_up" if gpu_needed > self.data[-1]["gpu_count"]
else "maintain"
})
return recommendations
def estimate_lead_time(self, action="scale_up"):
"""
估算扩容提前期
- 云上按需: ~5 分钟
- 云上预留: 1-7 天
- 自建采购: 4-12 周
"""
lead_times = {
"cloud_ondemand": timedelta(minutes=5),
"cloud_reserved": timedelta(days=7),
"self_hosted": timedelta(weeks=8),
}
return lead_times集群健康度评分
# cluster_health.py
def compute_cluster_health_score(metrics):
"""
综合集群健康度评分 (0-100)
"""
scores = {}
# 1. 节点健康度 (30%)
total_nodes = metrics["total_nodes"]
healthy_nodes = metrics["healthy_nodes"]
scores["node_health"] = (healthy_nodes / total_nodes) * 100
# 2. GPU 利用率 (25%)
avg_gpu_util = metrics["avg_gpu_utilization"]
# 理想范围 60-85%,过高或过低都扣分
if 60 <= avg_gpu_util <= 85:
scores["gpu_utilization"] = 100
elif avg_gpu_util > 85:
scores["gpu_utilization"] = max(0, 100 - (avg_gpu_util - 85) * 3)
else:
scores["gpu_utilization"] = max(0, avg_gpu_util / 60 * 100)
# 3. 错误率 (20%)
error_rate = metrics["error_rate"]
scores["error_rate"] = max(0, 100 - error_rate * 1000) # 0.1% 错误率 → 90 分
# 4. 资源碎片化 (15%)
fragmentation = metrics["resource_fragmentation_ratio"] # 0-1
scores["fragmentation"] = (1 - fragmentation) * 100
# 5. 成本效率 (10%)
cost_per_1k_inference = metrics["cost_per_1k_inference"]
target_cost = metrics["target_cost_per_1k_inference"]
if cost_per_1k_inference <= target_cost:
scores["cost_efficiency"] = 100
else:
scores["cost_efficiency"] = max(0, target_cost / cost_per_1k_inference * 100)
# 加权总分
weights = {
"node_health": 0.30,
"gpu_utilization": 0.25,
"error_rate": 0.20,
"fragmentation": 0.15,
"cost_efficiency": 0.10,
}
total_score = sum(scores[k] * weights[k] for k in weights)
return {
"total_score": round(total_score, 1),
"breakdown": {k: round(v, 1) for k, v in scores.items()},
"grade": (
"A" if total_score >= 90 else
"B" if total_score >= 75 else
"C" if total_score >= 60 else
"D"
),
}容量规划的经验法则: 保持 30% 的冗余容量。这意味着在高峰期,你的集群仍然有 30% 的 GPU 可以立即分配。低于 20% 时触发扩容告警,低于 10% 时触发紧急扩容流程。
17.7 小结
ML 系统的可观测性是一个多层次工程:
- 基础设施监控确保节点和 GPU 健康
- 应用监控确保推理服务可用
- 模型监控确保预测质量没有退化
- 成本监控确保资源使用效率
关键原则:只告警可操作的问题。每一条告警都应该有对应的 Runbook,告诉值班人员该做什么。没有 Runbook 的告警只是噪音。
可观测性建设是一个持续迭代的过程。每发生一次线上事故,都应该问自己:“我们能在多早的阶段发现这个问题?” 然后把那个阶段的监控补上。
延伸阅读
- Observability for Machine Learning (Huyen, 2023) — ML 系统可观测性的系统化框架
- Prometheus: Up & Running (Wilhelm & Rabkin, 2018) — 监控系统实践指南
- Google SRE Workbook — 第 6 章「Monitoring Distributed Systems」是告警策略的圣经
- OpenTelemetry 官方文档 opentelemetry.io — 统一的可观测性标准
- NVIDIA DCGM docs.nvidia.com/datacenter/dcgm — GPU 监控最佳实践
- Evidently AI Blog evidentlyai.com/blog — ML 监控和数据漂移检测的深度文章