引言
在机器学习生产中,工程师最常踩的坑不是模型精度,而是"训练和推理用的特征不一致"。训练时从离线表取历史特征,上线时从在线服务实时算特征,两者往往相差甚远——这就是著名的 Training-Serving Skew(训练服务偏差)。特征存储(Feature Store)正是为解决这类问题而生的基础设施:它统一管理特征的定义、计算、存储、检索与血缘,让特征既能批量生成训练样本,也能低延迟提供在线推理。本文将系统拆解特征平台的架构设计、开源方案 Feast、一致性保证与生产落地。
特征存储的本质是把"特征"从脚本里的临时变量,升级为与数据表同等重要的、可复用、可治理的一等资产。
一、特征平台核心概念
1.1 特征平台解决的问题
| 问题 | 传统做法 | 特征平台方案 |
|---|---|---|
| 特征重复造轮子 | 每个团队各自算 | 特征注册中心共享复用 |
| 训练/推理不一致 | 手工同步逻辑 | 同一份定义双写 |
| 特征不可发现 | 散落在 notebooks | 目录 + 血缘 |
| 特征不可回溯 | 无版本记录 | 版本化 + 时间旅行 |
| 在线延迟 | 全量重算 | 在线存储预计算 |
1.2 核心组件
一个特征平台通常由四层组成:定义层(FeatureView/Feature 元数据)、计算层(批/流管道)、存储层(离线 + 在线 + 向量)、服务层(在线检索 API + 训练样本生成)。
┌──────────────────────────────────────────────┐
│ 定义层:FeatureView / Entity / 特征注册 │
├──────────────┬───────────────┬───────────────┤
│ 计算层:批 │ 计算层:流 │ 计算层:Lambda │
│ Spark 日批 │ Flink 实时 │ 批+流对账 │
├──────────────┴───────────────┴───────────────┤
│ 存储层:离线(湖/仓) + 在线(Redis/向量) │
├──────────────────────────────────────────────┤
│ 服务层:训练样本生成 + get_online_features │
└──────────────────────────────────────────────┘
1.3 离线与在线双存储
离线存储(Object Storage / 数仓)面向训练的大规模点查与时间回溯;在线存储(Redis 等)面向推理的毫秒级单条/批量检索。两者通过"同一份定义 + 物化(Materialize)“保持语义一致。
二、在线/离线一致性
2.1 Training-Serving Skew 的成因
| 偏差类型 | 成因 | 示例 |
|---|---|---|
| 时间偏差 | 训练用历史值,在线用当前值 | 用户昨日消费 vs 实时累计 |
| 逻辑偏差 | 训练/在线各写一份计算 | 折扣口径不同 |
| 数据偏差 | 在线缺特征时用了兜底值 | 默认 0 vs 真实空值 |
| 分布偏差 | 训练分布已过时 | 模型老化 |
2.2 一致性保证机制
Feast 等平台用"单一特征定义 + 离线/在线双写"消除逻辑偏差,用 Point-in-Time Correct Join 消除时间偏差。
# feature_definition.py
from feast import Entity, FeatureView, Field, FileSource
from feast.types import Float32, Int64
from datetime import timedelta
customer = Entity(name="customer", join_keys=["customer_id"], value_type=Int64)
customer_stats = FeatureView(
name="customer_stats",
entities=[customer],
ttl=timedelta(days=30), # 在线存储过期时间
schema=[
Field(name="lifetime_value", dtype=Float32),
Field(name="order_count_30d", dtype=Int64),
],
source=FileSource(
path="s3://features/customer_stats.parquet",
timestamp_field="event_ts",
),
online=True, # 同时写入在线存储
)
2.3 Point-in-Time 正确性
训练样本必须"只用当时已知的信息”,否则会引入泄漏。Feast 的 get_historical_features 自动做 Point-in-Time Join。
# training_data.py
from feast import FeatureStore
from datetime import datetime
store = FeatureStore(repo_path="feature_repo")
# 只有 label 是"未来",所有特征都取自 label 之前的最新值
entity_df = store.get_historical_features(
entity_df="SELECT customer_id, order_ts AS event_timestamp, label FROM training_entities",
features=[
"customer_stats:lifetime_value",
"customer_stats:order_count_30d",
"user_profile:tier",
],
).to_df()
三、特征注册与管理
3.1 特征元数据模型
| 概念 | 英文 | 说明 | 类比 |
|---|---|---|---|
| 实体 | Entity | 特征的 join key 维度 | 主键 |
| 特征组 | FeatureGroup | 同一来源的特征集合 | 表 |
| 特征视图 | FeatureView | 特征组的在线/离线双写视图 | 物化视图 |
| 特征版本 | Feature Version | 定义变更的版本化 | Schema 版本 |
| 特征仓库 | Feature Repository | 定义即代码 | 代码仓库 |
3.2 特征注册即代码
特征定义应当纳入版本控制,与应用代码同一流程发布。
# feature_repo/features.py
from feast import Entity, FeatureView, Field
from feast.infra.offline_stores.file_source import FileSource
order = Entity(name="order", join_keys=["order_id"], value_type=str)
user = Entity(name="user", join_keys=["user_id"], value_type=str)
order_features = FeatureView(
name="order_features",
entities=[order],
ttl=timedelta(days=7),
schema=[
Field(name="amount", dtype=Float32),
Field(name="status", dtype=str),
],
source=FileSource(path="s3://features/orders.parquet", timestamp_field="ts"),
)
# 特征组:同一主题的特征放在一起便于管理与权限控制
user_features = FeatureView(
name="user_features",
entities=[user],
ttl=timedelta(days=90),
schema=[
Field(name="tier", dtype=str),
Field(name="churn_score", dtype=Float32),
],
source=FileSource(path="s3://features/users.parquet", timestamp_field="ts"),
)
3.3 注册与版本控制
# 应用特征定义到特征存储(注册到 Registry)
feast apply
# 物化历史特征到在线存储(回填)
feast materialize-incremental 2026-09-01T00:00:00
# 查看已注册特征与版本
feast feature-views list
feast registry-dump | jq '.featureViews[].name'
四、存储与检索
4.1 在线存储选型
| 存储 | 特性 | 延迟 | 适用 |
|---|---|---|---|
| Redis | 内存 KV,特征标准选型 | 毫秒级 | 绝大多数场景 |
| DynamoDB/云 KV | 无服务器、可扩展 | 毫秒级 | 已用云厂商 |
| 向量库(Milvus/FAISS) | 相似度检索 | 毫秒级 | 向量特征(Embedding) |
| 本地内存 | 进程内缓存 | 微秒级 | 高性能局部 |
4.2 Redis 在线存储配置
Feast 通过 feature_store.yaml 声明在线存储,支持 Redis 集群。
# feature_store.yaml
project: rec_features
registry: s3://features/registry.db
provider: aws
online_store:
type: redis
connection_string: redis-cluster.xxxx.ap-southeast-1.amazonaws.com:6379
key_ttl_seconds: 2592000 # 30 天
offline_store:
type: file
path: s3://features/offline
4.3 在线检索 API
推理服务通过 SDK 批量取特征,一次调用返回多实体多特征。
# online_serving.py
from feast import FeatureStore
store = FeatureStore(repo_path="feature_repo")
# 批量在线取特征:推荐候选集评分前统一拉取
features = store.get_online_features(
features=[
"user_features:tier",
"user_features:churn_score",
"order_features:amount_rolling_7d",
],
entity_rows=[
{"user_id": 1001, "order_id": 99881},
{"user_id": 1002, "order_id": 99882},
],
).to_dict()
print(features)
4.4 向量特征检索
对于 Embedding 类特征,接入向量库做相似度检索,为召回阶段提供候选。
# vector_features.py
from pymilvus import Collection, connections
connections.connect(host="milvus", port="19530")
col = Collection("user_embedding_v3")
result = col.search(
data=[query_embedding],
anns_field="embedding",
param={"metric_type": "IP", "params": {"nprobe": 16}},
limit=100,
output_fields=["user_id"],
)
candidate_ids = [hit.entity.get("user_id") for hit in result[0]]
五、特征计算管道:批与流
5.1 批流特征计算
| 维度 | 批量特征 | 流式特征 |
|---|---|---|
| 引擎 | Spark 日批 | Flink 实时 |
| 延迟 | T+1 / 小时级 | 秒-分钟级 |
| 特征类型 | 长期统计、画像 | 实时行为、实时风险 |
| 写入 | 离线存储 + 物化在线 | 直接写在线存储 |
5.2 流式特征计算
用 Flink 计算"近 10 分钟加购次数"等实时特征,并直接写入在线存储。
-- flink_feature.sql
CREATE TABLE cart_events (
user_id BIGINT,
item_id BIGINT,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
) WITH ('connector' = 'kafka', 'topic' = 'cart_events', 'format' = 'json');
CREATE TABLE online_features (
user_id BIGINT PRIMARY KEY NOT ENFORCED,
cart_count_10m BIGINT,
compute_time TIMESTAMP(3)
) WITH ('connector' = 'jdbc', 'url' = 'jdbc:redis://redis:6379');
INSERT INTO online_features
SELECT
user_id,
COUNT(*) AS cart_count_10m,
CURRENT_TIMESTAMP AS compute_time
FROM TABLE(TUMBLE(TABLE cart_events, DESCRIPTOR(event_time), INTERVAL '1' MINUTE))
GROUP BY user_id, TUMBLE(event_time, INTERVAL '1' MINUTE);
5.3 批流对账
流式特征与批量特征必须对账收敛,否则会出现"批量说 5 单、实时说 3 单"的口径冲突。
# reconcile_features.py
def reconcile(user_id: int, batch_cnt: int, stream_cnt: int, tol: float = 0.05):
diff = abs(batch_cnt - stream_cnt) / max(batch_cnt, 1)
if diff > tol:
emit_alert(f"feature skew user={user_id} batch={batch_cnt} stream={stream_cnt}")
return diff
六、Feast 与定制方案对比
6.1 开源 vs 自研
| 方案 | 优点 | 局限 | 适合 |
|---|---|---|---|
| Feast | 开源、定义即代码、与 Airflow/K8s 好集成 | 存储与服务需自管 | 中大型团队自建 |
| Hopsworks/Tecton | 商业化、功能完整 | 成本/锁定 | 预算充足企业 |
| 定制自研 | 完全贴合业务 | 研发与维护成本高 | 特征场景极特殊 |
6.2 Feast 架构要点
Feast 将"定义、Registry、离线/在线存储、服务"分离,服务层通过 gRPC 暴露特征检索接口,可以独立部署。
# feast_serving.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: feast-serving
spec:
replicas: 3
selector:
matchLabels: {app: feast-serving}
template:
metadata:
labels: {app: feast-serving}
spec:
containers:
- name: feast-serving
image: feastdev/feature-server:0.40
args: ["feast", "serve", "-h", "0.0.0.0", "-p", "6566"]
env:
- name: FEATURE_STORE_YAML
value: /mnt/feature_store.yaml
ports:
- containerPort: 6566
6.3 自研的最小闭环
如果选择自研,最小闭环应包含:特征定义注册、离线表、在线 KV、物化任务、检索 API 五件事。
# minimal_fs.py
class MinimalFeatureStore:
def __init__(self, redis, offline):
self.redis, self.offline = redis, offline
def register(self, name: str, sql: str):
self.offline.save_definition(name, sql)
def materialize(self, name: str, ts_from: str):
df = self.offline.run(self.offline.get_definition(name), ts_from)
for row in df.to_dict("records"):
self.redis.set(f"feat:{name}:{row['user_id']}", encode(row))
def get(self, name: str, keys: list):
pipe = self.redis.pipeline()
for k in keys:
pipe.get(f"feat:{name}:{k}")
return [decode(v) for v in pipe.execute()]
七、特征血缘与监控
7.1 特征血缘
特征血缘记录"特征 → 特征视图 → 源表/源流 → 生产该特征的作业",是特征审计与变更影响分析的基础。
{
"feature": "user_features:churn_score",
"version": "v3",
"owner": "growth-ml",
"lineage": {
"upstream": [
{"type": "table", "name": "dws.user_health"},
{"type": "stream", "name": "kafka:user_actions"},
{"type": "job", "name": "spark:churn_score_daily"}
],
"downstream": [
{"type": "model", "name": "churn_prediction_v2"},
{"type": "serving", "name": "rec-service"}
]
},
"owner_of_record": "growth-ml",
"sla": {"freshness_minutes": 30}
}
7.2 特征监控
特征也需要像数据表一样监控分布漂移、缺失率与延迟。
| 指标 | 检测对象 | 告警阈值 |
|---|---|---|
| 分布漂移 | 特征取值分布 vs 训练基线 | PSI > 0.2 |
| 缺失率 | 在线返回 null 比例 | > 1% |
| 新鲜度 | 特征最新时间滞后 | > 30 分钟 |
| 服务延迟 | 在线检索 P99 | > 20ms |
7.3 特征漂移检测
用 PSI(Population Stability Index)检测特征分布漂移,是模型老化的早期信号。
# psi_monitor.py
import numpy as np
def compute_psi(expected: np.ndarray, actual: np.ndarray, buckets: int = 10) -> float:
e_hist, _ = np.histogram(expected, bins=buckets)
a_hist, _ = np.histogram(actual, bins=buckets, range=(expected.min(), expected.max()))
e_ratio = np.clip(e_hist / e_hist.sum(), 1e-6, 1)
a_ratio = np.clip(a_hist / a_hist.sum(), 1e-6, 1)
return float(np.sum((a_ratio - e_ratio) * np.log(a_ratio / e_ratio)))
print("PSI:", round(compute_psi(train_churn_score, online_churn_score), 3))
八、推荐系统实战案例
8.1 案例:某内容平台的推荐特征平台
某内容平台为推荐系统建设特征平台,覆盖 3 亿用户的召回、粗排与精排。
| 阶段 | 动作 | 结果 |
|---|---|---|
| 定义 | 注册 1200+ 特征到 40 个 FeatureView | 特征复用率 60% |
| 双写 | Spark 批特征 + Flink 流特征统一物化 | Skew 归零 |
| 在线 | Redis 集群承载 20 万 QPS 特征检索 | P99 8ms |
| 血缘 | 特征→模型映射,变更自动通知 | 事故率下降 70% |
| 治理 | 特征 Owner + 版本 + 下线审批 | 特征资产化 |
8.2 精排特征管道
精排特征管道串联"实时行为 + 用户画像 + 物品 Embedding",一次推理请求取回全部特征。
# ranking_features.py
from feast import FeatureStore
store = FeatureStore(repo_path="feature_repo")
def get_ranking_features(user_id: int, candidate_items: list) -> dict:
# 用户侧特征 + 物品侧特征 + 交叉特征一次取回
online = store.get_online_features(
features=[
"user_features:recent_cat_weights",
"user_features:avg_ctr_7d",
"item_features:item_embedding",
"item_features:item_cat",
],
entity_rows=[{"user_id": user_id, "item_id": i} for i in candidate_items],
).to_dict()
return {k: online[k] for k in ["item_embedding", "avg_ctr_7d", "item_cat"]}
8.3 实验与回放
特征平台的价值还体现在实验回放:切换特征版本后,可以基于同一批历史样本重放,快速评估新特征对模型的贡献。
# feature_backtest.py
def backtest_feature(store, feature_name: str, entity_df):
# 使用历史时间点重放特征,评估特征有效性
hist = store.get_historical_features(
entity_df=entity_df,
features=[f"{feature_name}"],
).to_df()
return evaluate_importance(hist)
九、常见问题与最佳实践
Q1: 什么时候需要 Feature Store,什么时候不需要?
当出现以下任一信号时,就该引入特征平台:多个模型共享同一批特征、训练与在线特征逻辑难以保持一致、特征不可回溯导致实验无法复现。反之,如果只有一两个模型的少量特征,直接写在管道里更轻量,不必过早引入基础设施。
Q2: 在线存储总是读不到特征怎么办?
这是 Skew 的典型表现,根因通常是"物化未及时跟上"或"TTL 过期"。最佳实践是:为每个特征视图配置合理 TTL 与物化频率,并把"在线缺失率"作为核心 SLO 监控;缺失时用可解释的兜底值(而非默默填 0)并打日志,便于审计与排查。
Q3: 特征血缘要做到什么粒度?
特征级血缘是底线,字段级血缘是加分项。至少要让"某个模型用了哪些特征、特征来自哪张表/哪个流、谁是 Owner"可查询。变更发布前必须走影响分析:特征改动会波及哪些模型与线上服务。
Q4: 批量与流式特征如何对账?
建立"同一口径的双份计算"对账机制:流式特征按小时快照落入离线表,与批量特征按同一口径比较,偏差超过阈值(如 5%)即告警。对账不是最终状态而是一致性的护栏,目的是在 Skew 影响模型前发现它。
总结
| 能力 | 推荐方案 | 关键实践 |
|---|---|---|
| 特征定义 | Feast FeatureView / 定义即代码 | 版本化 + 代码评审 |
| 一致性 | Point-in-Time Join + 双写 | Skew 监控 |
| 在线存储 | Redis + 向量库 | P99 < 20ms |
| 计算管道 | Spark(批) + Flink(流) | 批流对账 |
| 血缘治理 | 特征级血缘 + Owner | 变更影响分析 |
| 质量监控 | PSI + 缺失率 + 新鲜度 | 分布漂移预警 |
特征存储不是"又一个数据平台组件",而是连接数据工程与机器学习工程的枢纽。它的落地成败不取决于技术选型,而取决于三个工程纪律:特征定义单一化、训练服务一致性可验证、特征资产可治理。先从一个业务场景(比如推荐)把"定义-计算-存储-检索-监控"闭环跑通,再横向复制到更多场景,特征平台就会成为组织 ML 能力的真正底座。
参考与延伸阅读
- Feast 官方文档:FeatureStore、FeatureView、在线/离线存储与 gRPC 服务
- Uber Michelangelo / Netflix Feature Store 工程博客中的一致性设计
- 微软与 AWS 关于 Training-Serving Skew 的工程实践
- Tecton 关于特征平台演进与 Lambda 架构的行业分析
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。