1. 项目概述:这不是一次“部署”,而是一场从实验室到产线的系统性迁移

“From Notebook to Production: Running ML in the Real World (Part 4)”——这个标题里藏着一个被无数数据科学家反复咀嚼、又悄悄咽下的苦涩真相: Jupyter Notebook 从来就不是生产环境的入口,它只是我们理解问题的第一张草稿纸。 我在带团队做模型交付的六年里,亲手把超过37个模型从本地笔记本推上生产服务,其中21个在上线首周就因“看似合理、实则致命”的设计缺陷被紧急回滚。Part 4 不是系列文章的收尾,恰恰是真正硬仗的开始:它聚焦的是模型在真实业务流中持续存活的能力——不是“能跑”,而是“稳跑”、“可查”、“可调”、“可退”。核心关键词非常明确: ML 生产化(MLOps)、模型监控、数据漂移检测、在线推理服务、A/B 测试框架、可观测性(Observability) 。它解决的不是“怎么把 pickle 文件扔进 Flask API”,而是当用户点击“推荐商品”按钮时,后端服务在每秒处理 800+ 请求、日均调用 2300 万次、上游数据源每小时更新 12 个特征表、下游业务方每天提 5 条新指标需求的情况下,如何让模型预测结果依然可信、可解释、可归因。适合三类人深度参考:刚从 Kaggle 赛道转战工业界的算法工程师(别再只盯着 AUC 了);正在搭建公司级 MLOps 基础设施的平台工程师(别再只堆 Kubernetes YAML 了);以及技术背景扎实、需要对模型风险有实质管控权的数据产品负责人(别再只问“准确率多少”了)。这篇文章不讲理论推导,只讲我在金融风控、电商搜索、IoT 设备预测三个高并发、强监管、低容错场景中,踩过坑、验证过、现在还在用的实操路径。

2. 内容整体设计与思路拆解:为什么 Part 4 必须放弃“单点部署思维”

2.1 从“单次发布”到“持续闭环”的范式转移

很多团队卡在 Part 3(模型封装与 API 化)就以为大功告成,结果上线三天后发现:线上预测延迟从 12ms 涨到 320ms,但日志里只有 INFO:root:Request processed 这一行;A/B 测试组反馈新模型点击率 +2.3%,但风控同事同步指出逾期率上升 0.8 个百分点,而这两个指标在训练集里根本没做过联合分析;更糟的是,当业务方要求“把上周三下午的预测结果全拉出来复盘”时,你才发现所有原始请求体和响应体都没存——因为当初觉得“太占磁盘”。Part 4 的设计起点,就是彻底抛弃“模型部署即终点”的幻觉。我把它拆解为四个不可割裂的子系统:

  • 可观测性层(Observability Layer) :不是简单加 Prometheus metrics,而是构建“请求-特征-预测-结果”四维追踪链。比如在电商搜索中,每个商品曝光请求必须携带 request_id user_segment_id query_hash device_type 四个关键上下文标签,这些标签要贯穿从 Nginx access log → 特征服务 → 模型推理 → 结果缓存 → 用户行为埋点的全链路。这样当某类用户点击率异常时,才能精准下钻到“iOS 17.4 用户 + 长尾查询词 + 高价值商品”这个组合维度。

  • 监控与告警层(Monitoring & Alerting) :拒绝“CPU > 90% 就告警”的粗暴逻辑。我们定义三类黄金指标:

    • SLO 指标 :P95 推理延迟 ≤ 80ms(业务 SLA 要求),错误率 < 0.05%(基于历史 30 天基线动态计算);
    • 数据健康指标 :输入特征分布偏移(KS 检验 p-value < 0.01)、缺失值率突增 > 5 倍标准差、特征值域越界(如年龄字段出现负数);
    • 模型表现指标 :在线 AUC 滑动窗口下降 > 0.015(对比过去 7 天均值)、预测置信度分布偏移(KL 散度 > 0.3)。
      这些指标全部接入 Grafana,但告警规则写在代码里(不是 Grafana UI 里点点点),确保可版本化、可测试、可审计。
  • 自动化响应层(Auto-Response Layer) :监控触发 ≠ 人工救火。我们在支付风控模型中实现了三级响应:一级(p-value < 0.001)自动冻结该特征输入,改用历史均值填充;二级(AUC 下降 > 0.02)触发影子流量(Shadow Traffic),将 5% 真实请求同时打给新旧两个模型,生成差异报告;三级(错误率 > 0.1%)自动执行熔断(Circuit Breaker),将流量切回前一稳定版本,并发邮件给模型 owner 和 SRE。整个过程平均耗时 47 秒,比人工介入快 11 倍。

  • 实验治理层(Experiment Governance) :所有线上实验必须通过统一网关注册,强制填写:实验目标(如“提升 GMV”)、评估周期(≥ 7 天)、对照组/实验组流量分配策略、核心评估指标(必须含业务指标+模型指标)、回滚条件(如“GMV 下降 > 1% 或逾期率上升 > 0.3%”)。我们用 Airflow DAG 实现了“实验注册 → 特征版本锁定 → 模型版本绑定 → 流量配置下发 → 指标自动采集 → 报告生成 → 决策建议”的全自动流水线。去年 Q3 共运行 42 个实验,其中 17 个因未达预期被自动终止,避免了无效模型上线。

