AI 工程化的基础设施挑战
当企业的 AI 应用从个位数增长到数百个,从实验性项目演变为核心业务系统,基础设施的短板会迅速暴露。数据科学家在本地训练好的模型,到了生产环境发现特征计算逻辑不一致——训练时用的是实时聚合的点击次数,推理时用的是 T+1 离线批处理的结果,模型效果断崖式下跌。
这种「训练-服务偏移」(Training-Serving Skew)是 ML 系统最常见的故障模式之一。根因在于特征工程逻辑在训练管道和推理服务中被重复实现了两遍:Python Pandas 脚本用于离线训练,Java 服务用于在线推理,两者的时区处理、空值填充、分位数计算细节不可避免地产生分歧。
AI 工程化基础设施的目标正是系统性解决这类问题:通过特征平台统一离线和在线特征计算、通过模型注册中心管控版本血缘、通过推理流水线编排从数据到预测的完整链路。这三大支柱构成了企业级 MLOps 的底座。
特征平台:训练-服务一致性的关键
特征计算的双轨问题
传统的 ML 特征工程通常有两条独立的代码路径:
| 维度 | 离线训练 | 在线推理 |
|---|---|---|
| 计算引擎 | Spark / Pandas | Python / Java 服务 |
| 调度方式 | Airflow 定时批处理 | 实时 API 调用 |
| 数据源 | Hive / Data Lake | Redis / MySQL / Kafka |
| 时间窗口 | 回溯历史范围 | 当前时间点 |
| 代码仓库 | feature_engineering/train.py | serving/features.py |
任何一列的差异都可能导致特征值不一致。例如,离线计算用户近 7 天点击次数时包含当天数据,而在线服务由于数据延迟只能统计到昨天,同一个特征在两个环境中含义不同。
特征平台的核心架构
特征平台(Feature Store)将特征定义为统一的逻辑实体,同时支撑离线训练和在线 Serving:
┌─────────────────────────────────────────────────────────────┐
│ Feature Definition │
│ FeatureView(name="user_click_count_7d") │
│ entities=[User] │
│ ttl=timedelta(days=7) │
│ schema={"click_count": Int64} │
│ source=click_stream_table │
│ aggregation=Aggregation(window="7d", function="COUNT") │
└─────────────────────────────────────────────────────────────┘
│ │
▼ ▼
┌──────────────────────┐ ┌──────────────────────┐
│ Offline Store │ │ Online Store │
│ (Parquet / Delta) │ │ (Redis / DynamoDB) │
│ 历史特征值批量生成 │ │ 最新特征值低延迟读取 │
│ 用于模型训练 │ │ 用于实时推理 │
└──────────────────────┘ └──────────────────────┘
特征平台的三个核心能力:
1. 统一特征定义(Feature View)
Feature View 是特征的逻辑描述,声明特征的实体(Entity)、数据源、变换逻辑和 TTL。这是特征平台区别于简单键值存储的关键——特征不是原始数据的被动镜像,而是带有语义和计算逻辑的实体。
from feast import Entity, Feature, FeatureView, ValueType
from feast.types import Int64, Float32
from datetime import timedelta
# 定义实体
user = Entity(name="user_id", value_type=ValueType.STRING)
# 定义特征视图
user_stats_view = FeatureView(
name="user_stats_1d",
entities=[user],
ttl=timedelta(days=1),
schema=[
Field(name="click_count", dtype=Int64),
Field(name="purchase_amount", dtype=Float32),
Field(name="last_session_duration", dtype=Float32),
],
source=user_click_source,
online=True,
offline=True,
)
2. 离线与在线存储双写
特征平台在特征值计算完成后,同时写入:
- 离线存储(通常为对象存储上的 Parquet/Delta/Iceberg 格式):用于批量历史回溯训练,支持时间旅行查询(「在 2024-03-15 这一天,用户的特征值是什么?」)
- 在线存储(通常为 Redis/DynamoDB/Redis Cluster):用于毫秒级的推理特征获取,支持实体 ID 的点查询
双写机制确保在线和离线特征值完全一致,因为它们来自同一套计算逻辑。
3. 实时特征流处理
对于需要近实时更新的特征(如最近 5 分钟点击次数),特征平台集成流处理引擎(Flink / Spark Streaming / Kafka Streams),将特征增量计算结果实时推送到在线存储:
# Feast 的推模式(Push Mode)
from feast import PushSource
push_source = PushSource(
name="click_events_push",
schema=["user_id", "click_count", "event_timestamp"]
)
# 实时流水线将 Kafka 消费结果推送
store.push("user_stats_1d", df=streaming_df)
特征平台选型对比
| 特征平台 | 开源 | 实时能力 | 离线存储 | 在线存储 | 企业级特性 |
|---|---|---|---|---|---|
| Feast | 是 | 推模式 | GCS/S3 + Parquet | Redis/SQLite/DynamoDB | 社区活跃,扩展性强 |
| Tecton | 否 | 原生流批一体 | Delta Lake | DynamoDB/Redis | 全托管,企业级 SLA |
| Databricks Feature Store | 否 | 与 Delta Live Tables 集成 | Delta Lake | Databricks 在线存储 | 与 Spark 生态深度整合 |
| AWS SageMaker Feature Store | 否 | 流式摄取 | S3 + Glue | DynamoDB / In-Memory | 与 AWS ML 栈无缝集成 |
| Vertex AI Feature Store | 否 | 流式摄取 | BigQuery / GCS | Bigtable / Online Serving | 与 GCP 生态集成 |
选型建议:技术实力强且追求灵活性的团队选择 Feast 自建;需要快速上线且预算充足的企业选择 Tecton 或云厂商托管方案。关键在于确保特征平台与企业现有数据栈(数据湖、消息队列、调度系统)的兼容性。
模型注册中心与治理
模型版本管理的必要性
生产环境中同时运行着数十甚至上百个模型版本:推荐系统的召回模型 v3.2、精排模型 v4.1-beta、风控模型的 A/B 实验版本、以及上一个稳定版本作为热备用。没有集中管理的模型注册中心,团队将陷入「哪个版本部署在哪个环境」的混乱。
模型注册中心(Model Registry)的核心功能:
版本追踪:记录每次训练生成的模型版本,关联训练数据版本、代码 Git Commit、超参数配置和评估指标。
import mlflow
with mlflow.start_run():
mlflow.log_param("learning_rate", 0.001)
mlflow.log_param("batch_size", 256)
mlflow.log_metric("val_auc", 0.89)
mlflow.sklearn.log_model(model, "model")
# 注册到 Model Registry
model_version = mlflow.register_model(
model_uri=f"runs:/{mlflow.active_run().info.run_id}/model",
name="recommendation_ranking_model"
)
状态流转:模型版本经历「Staging → Production → Archived」的生命周期。只有标记为 Production 的版本才能被推理服务加载,Staging 版本用于预发环境测试,Archived 版本保留供审计但不可服务。
血缘追踪:追溯模型的上游依赖——训练数据来自哪个 Hive 表、经过哪些特征变换、使用哪份代码的哪个 Commit。当上游数据质量出现问题时,可以快速定位受影响的模型版本。
A/B 实验绑定:将模型版本与实验 ID 关联,记录每个版本的流量占比、业务指标表现和统计显著性。
模型部署模式
| 部署模式 | 说明 | 适用场景 |
|---|---|---|
| 影子模式(Shadow Mode) | 新模型接收生产流量但不影响实际决策,输出仅用于对比评估 | 高风险模型上线前的效果验证 |
| 金丝雀发布(Canary) | 新模型承接 5%~10% 流量,观察指标后逐步扩大 | 常规迭代上线 |
| 蓝绿部署(Blue-Green) | 两个完全对等的部署集群,瞬间切换流量 | 需要零停机回滚的关键业务 |
| 多臂老虎机(MAB) | 动态分配流量至多个模型,根据实时效果自动调整比例 | 持续在线自动优化 |
模型治理合规
金融、医疗等监管严格的行业对模型治理有额外要求:
- 可解释性报告:每个 Production 模型需提供特征重要性分析、SHAP 值全局解释和关键决策的个案解释
- 公平性审计:定期检查模型在不同人群(性别、年龄、地域)上的性能差异,确保无歧视性偏向
- 模型备案:记录模型训练数据的来源、数据脱敏方法、模型架构和训练流程,满足监管审查要求
推理流水线编排
从训练到推理的完整链路
一个完整的 ML 推理流程通常包含多个步骤,需要可靠地编排:
原始事件 → 特征计算 → 特征获取 → 模型推理 → 后处理 → 业务 Action
│ │ │ │ │ │
│ 实时聚合 特征平台 模型服务 阈值判断 写入 DB
│ (Flink) (Feast) (Triton) (规则引擎) (发送消息)
│
└── Kafka / Pulsar 消息流
手动编排这些步骤极易出错:特征计算超时怎么办、模型推理返回格式异常怎么办、后处理阈值调整后如何热更新。推理流水线框架(Kubeflow Pipelines、Apache Airflow、Prefect)提供了声明式的编排能力。
Kubeflow Pipelines 实战
Kubeflow Pipelines(KFP)是 Kubernetes 原生的 ML 工作流引擎,适用于需要弹性扩缩容和混合负载(CPU 预处理 + GPU 推理)的场景。
from kfp import dsl
from kfp.client import Client
@dsl.component(base_image="python:3.11")
def fetch_features(user_id: str) -> dict:
import feast
store = feast.FeatureStore(repo_path=".")
features = store.get_online_features(
features=["user_stats_1d:click_count", "user_stats_1d:purchase_amount"],
entity_rows=[{"user_id": user_id}]
).to_dict()
return features
@dsl.component(base_image="nvcr.io/nvidia/tritonserver:23.10-py3")
def model_inference(features: dict, model_name: str) -> dict:
import tritonclient.http as httpclient
client = httpclient.InferenceServerClient(url="triton:8000")
# 构造输入张量,执行推理
return {"score": 0.85, "label": "high_value"}
@dsl.component(base_image="python:3.11")
def business_action(prediction: dict) -> str:
if prediction["score"] > 0.8:
return "send_push_notification"
return "no_action"
@dsl.pipeline(name="recommendation_pipeline")
def recommendation_pipeline(user_id: str):
features = fetch_features(user_id=user_id)
prediction = model_inference(features=features.output, model_name="ranking_v3")
action = business_action(prediction=prediction.output)
# 提交执行
client = Client(host="kubeflow-pipeline-ui.example.com")
client.create_run_from_pipeline_func(recommendation_pipeline, arguments={"user_id": "user_12345"})
KFP 的优势在于:
- Kubernetes 原生:自动调度 Pod 到合适的节点(GPU / CPU),利用 K8s 的弹性扩缩容
- artifact 追踪:每个组件的输入输出自动存入 MinIO/S3,支持完整的数据血缘追踪
- 缓存复用:若上游数据和代码未变更,自动复用上一次的输出结果,节省重复计算
- 可视化:Pipeline UI 清晰展示每个步骤的执行状态、耗时和日志
实时推理流水线设计
对于延迟敏感的场景(如实时推荐、 fraud detection),需要构建基于事件驱动的推理流水线:
事件溯源架构:
- 用户行为事件(点击、加购、搜索)写入 Kafka
- Flink 实时聚合出特征,推送到在线特征存储(Redis)
- 特征更新触发推理请求,调用 Triton 模型服务
- 推理结果写入事件流,推送给推荐服务或风控引擎
性能优化要点:
- 特征预聚合:将小时级、日级的聚合特征预计算到特征平台,在线只做点查询
- 请求合并:对于同一用户的多个特征请求,在 API Gateway 层合并为单次 batch 请求
- 异步非阻塞:推理服务使用异步框架(FastAPI + asyncio、Tornado),避免线程阻塞等待模型返回
- 本地缓存:对热点用户的特征和推理结果做进程内缓存(如 LRU Cache),命中缓存时延迟降至亚毫秒级
数据血缘与影响分析
当上游数据表结构发生变更时,哪些特征、哪些模型、哪些业务服务会受到影响?数据血缘(Data Lineage)回答了这个问题。
血缘追踪三层
基础设施层血缘:Apache Atlas、DataHub、OpenLineage 追踪数据在存储系统间的流动路径——Hive 表 → Spark Job → Parquet 文件 → 特征平台 → 模型训练。
特征层血缘:特征平台记录每个特征的计算逻辑和上下游依赖。当原始数据表(如 user_click_log)发生 schema 变更时,特征平台自动标记所有依赖该表的特征为「待审核」状态。
模型层血缘:模型注册中心记录模型的输入特征列表、训练数据集版本和评估结果。当某个特征被标记为「数据漂移」或「下线」时,自动通知所有依赖该特征的模型 owner。
影响分析的实践价值
某电商公司每周发布上百次特征更新。上线前的影响分析显示:「用户近 7 天订单金额」的特征更新将影响 12 个生产模型、3 个 A/B 实验和 2 个下游报表。基于这个分析,数据团队将更新时间调整到低峰时段,并通知相关业务方做好监控准备。
总结
AI 工程化基础设施不是奢侈品,而是规模化部署 AI 应用的必需品。特征平台解决了训练-服务一致性的顽疾、模型注册中心提供了版本治理的秩序、推理流水线编排实现了从数据到决策的自动化流转。
构建这套基础设施的关键在于:
- 从第一天起统一特征定义:避免训练和服务两套特征代码的割裂
- 模型版本即产物:将模型视为一级软件制品,经历完整的 Dev(MLOps) 生命周期
- 流水线即代码:用声明式配置定义推理流程,支持版本控制、代码审查和自动化测试
- 血缘即安全网:任何数据或模型的变更必须通过影响分析,确保下游服务不受意外冲击
当特征平台、模型注册中心和推理流水线三位一体地运转时,AI 团队才能真正从重复的工程救火中解放出来,专注于算法创新和业务价值的创造。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。