特征存储(Feature Store)架构:从一致性到在线检索的完整实践

深入解析特征存储(Feature Store)的架构与实践:特征平台核心概念、训练服务一致性(Training-Serving Skew)、特征注册与管理(FeatureView/版本)、在线离线存储(Redis/向量检索)、批流特征计算管道、Feast 与定制方案对比、特征血缘与监控,以及推荐系统端的完整实战案例。

引言

在机器学习生产中,工程师最常踩的坑不是模型精度,而是"训练和推理用的特征不一致"。训练时从离线表取历史特征,上线时从在线服务实时算特征,两者往往相差甚远——这就是著名的 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 架构的行业分析

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 反向 ETL 与数据激活:让数据仓库的价值回到业务系统
  2. 数据可观测性:从管道监控到数据宕机的全方位保障
  3. 湖仓一体架构:Iceberg、Delta Lake 与 Hudi 的统一数据底座