提示:不要试图用一个工具解决所有问题。我们用 OpenTelemetry 做分布式追踪,用 Evidently 做数据漂移检测,用 Prometheus+Alertmanager 做指标告警,用 Argo Workflows 做实验编排——它们之间用轻量级 gRPC 接口通信,而不是强行塞进一个“全能平台”。

2.2 为什么必须放弃“模型为中心”的架构,转向“服务契约为中心”

新手常犯的致命错误,是把模型文件(.pkl/.onnx)当作部署单元。但在真实世界里,模型只是服务的一个组件。真正的部署单元,是 Service Contract(服务契约) ——一份明确定义了输入输出 Schema、SLA、依赖关系、降级策略的协议文档。以我们为物流调度系统开发的 ETA(预计到达时间)模型为例,其契约包含:

字段 类型 描述 示例 强制性
order_id string 订单唯一标识 "ORD-2024-789456" 必填
pickup_latlng array[float] 取货点经纬度 [31.2304, 121.4737] 必填
dropoff_latlng array[float] 送货点经纬度 [31.1923, 121.3985] 必填
vehicle_type string 车辆类型枚举 "motorcycle", "van" 必填
weather_condition string 天气状况(可选) "rainy", "sunny" 可选
response_time_ms int P95 响应延迟 65 SLA
fallback_strategy string 降级策略 "return_static_table" 必填

这份契约决定了:前端 SDK 必须校验 order_id 格式;特征服务必须在 15ms 内返回 weather_condition ,超时则走默认值;模型服务收到非法 vehicle_type 时必须返回 HTTP 400 并附带错误码 INVALID_VEHICLE_TYPE ;当预测失败时,网关必须按契约调用静态查表服务返回兜底 ETA。 契约不是文档,而是代码——我们用 Protobuf 定义它,并自动生成 Python/Java/Go 的客户端和服务端 stub,所有校验逻辑嵌入生成代码中。 这样,当业务方提出“增加电动车类型支持”时,改动不是改模型代码,而是更新契约中的 vehicle_type 枚举值,重新生成代码,所有上下游自动适配。我们因此将模型迭代平均周期从 11 天压缩到 3.2 天。

2.3 成本与复杂度的现实平衡:不做“银弹”,只做“够用”

很多团队一上来就想建 Feature Store、Model Registry、Experiment Platform 三大件,结果半年过去,Feature Store 只存了 3 个特征,Model Registry 里全是过期模型,Experiment Platform 因权限问题没人敢用。Part 4 的务实原则是: 用最小可行模块(MVP Module)解决最痛的单点问题,再逐步编织成网。 我们的真实演进路径是:

  • 第 1 个月 :只做一件事——在模型服务入口加一层“请求镜像中间件”,把所有入参和出参(含 timestamp、request_id、model_version)以 Parquet 格式写入 S3,每天自动分区。成本:2 人日,收益:首次获得可回溯的线上预测数据,支撑了 80% 的模型诊断需求。

  • 第 3 个月 :基于镜像数据,用 PySpark 写一个离线脚本,每小时计算各特征的空值率、分布偏移(KS)、值域异常(IQR 法),结果写入 MySQL 表,Grafana 直连展示。成本:3 人日,收益:提前 2 天发现天气特征数据源中断,避免模型性能下滑。

  • 第 6 个月 :把离线脚本改造成实时 Flink 作业,接入 Kafka,实现秒级漂移检测;同时将告警规则从 Grafana 迁移到 Alertmanager,支持多通道通知和静默期。成本:5 人日,收益:数据漂移平均响应时间从 4 小时缩短至 92 秒。

  • 第 12 个月 :才引入 Feast 作为 Feature Store,但只接入最关键的 12 个实时特征(如用户实时点击序列、库存变化),其余 200+ 批处理特征仍走 Hive。成本:15 人日,收益:新模型接入特征开发时间从 5 天降至 4 小时。

