离线特征存储的设计方案:Feast与离线Parquet的工程取舍
一、特征存储的工程定位
特征存储(Feature Store)是ML基础设施中较晚被标准化的组件。在2017年Uber发布Michelangelo的Feature Store概念之前,大多数团队的特征管理处于"临时脚本+CSV文件"的状态。特征存储的核心价值主张是:统一离线训练和在线推理的特征定义,避免同一个特征在两个环境中被不同代码实现而导致训练-推理偏差。
但"是否需要引入一个专门的特征存储系统(如Feast)"的答案并非总是"Yes"。对于数据规模在GB级别、特征数量在100以内的小型ML项目,一套结构化的Parquet文件系统可能比部署Feast更经济。两者的取舍涉及数据规模、团队能力、延迟SLA和运维成本等多个维度的考量。
二、离线Parquet方案:极简主义的设计与实现
离线Parquet方案的核心思想是:用文件系统的目录结构和Parquet的列式存储来组织特征数据,用YAML/JSON文件管理特征元数据,用Redis/LevelDB作为在线特征缓存层。
这一方案的优势在于零运维成本(不需要部署和维护额外的服务)、与现有数据栈的无缝对接(Spark、pandas、Polars都原生支持Parquet)、极低的入门门槛(任何熟悉文件系统的开发者都能理解数据组织方式)。
典型的数据组织方式是按日期分区的目录结构:
features/ ├── metadata/ │ ├── feature_registry.yaml # 特征注册表 │ └── feature_sets/ │ ├── user_features.yaml # 用户特征集定义 │ └── item_features.yaml # 商品特征集定义 ├── offline/ │ ├── user_features/ │ │ ├── dt=2024-01-01/ │ │ │ └── part-00000.parquet │ │ └── dt=2024-01-02/ │ │ └── part-00000.parquet │ └── item_features/ │ └── ... └── online_export/ ├── user_features_latest.parquet └── item_features_latest.parquet""" 离线Parquet特征存储的轻量级实现:特征读取与在线同步 """ import os import yaml import pandas as pd import redis from pathlib import Path from datetime import datetime, timedelta from typing import Optional, List class LightweightFeatureStore: """基于Parquet和Redis的轻量级特征存储。 适用于特征数量有限(<500)、团队规模较小(<10人)的场景。 核心特点:零额外服务部署、文件系统即存储层、Redis仅作缓存。 """ def __init__( self, feature_root: str, redis_host: str = "localhost", redis_port: int = 6379, redis_prefix: str = "fs:", online_ttl_seconds: int = 86400 # 在线特征默认24小时过期 ): """ Args: feature_root: 特征数据的根目录 redis_host: Redis主机地址 redis_port: Redis端口 redis_prefix: Redis键前缀(用于命名空间隔离) online_ttl_seconds: 在线特征的TTL(秒) """ self.root = Path(feature_root) self.redis_client = redis.Redis( host=redis_host, port=redis_port, decode_responses=False ) self.redis_prefix = redis_prefix self.online_ttl = online_ttl_seconds # 加载特征元数据 self.metadata = self._load_metadata() def _load_metadata(self) -> dict: """加载特征注册表元数据。 Returns: dict: 特征集的定义信息 """ registry_path = self.root / "metadata" / "feature_registry.yaml" if not registry_path.exists(): raise FileNotFoundError(f"特征注册表不存在: {registry_path}") with open(registry_path) as f: registry = yaml.safe_load(f) # 加载每个特征集的详细定义 for fs_name in registry.get("feature_sets", []): fs_path = self.root / "metadata" / "feature_sets" / f"{fs_name}.yaml" if fs_path.exists(): with open(fs_path) as f: registry["feature_sets_detail"] = registry.get( "feature_sets_detail", {} ) registry["feature_sets_detail"][fs_name] = yaml.safe_load(f) return registry def get_offline_features( self, feature_set: str, entity_ids: Optional[List[str]] = None, date_range: Optional[tuple[str, str]] = None, ) -> pd.DataFrame: """读取离线特征(用于模型训练)。 按日期分区读取Parquet文件,可选择按实体ID和时间范围过滤。 基于Parquet的谓词下推,只读取需要的行列。 Args: feature_set: 特征集名称(如 "user_features") entity_ids: 要读取的实体ID列表(None表示全部) date_range: 日期范围 (start_date, end_date) Returns: pd.DataFrame: 特征数据 """ feature_path = self.root / "offline" / feature_set if not feature_path.exists(): raise ValueError(f"特征集路径不存在: {feature_path}") # 使用Parquet的分区过滤功能(谓词下推) # pandas的read_parquet支持filters参数直接在文件层面过滤 filters = [] if date_range: start, end = date_range filters.append(("dt", ">=", start)) filters.append(("dt", "<=", end)) if entity_ids: # 对于实体ID过滤,先用pyarrow的dataset API # 它可以利用Parquet的row group统计信息跳过不相关文件 import pyarrow.dataset as ds dataset = ds.dataset(feature_path, format="parquet", partitioning="hive") # 构建过滤器表达式 import pyarrow.compute as pc expr = pc.field("entity_id").isin(entity_ids) if date_range: expr = expr & ( (pc.field("dt") >= date_range[0]) & (pc.field("dt") <= date_range[1]) ) table = dataset.to_table(filter=expr) return table.to_pandas() # 简单情况:使用pandas直接读取(利用分区过滤) return pd.read_parquet(feature_path, filters=filters if filters else None) def sync_to_online( self, feature_set: str, entity_ids: Optional[List[str]] = None, ) -> int: """将最新的离线特征同步到Redis在线存储。 策略:读取当日最新的Parquet分区,逐条写入Redis(Hash结构)。 适用于T+1更新的场景(每天批量同步一次)。 Args: feature_set: 特征集名称 entity_ids: 要同步的实体ID(None表示全量) Returns: int: 成功写入的实体数量 """ # 使用今天的日期作为最新分区 today = datetime.now().strftime("%Y-%m-%d") try: df = self.get_offline_features( feature_set, entity_ids=entity_ids, date_range=(today, today) ) except Exception: # 如果今日数据尚未生成,回退到昨日 yesterday = (datetime.now() - timedelta(days=1)).strftime("%Y-%m-%d") df = self.get_offline_features( feature_set, entity_ids=entity_ids, date_range=(yesterday, yesterday) ) if df.empty: return 0 count = 0 # 使用pipeline批量写入以降低网络往返次数 pipe = self.redis_client.pipeline() for _, row in df.iterrows(): entity_id = row["entity_id"] key = f"{self.redis_prefix}{feature_set}:{entity_id}" # 将特征值序列化为hash字段 feature_dict = { col: row[col] for col in df.columns if col not in ("entity_id", "dt") } # Redis HSET: 设置hash的多个字段 pipe.hset(key, mapping={ k: str(v).encode() for k, v in feature_dict.items() }) pipe.expire(key, self.online_ttl) count += 1 # 每1000条执行一次(避免pipeline过大) if count % 1000 == 0: pipe.execute() pipe = self.redis_client.pipeline() # 执行剩余的 if count % 1000 != 0: pipe.execute() return count三、Feast:何时值得引入一个特征平台
Feast(Feature Store)是由Google Cloud和Gojek共同维护的开源特征存储。它提供的核心能力超越Parquet方案的地方在于:point-in-time正确性(保证训练数据不会使用未来信息)、在线服务的低延迟(通过gRPC在线服务确保<10ms的特征读取)、特征版本化和回溯(可以复现任意历史时间点的训练数据集)。
但Feast的引入也带来显著的成本:需要部署Feast Server(在线服务)、需要维护Offline Store(BigQuery/Redshift/文件)和Online Store(Redis/Datastore)之间的数据同步、需要学习Feast的概念体系(FeatureView、Entity、FeatureService等)。
决策的关键问题是:你的场景中是否存在"point-in-time join"的刚需?如果特征和标签的时间对齐不是问题(例如所有特征都是T+1生成的静态快照),Parquet方案就能满足需求。只有当特征具有不同的时间戳、需要精确地按事件时间进行join时,Feast的point-in-time正确性保证才成为不可替代的差异化价值。
四、渐进式迁移路径
对于不确定是否需要Feast的团队,一条务实的路径是从Parquet方案开始,在需求触发时渐进迁移:
第一阶段(当前):Parquet + YAML元数据 + Redis缓存。满足所有特征的离线训练和在线服务需求,但缺乏point-in-time join和历史版本回溯。
第二阶段(触发条件:标签泄漏风险):将离线部分的Parquet数据导入Feast的Offline Store,使用Feast的historical retrieval API获取训练数据,但保留自建的Redis在线服务。这是"用Feast做训练数据生成,用自己的Redis做在线推理"的混合架构。
第三阶段(触发条件:在线延迟或特征一致性要求提升):将在线服务迁移到Feast的Online Serving API,统一使用Feast管理整个特征生命周期。
五、总结
离线Parquet方案和Feast特征平台不是替代关系,而是特征存储成熟度谱系上的两个节点。Parquet方案以文件系统和YAML元数据实现了特征存储的最核心需求——统一特征定义、消除离在线不一致——同时保持了零运维成本的优势。Feast在point-in-time正确性、在线服务低延迟和历史版本管理上提供了企业级保障,但以引入额外的服务组件和概念复杂性为代价。对于大多数特征数量有限、数据T+1更新、团队规模不大的ML项目,从Parquet方案起步然后在需求明确时迁移到Feast,是一条比"一开始就上Feast"更经济的实践路径。