注意:永远先问“这个问题不解决,业务会死吗?”——如果答案是否定的,那就暂缓。我们曾砍掉一个“模型血缘追踪”需求,因为当时 90% 的模型变更都由同一人负责,口头沟通比建系统更快。直到团队扩到 12 人、模型数破 50,才重拾这个需求。

3. 核心细节解析与实操要点:监控、漂移、A/B 测试的落地陷阱

3.1 模型监控不是“看图说话”,而是建立可操作的决策树

很多团队的监控 Dashboard 看似华丽:折线图、热力图、仪表盘一应俱全,但当告警响起时,工程师第一反应是“这图说明什么?下一步该查哪?”——这就是监控失效的标志。Part 4 的监控设计,核心是把每个指标映射到明确的、可执行的运维动作。以我们最常遇到的“预测延迟突增”为例,我们构建了如下决策树:

延迟 P95 > 80ms?
├─ 是 → 查看 CPU/内存使用率
│   ├─ CPU > 90%? → 检查是否有慢查询阻塞线程池(查 PostgreSQL pg_stat_activity)
│   │   └─ 发现慢查询 → 通知 DBA 优化,同时临时扩容 CPU
│   └─ 内存 > 85%? → 检查特征缓存命中率
│       └─ 命中率 < 60%? → 调整 Redis 缓存 TTL 或增加缓存 key 粒度
└─ 否 → 检查网络延迟(curl -w "@curl-format.txt" -o /dev/null -s http://model-service/predict)
    ├─ 网络延迟 > 20ms? → 检查 Service Mesh(Istio)Sidecar 日志
    └─ 网络正常? → 检查模型推理耗时(OpenTelemetry trace 中 model_inference_span)
        └─ 模型耗时 > 50ms? → 触发模型性能剖析(Py-Spy 采样)
            └─ 发现 numpy 向量化不足 → 重构特征工程代码

这个决策树不是写在 Wiki 里,而是直接编码进我们的监控告警脚本中。当 Prometheus 触发 model_latency_p95_over_threshold 告警时,Alertmanager 会调用一个 Python webhook,该脚本自动执行上述分支判断,最终生成一条 Slack 消息,内容类似:“⚠️ 延迟告警:P95=112ms。已定位:Redis 缓存命中率 42%(基线 78%)。建议:检查 feature_cache_key 生成逻辑。执行命令: kubectl exec -it redis-pod -- redis-cli info | grep keyspace_hits ”。 工程师拿到的不是问题,而是带着上下文的解决方案。 我们因此将平均故障定位时间(MTTD)从 28 分钟降至 3.7 分钟。

3.2 数据漂移检测:别迷信 KS 检验,要结合业务语义

KS 检验(Kolmogorov-Smirnov Test)是数据漂移检测的标配,但它有个致命缺陷:对长尾分布极不敏感。在电商场景中,用户购买金额(order_amount)天然右偏,KS 检验可能显示 p-value=0.2(无显著漂移),但实际业务中,1000 元以上订单占比从 0.8% 涨到 3.2%,这直接导致模型对高价值用户的预测偏差放大。Part 4 的实操方案是“三叉戟”检测法:

  • 统计检验(Statistical Test) :对连续变量用 KS 检验(p-value < 0.01),对离散变量用卡方检验(χ² p-value < 0.05)。这是基础门槛,但仅作初筛。

  • 业务阈值(Business Threshold) :为关键特征设定硬性业务规则。例如:

    • user_age :必须 ∈ [16, 100],否则告警;
    • click_through_rate :7 天滑动窗口均值波动 > ±15% 且持续 3 小时,告警;
    • device_os_version :iOS/Android 主版本号占比突变 > 20%,告警(预示新系统上线影响)。
  • 模型敏感度(Model Sensitivity) :用 SHAP 值或 LIME 解释单个预测,观察当某特征值变化时,预测结果的变化幅度。我们训练了一个轻量级“敏感度代理模型”(Sensitivity Proxy Model),输入是特征向量,输出是该样本对各特征的梯度模长。当某特征的梯度模长均值突增 3 倍,即视为该特征对模型输出变得异常敏感,需重点检查其数据质量。在一次 iOS 17 升级后, device_os_version 的敏感度模长从 0.12 涨到 0.45,我们立即发现模型对新系统用户的预测置信度普遍偏低,及时触发了特征工程优化。

实操心得:漂移检测的报警阈值绝不能拍脑袋定。我们采用“滚动基线法”:每天用过去 14 天的数据计算各指标的均值 μ 和标准差 σ,告警阈值设为 μ ± 3σ。这样既能适应业务自然波动(如双十一大促期间 CTR 本就会升高),又能捕捉异常突变。上线后误报率下降 68%。

3.3 A/B 测试框架:流量分桶不是随机,而是可控的“业务实验”

很多团队的 A/B 测试就是“50% 流量给 A,50% 给 B”,结果发现 A 组转化率高,B 组 GMV 高,老板问“到底该用哪个?”,没人答得出来。Part 4 的 A/B 框架核心是 Multi-Armed Bandit(MAB)+ 业务目标加权 。我们不追求统计学上的“绝对最优”,而是追求业务上的“综合收益最大”。

  • 流量分桶策略 :放弃全局随机,采用“分层哈希”。以电商为例,我们将用户按 user_id % 100 分成 100 个桶,但每个桶内再按 user_segment (新客/老客/高价值/低价值)和 device_type (iOS/Android/Web)做二次哈希,确保每个实验组在各用户群体中分布均衡。这样,当测试“新推荐算法”时,我们能单独分析“对 iOS 新客的提升效果”,而不是被整体数据淹没。

  • 评估指标体系 :强制定义三层指标:

    • 北极星指标(North Star) :直接反映商业目标,如“GMV”、“LTV/CAC”;
    • 护栏指标(Guardrail) :必须守住的底线,如“用户投诉率 < 0.01%”、“服务器错误率 < 0.05%”;
    • 探索指标(Exploratory) :辅助归因,如“长尾商品曝光占比”、“跨品类点击率”。
  • 决策引擎 :我们用 Thompson Sampling 算法动态调整流量分配。初始各组 25% 流量,每小时根据各组的北极星指标(GMV)和护栏指标(投诉率)计算“综合得分”:
    Score = GMV_per_user * 0.7 - Complaints_per_1000_users * 100
    得分高的组自动获得更高流量权重。在一次搜索排序实验中,算法在 36 小时内将最优组流量从 25% 提升至 68%,最终 GMV 提升 4.2%,而投诉率未超阈值。相比传统固定分流,实验周期缩短 55%。

注意:A/B 测试必须隔离“技术变更”和“业务变更”。我们严格规定:任何模型版本升级,必须搭配一个“Baseline”实验组(即旧模型),且 Baseline 组的代码、配置、依赖必须与上线前完全一致。曾有一次,新模型上线后 GMV 下降,排查发现是 CDN 缓存配置变更导致页面加载变慢,而非模型问题——正因 Baseline 组也受此影响,数据对比才暴露了真因。

4. 实操过程与核心环节实现:从零搭建可观测性与自动化响应

4.1 第一步:构建“请求-特征-预测”全链路追踪(OpenTelemetry 实战)

可观测性的基石是追踪(Tracing),没有追踪,监控就是盲人摸象。我们选择 OpenTelemetry(OTel)而非 Zipkin 或 Jaeger,因为它原生支持多语言、云原生集成,并且 Trace 数据可直接导出为 Prometheus metrics。以下是 Python 模型服务的 OTel 集成实操步骤(基于 FastAPI):

Step 1:安装与初始化

pip install opentelemetry-api opentelemetry-sdk opentelemetry-exporter-otlp-proto-http opentelemetry-instrumentation-fastapi

Step 2:配置 OTel Collector(Docker Compose)

# otel-collector-config.yaml
receivers:
  otlp:
    protocols:
      http:

exporters:
  otlphttp:
    endpoint: "http://prometheus:9090"
  logging:

service:
  pipelines:
    traces:
      receivers: [otlp]
      exporters: [otlphttp, logging]

Step 3:在 FastAPI 应用中注入追踪

from fastapi import FastAPI, Request, Response
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPHTTPSpanExporter
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor

# 初始化 tracer
provider = TracerProvider()
processor = BatchSpanProcessor(OTLPHTTPSpanExporter(endpoint="http://otel-collector:4318/v1/traces"))
provider.add_span_processor(processor)
trace.set_tracer_provider(provider)

app = FastAPI()

# 自动注入 FastAPI 追踪
FastAPIInstrumentor.instrument_app(app)

@app.post("/predict")
async def predict(request: Request, response: Response):
    # 1. 创建 span,显式添加业务上下文
    tracer = trace.get_tracer(__name__)
    with tracer.start_as_current_span("model_predict") as span:
        # 添加 request_id(从 header 或生成)
        req_id = request.headers.get("X-Request-ID", str(uuid.uuid4()))
        span.set_attribute("http.request_id", req_id)
        
        # 2. 记录输入特征(脱敏后)
        body = await request.json()
        features = {k: v for k, v in body.items() if k not in ["user_pii", "token"]}
        span.set_attribute("input_features", json.dumps(features))
        
        # 3. 记录模型版本与关键元数据
        span.set_attribute("model.version", "v2.3.1")
        span.set_attribute("model.framework", "sklearn")
        
        # 4. 执行预测(此处省略模型加载逻辑)
        prediction = model.predict([list(features.values())])
        
        # 5. 记录输出与耗时
        span.set_attribute("output.prediction", float(prediction[0]))
        span.set_attribute("output.confidence", float(confidence))
        
        return {"prediction": float(prediction[0]), "confidence": float(confidence)}

关键细节与避坑:

  • Span 命名规范 :不用 predict 这种泛称,而用 model_{model_name}_predict (如 model_eta_v2_predict ),便于在 Jaeger UI 中按模型过滤。
  • 属性(Attribute)精简 :不记录原始请求体(隐私与体积),只记录脱敏后的特征名和值类型(如 "age": "int" ),具体值用 set_attribute("feature.age.value", age) 单独记录。
  • 错误捕获 :在 try/except 中,对 span.record_exception(e) ,并设置 span.set_status(Status(StatusCode.ERROR)) ,这样错误率指标可自动聚合。
  • 采样率控制 :生产环境不全量采样,用 ParentBased(trace_id_ratio_sample_rate=0.01) 设置 1% 采样率,既保关键链路,又控开销。

部署后,在 Jaeger UI 中输入 model_eta_v2_predict ,即可看到完整调用链: Nginx → API Gateway → Feature Service → Model Service → Cache ,每个环节的耗时、状态、错误信息一目了然。当某次延迟突增时,我们直接定位到“Feature Service 调用外部天气 API 超时”,而非在模型服务里大海捞针。

4.2 第二步:数据漂移检测流水线(Evidently + Airflow)

Evidently 是目前最轻量、最易集成的漂移检测库,它不依赖数据库,直接读取 Pandas DataFrame 即可生成 HTML 报告或 JSON 指标。我们将其嵌入 Airflow DAG,实现每小时自动检测。

Step 1:准备数据

# airflow_dag.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
import pandas as pd
from evidently.report import Report
from evidently.metrics import DataDriftTable, DatasetSummaryMetric

def run_drift_detection(**context):
    # 读取昨日数据(生产数据湖路径)
    ref_df = pd.read_parquet("s3://data-lake/features/ref/2024-05-20/")
    cur_df = pd.read_parquet("s3://data-lake/features/cur/2024-05-21/")
    
    # 构建 Evidently Report
    report = Report(metrics=[
        DataDriftTable(), 
        DatasetSummaryMetric()
    ])
    report.run(reference_data=ref_df, current_data=cur_df)
    
    # 导出为 JSON(供告警脚本消费)
    result_json = report.as_dict()
    
    # 提取关键漂移指标
    drift_metrics = {}
    for metric in result_json["metrics"]:
        if metric["metric"] == "DataDriftTable":
            for col in metric["result"]["drift_by_columns"].values():
                if col["drift_detected"]:
                    drift_metrics[col["column_name"]] = {
                        "p_value": col["p_value"],
                        "distance": col["distance"]
                    }
    
    # 写入 S3 供告警
    s3_client.put_object(
        Bucket="ml-monitoring",
        Key=f"drift_reports/{datetime.now().strftime('%Y%m%d_%H')}.json",
        Body=json.dumps(drift_metrics)
    )

Step 2:Airflow DAG 定义

default_args = {
    'owner': 'ml-team',
    'depends_on_past': False,
    'start_date': datetime(2024, 5, 1),
    'email_on_failure': True,
    'retries': 1,
}

dag = DAG(
    'feature_drift_detection',
    default_args=default_args,
    description='Hourly data drift detection',
    schedule_interval='0 * * * *',  # 每小时执行
    catchup=False,
)

drift_task = PythonOperator(
    task_id='run_drift_detection',
    python_callable=run_drift_detection,
    dag=dag,
)

alert_task = PythonOperator(
    task_id='send_drift_alert',
    python_callable=send_alert_if_drift,
    op_kwargs={'threshold_pvalue': 0.001},
    dag=dag,
)

drift_task >> alert_task

Step 3:告警脚本(send_alert_if_drift)

def send_alert_if_drift(**context):
    # 从 S3 读取最新漂移报告
    obj = s3_client.get_object(Bucket="ml-monitoring", Key="drift_reports/latest.json")
    drift_data = json.loads(obj['Body'].read())
    
    # 检查关键特征
    critical_features = ["user_age", "order_amount", "click_count_7d"]
    alerts = []
    for feat in critical_features:
        if feat in drift_data and drift_data[feat]["p_value"] < 0.001:
            alerts.append(f"{feat} drift detected! p-value={drift_data[feat]['p_value']:.4f}")
    
    if alerts:
        # 发送 Slack 告警,带跳转链接到 Evidently HTML 报告
        slack_msg = {
            "text": "🚨 CRITICAL DATA DRIFT DETECTED",
            "blocks": [
                {"type": "section", "text": {"type": "mrkdwn", "text": "\n".join(alerts)}},
                {"type": "actions", "elements": [
                    {"type": "button", "text": {"type": "plain_text", "text": "View Full Report"}, 
                     "url": "https://evidently-reports.s3.amazonaws.com/latest.html"}
                ]}
            ]
        }
        requests.post(SLACK_WEBHOOK, json=slack_msg)

实操心得:

  • 数据采样 ref_df 用全量历史数据(如过去 30 天), cur_df 用最近 1 小时数据。Evidently 对样本量敏感, cur_df 至少 1000 行,否则漂移检测不可靠。
  • 特征筛选 :不检测所有特征,只检测 model.feature_names_ 中的 20 个核心特征。我们维护一个 critical_features.yaml 文件,由数据科学家和业务方共同确认。
  • 报告归档 :每次生成的 HTML 报告存入 S3,并用 s3_website 静态托管,URL 形如 https://evidently-reports.s3.amazonaws.com/20240521_14.html ,方便长期追溯。

4.3 第三步:自动化响应与熔断(Envoy + Custom Filter)

当监控发现严重问题时,“告警”只是开始,“自动响应”才是闭环。我们基于 Envoy Proxy 构建了模型服务的智能网关,通过自定义 WASM Filter 实现熔断与降级。

Step 1:编写 WASM Filter(Rust)

// src/lib.rs
use proxy_wasm::traits::*;
use proxy_wasm::types::*;

#[no_mangle]
pub fn _start() {
    proxy_wasm::set_log_level(LogLevel::Trace);
    proxy_wasm::set_root_context(|_| -> Box<dyn RootContext> {
        Box::new(HealthCheckRoot {})
    });
}

struct HealthCheckRoot {}

impl Context for HealthCheckRoot {}
impl RootContext for HealthCheckRoot {
    fn on_configure(&mut self, _plugin_configuration_size: usize) -> bool {
        true
    }
}

impl StreamContext for HealthCheckFilter {
    fn on_http_request_headers(&mut self, _num_headers: usize) -> Action {
        // 1. 从 Prometheus 查询当前模型错误率
        let error_rate = self.query_prometheus("rate(model_errors_total[5m])");
        
        // 2. 如果错误率 > 0.1%,触发熔断
        if error_rate > 0.1 {
            // 3. 重写请求头,指向降级服务
            self.set_http_request_header("x-fallback-service", "static-eta-service");
            return Action::Continue;
        }
        
        // 4. 正常转发
        Action::Continue
    }
}

Step 2:Envoy 配置

# envoy.yaml
static_resources:
  listeners:
  - name: model-listener
    address:
      socket_address: { address: 0.0.0.0, port_value: 8000 }
    filter_chains:
    - filters:
      - name: envoy.filters.network.http_connection_manager
        typed_config:
          "@type": type.googleapis.com/envoy.extensions.filters.network.http_connection_manager.v3.HttpConnectionManager
          route_config:
            name: local_route
            virtual_hosts:
            - name: backend
              domains: ["*"]
              routes:
              - match: { prefix: "/predict" }
                route: { cluster: model-cluster }
          http_filters:
          - name: envoy.filters.http.wasm
            typed_config:
              "@type": type.googleapis.com/envoy.extensions.filters.http.wasm.v3.Wasm
              config:
                root_id: "health-check-root"
                vm_config:
                  runtime: "envoy.wasm.runtime.v8"
                  code:
                    local:
                      filename: "/etc/envoy/filters/health_check.wasm"
                  allow_precompiled: true

Step 3:Prometheus 查询函数

def query_prometheus(query):
    url = f"http://prometheus:9090/api/v1/query?query={query}"
    try:
        res = requests.get(url, timeout=2)
        if res.status_code == 200:
            data = res.json()
            if data["data"]["result"]:
                return float(data["data"]["result"][0]["value"][1])
    except Exception as e:
        proxy_wasm.log_warn(f"Prometheus query failed: {e}")
    return 0.0

部署后,当模型服务错误率飙升时,Envoy 在 200ms 内完成检测、重写 Header、转发至降级服务,用户无感知。我们用这种方式,在一次 Kubernetes 节点故障中,将服务可用性从 92.3% 提升至 99.99%。

5. 常见问题与排查技巧实录:来自产线的 12 个血泪教训

5.1 “模型在测试环境跑得飞快,上线就超时”——内存泄漏的隐形杀手

现象 :模型服务 P95 延迟从 15ms 涨到 500ms,CPU 使用率正常,内存使用率缓慢爬升,重启后恢复,几小时后复发。

排查过程

  • psutil 在服务中添加内存监控: psutil.Process().memory_info().rss ,每分钟打点。
  • 发现内存占用呈线性增长,斜率约 2MB/分钟。
  • tracemalloc 追踪内存分配: tracemalloc.start(); time.sleep(300); snapshot = tracemalloc.take_snapshot()
  • 分析 snapshot,发现 pandas.core.internals.managers.BlockManager 占用 85% 内存,源头是 pd.concat() 被反复调用。

根因与修复

  • 问题代码:每次预测都 df = pd.concat([df, new_row]) 动态追加,而 concat 会创建新 DataFrame,旧对象未被及时 GC。
  • 修复方案:改用预分配数组 + pd.DataFrame.from_records() 。对于实时预测,我们维护一个长度为 1000 的环形缓冲区,满则批量处理,避免高频小对象分配。
  • 经验 :Python 服务内存泄漏 70% 源于 pandas 操作不当。上线前必做 tracemalloc 压测。

5.2 “A/B 测试结果矛盾:A 组点击率高,B 组转化率高”——辛普森悖论的业务陷阱

现象 :新推荐算法 A 组 CTR +5%,但 GMV -2%;旧算法 B 组 CTR -3%,GMV +1%。业务方无法决策。

排查过程

  • 拆解用户分层:发现 A 组在“新客”中 CTR +12%,但“老客”中 CTR -1%;B 组相反。
  • 进一步分析:A 组新客 GMV 提升,但老客因推荐了更多低价商品

更多推荐