主题切换
17.5 特征存储
名词解释:特征存储(Feature Store)
特征存储是机器学习工程中用于集中管理特征定义、计算、存储和服务的系统。在量化交易中,Feature Store 解决了"研究环境中的特征计算逻辑与生产环境中的特征计算逻辑不一致"这一核心痛点。它将特征的创建和使用解耦,确保在线推理时使用的特征值与回测时使用的特征值完全一致,消除训练-服务偏差(Training-Serving Skew)。
一、Feature Store 在量化中的应用
1.1 量化场景下的 Feature Store 价值
在量化交易中,Feature Store 解决以下关键问题:
| 问题 | 传统方式 | 使用 Feature Store |
|---|---|---|
| 训练-服务一致性 | 研究和生产用不同代码计算特征 | 共享特征定义,确保一致性 |
| 特征复用 | 每个策略独立计算相同的特征 | 特征一次定义,多次复用 |
| 时点正确性 | 容易引入前视偏差 | 内置时点查询,自动避免 |
| 特征版本管理 | 特征变更后回测不可复现 | 特征版本化,支持回溯 |
| 在线/离线一致性 | 在线特征可能滞后于离线 | 统一特征注册表保证一致性 |
1.2 量化特征的类型
此处有展示代码展开 ▼
python
from enum import Enum
from dataclasses import dataclass
from typing import List, Dict, Any, Optional
from datetime import datetime, timedelta
import pandas as pd
import numpy as np
class FeatureType(Enum):
"""量化特征的类型分类"""
PRICE_DERIVED = "price_derived" # 基于价格计算:收益率、波动率、动量
FUNDAMENTAL = "fundamental" # 基本面:PE、PB、ROE
ALTERNATIVE = "alternative" # 另类数据:情绪、卫星图、供应链
MACRO = "macro" # 宏观:利率、PMI、CPI
MICROSTRUCTURE = "microstructure" # 微观结构:价差、VPIN、订单不平衡
RISK = "risk" # 风险:Beta、VaR、因子暴露
CUSTOM = "custom" # 自定义/组合特征
@dataclass
class FeatureDefinition:
"""
特征的完整定义,包括计算逻辑、数据依赖和元数据。
"""
name: str # 特征名称(全局唯一)
description: str # 描述
feature_type: FeatureType # 特征类型
entity: str # 实体(如 'stock', 'futures')
data_sources: List[str] # 数据源依赖
computation_function: str # 计算函数名
parameters: Dict[str, Any] # 计算参数
output_type: str # 输出类型(float, int, category)
refresh_frequency: str # 刷新频率('realtime', 'daily', 'monthly')
owner: str # 负责人
version: int = 1 # 版本号
created_at: str = None # 创建时间
def __post_init__(self):
if self.created_at is None:
self.created_at = datetime.now().isoformat()点击展开可浏览运行结果
📘 本段代码定义了 2 个函数/类:类 `FeatureType`(量化特征的类型分类)、类 `FeatureDefinition`(特征的完整定义,包括计算逻辑、数据依赖和元数据)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
二、在线/离线特征一致性
2.1 训练-服务偏差的来源
训练-服务偏差(Training-Serving Skew)是量化策略从回测到实盘的最大隐患之一。常见来源:
- 特征计算逻辑不同:回测用 Pandas rolling,实盘用增量计算
- 数据源不同:回测用调整后数据,实盘用原始数据
- 时间窗口边界不同:回测月频调仓,实盘日频更新
- 缺失值处理不同:回测用完整数据回填,实盘只能向前填充
2.2 一致性保障设计
此处有展示代码展开 ▼
python
class FeatureRegistry:
"""
特征注册中心:统一管理在线和离线特征的定义和版本。
确保同一特征在训练和推理时使用完全相同的逻辑。
"""
def __init__(self):
self.features: Dict[str, FeatureDefinition] = {}
self.feature_transformations: Dict[str, callable] = {}
def register(self,
feature_def: FeatureDefinition,
offline_compute_fn: callable,
online_compute_fn: Optional[callable] = None):
"""
注册一个特征,同时提供在线和离线计算函数。
参数:
feature_def: 特征定义
offline_compute_fn: 离线(回测)计算函数
online_compute_fn: 在线(实时)计算函数,若为None则使用离线函数
"""
self.features[feature_def.name] = feature_def
self.feature_transformations[feature_def.name] = {
'offline': offline_compute_fn,
'online': online_compute_fn or offline_compute_fn
}
def compute_offline(self,
feature_name: str,
data: pd.DataFrame,
**kwargs) -> pd.Series:
"""使用离线函数计算特征"""
if feature_name not in self.feature_transformations:
raise ValueError(f"Feature '{feature_name}' not registered")
fn = self.feature_transformations[feature_name]['offline']
return fn(data, **kwargs)
def compute_online(self,
feature_name: str,
data: Dict[str, Any],
**kwargs) -> float:
"""使用在线函数计算特征"""
if feature_name not in self.feature_transformations:
raise ValueError(f"Feature '{feature_name}' not registered")
fn = self.feature_transformations[feature_name]['online']
return fn(data, **kwargs)
# 示例:注册动量特征
def momentum_offline(data: pd.DataFrame, window: int = 20) -> pd.Series:
"""离线计算:Pandas rolling window"""
return data['close'].pct_change(window)
def momentum_online(data: Dict[str, Any], window: int = 20) -> float:
"""在线计算:增量更新(使用缓存的历史值)"""
# data 包含 'current_close' 和 'close_N_days_ago'
if data['close_N_days_ago'] is None or data['close_N_days_ago'] == 0:
return 0.0
return data['current_close'] / data['close_N_days_ago'] - 1.0
# 注册
registry = FeatureRegistry()
registry.register(
FeatureDefinition(
name="momentum_20d",
description="20日价格动量",
feature_type=FeatureType.PRICE_DERIVED,
entity="stock",
data_sources=["market_data"],
computation_function="momentum",
parameters={"window": 20},
output_type="float",
refresh_frequency="daily"
),
offline_compute_fn=lambda data: momentum_offline(data, window=20),
online_compute_fn=lambda data: momentum_online(data, window=20)
)点击展开可浏览运行结果
📘 本段代码定义了 3 个函数/类:类 `FeatureRegistry`(特征注册中心:统一管理在线和离线特征的定义和版本)、函数 `momentum_offline`(离线计算:Pandas rolling window)、函数 `momentum_online`(在线计算:增量更新(使用缓存的历史值))。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
三、特征版本管理与回滚
3.1 特征版本化
此处有展示代码展开 ▼
python
class FeatureVersionManager:
"""
特征版本管理器。
支持:
1. 特征定义的版本化存储
2. 回滚到历史版本
3. 版本间差异比较
"""
def __init__(self, backend=None):
self.backend = backend or {} # 简化:使用内存字典
self.feature_versions: Dict[str, List[Dict]] = {}
def save_version(self,
feature_name: str,
definition: FeatureDefinition,
metadata: Dict[str, Any] = None):
"""
保存特征的新版本。
使用语义化版本号:MAJOR.MINOR.PATCH
- MAJOR: 计算逻辑发生重大变更(可能改变特征含义)
- MINOR: 参数调整或数据源变更
- PATCH: 错误修复不影响特征语义
"""
if feature_name not in self.feature_versions:
self.feature_versions[feature_name] = []
# 版本号递增
if self.feature_versions[feature_name]:
prev_version = self.feature_versions[feature_name][-1]['version']
new_version = self._increment_version(
prev_version, metadata.get('change_type', 'PATCH')
)
else:
new_version = '1.0.0'
version_record = {
'version': new_version,
'definition': definition,
'metadata': metadata or {},
'timestamp': datetime.now().isoformat(),
'change_type': metadata.get('change_type', 'PATCH') if metadata
else 'PATCH'
}
self.feature_versions[feature_name].append(version_record)
return new_version
def get_version(self,
feature_name: str,
version: str = None,
as_of_timestamp: str = None) -> Dict:
"""
获取指定版本的特征定义。
参数:
feature_name: 特征名称
version: 版本号,若指定 as_of_timestamp 则忽略
as_of_timestamp: 获取该时间点之前的最新版本
"""
versions = self.feature_versions.get(feature_name, [])
if version:
for v in versions:
if v['version'] == version:
return v
return None
if as_of_timestamp:
ts = pd.Timestamp(as_of_timestamp)
valid_versions = [v for v in versions
if pd.Timestamp(v['timestamp']) <= ts]
return valid_versions[-1] if valid_versions else None
return versions[-1] if versions else None
def rollback(self, feature_name: str, target_version: str) -> Dict:
"""
回滚到指定的历史版本。
回滚会创建一个新版本,其计算逻辑与目标版本相同。
"""
target = self.get_version(feature_name, target_version)
if target is None:
raise ValueError(f"Version '{target_version}' not found "
f"for feature '{feature_name}'")
# 创建回滚版本(更新版本号,继承计算逻辑)
rolled = self.save_version(
feature_name,
target['definition'],
{
'change_type': 'ROLLBACK',
'rolled_from_version': target_version,
'reason': 'Rollback requested'
}
)
return self.get_version(feature_name, rolled)
def compare_versions(self,
feature_name: str,
version_a: str,
version_b: str) -> Dict:
"""
比较两个版本的特征定义差异。
"""
a = self.get_version(feature_name, version_a)
b = self.get_version(feature_name, version_b)
if a is None or b is None:
return {'error': 'Version not found'}
diffs = {}
# 比较参数
params_a = a['definition'].parameters
params_b = b['definition'].parameters
all_params = set(list(params_a.keys()) + list(params_b.keys()))
for param in all_params:
val_a = params_a.get(param)
val_b = params_b.get(param)
if val_a != val_b:
diffs[f'parameter.{param}'] = {
'old': val_a,
'new': val_b
}
return {
'version_a': version_a,
'version_b': version_b,
'differences': diffs,
'has_changes': len(diffs) > 0
}
def _increment_version(self,
current: str,
change_type: str) -> str:
"""语义化版本号递增"""
major, minor, patch = map(int, current.split('.'))
if change_type == 'MAJOR':
return f"{major + 1}.0.0"
elif change_type == 'MINOR':
return f"{major}.{minor + 1}.0"
else: # PATCH or ROLLBACK
return f"{major}.{minor}.{patch + 1}"点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `FeatureVersionManager`(特征版本管理器)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
四、开源方案简介
4.1 Feast(Feature Store)
Feast(Feature Store)是 Google 开源的 Feature Store 框架,已捐赠给 Linux Foundation。
核心概念:
- FeatureView:将数据源与特征定义绑定,定义特征的逻辑视图
- Entity:特征的归属实体(如 driver, customer, stock)
- FeatureService:为模型提供的一组特征的组合
- Online Store:低延迟特征服务(Redis/Datastore)
- Offline Store:历史特征仓库(BigQuery/Parquet/Snowflake)
此处有展示代码展开 ▼
python
# Feast 在量化场景下的概念映射(用纯 Python 表达,便于校验与生成 Feast 定义文件)
import pandas as pd
ENTITIES = {
'stock': 'join_key=symbol,A 股 6 位代码',
'futures_contract': 'join_key=contract_id,如 IF2403',
}
FEATURE_VIEWS = {
'stock_price_features': {
'entity': 'stock',
'source': 'daily_market_data',
'ttl_days': 3,
'features': ['daily_return', 'volume', 'vwap', 'volatility_20d'],
'note': '日频行情,T+0 收盘后即可用',
},
'stock_fundamental_features': {
'entity': 'stock',
'source': 'quarterly_financials',
'ttl_days': 120,
'features': ['pe_ratio', 'pb_ratio', 'roe', 'debt_to_equity'],
'note': '季频财报,必须按「公告日」而非「报告期」做 point-in-time join',
},
}
FEATURE_SERVICES = {
'momentum_strategy_features': ['daily_return', 'volatility_20d', 'volume', 'pe_ratio'],
'mean_reversion_features': ['daily_return', 'z_score_20d', 'rsi_14d'],
}
ONLINE_STORE = {'type': 'Redis', 'key_pattern': 'stock:{symbol}:feature:{feature_name}'}
OFFLINE_STORE = {'type': 'Parquet on S3/MinIO', 'partition': 'feature_name / date'}
print("=== Entity ===")
for name, desc in ENTITIES.items():
print(f" {name:<20}{desc}")
print("\n=== FeatureView ===")
for name, fv in FEATURE_VIEWS.items():
print(f" {name} (entity={fv['entity']}, source={fv['source']}, ttl={fv['ttl_days']}d)")
print(f" features: {', '.join(fv['features'])}")
print(f" note: {fv['note']}")
print("\n=== FeatureService(模型消费的特征组合)===")
declared = {f for fv in FEATURE_VIEWS.values() for f in fv['features']}
for svc, feats in FEATURE_SERVICES.items():
missing = [f for f in feats if f not in declared]
status = 'OK' if not missing else f"缺少定义: {missing}"
print(f" {svc:<32}{len(feats)} 个特征 [{status}]")
print(f"\nOnline Store: {ONLINE_STORE['type']},key = {ONLINE_STORE['key_pattern']}")
print(f"Offline Store: {OFFLINE_STORE['type']},分区 = {OFFLINE_STORE['partition']}")
# ===== point-in-time join:Feast get_historical_features 的核心语义 =====
# 财报「报告期」是 2024Q1,但真正可用的时间是「公告日」。
# 用报告期对齐 = 未来函数;用公告日对齐 = 正确。
fundamentals = pd.DataFrame({
'symbol': ['600519'] * 3,
'report_period': pd.to_datetime(['2023-12-31', '2024-03-31', '2024-06-30']),
'announce_date': pd.to_datetime(['2024-03-29', '2024-04-27', '2024-08-09']),
'roe': [0.310, 0.082, 0.171],
})
signal_dates = pd.DataFrame({
'symbol': ['600519'] * 4,
'event_timestamp': pd.to_datetime(['2024-04-01', '2024-04-26', '2024-05-06', '2024-08-12']),
})
wrong = pd.merge_asof(signal_dates.sort_values('event_timestamp'),
fundamentals.sort_values('report_period'),
left_on='event_timestamp', right_on='report_period',
by='symbol', direction='backward')
right = pd.merge_asof(signal_dates.sort_values('event_timestamp'),
fundamentals.sort_values('announce_date'),
left_on='event_timestamp', right_on='announce_date',
by='symbol', direction='backward')
cmp = pd.DataFrame({
'取数日': signal_dates['event_timestamp'].dt.date,
'按报告期(错)': wrong['roe'].values,
'按公告日(对)': right['roe'].values,
})
print("\n=== point-in-time join 对照 ===")
print(cmp.to_string(index=False))
n_bad = int((cmp['按报告期(错)'] != cmp['按公告日(对)']).sum())
print(f"\n{n_bad}/{len(cmp)} 个取数日出现未来函数:按报告期对齐会提前拿到尚未公告的 ROE。")
print("Feast 的 get_historical_features 靠 event_timestamp + ttl 自动规避这个坑。")点击展开可浏览运行结果
=== Entity ===
stock join_key=symbol,A 股 6 位代码
futures_contract join_key=contract_id,如 IF2403
=== FeatureView ===
stock_price_features (entity=stock, source=daily_market_data, ttl=3d)
features: daily_return, volume, vwap, volatility_20d
note: 日频行情,T+0 收盘后即可用
stock_fundamental_features (entity=stock, source=quarterly_financials, ttl=120d)
features: pe_ratio, pb_ratio, roe, debt_to_equity
note: 季频财报,必须按「公告日」而非「报告期」做 point-in-time join
=== FeatureService(模型消费的特征组合)===
momentum_strategy_features 4 个特征 [OK]
mean_reversion_features 3 个特征 [缺少定义: ['z_score_20d', 'rsi_14d']]
Online Store: Redis,key = stock:{symbol}:feature:{feature_name}
Offline Store: Parquet on S3/MinIO,分区 = feature_name / date
=== point-in-time join 对照 ===
取数日 按报告期(错) 按公告日(对)
2024-04-01 0.082 0.310
2024-04-26 0.082 0.310
2024-05-06 0.082 0.082
2024-08-12 0.171 0.171
2/4 个取数日出现未来函数:按报告期对齐会提前拿到尚未公告的 ROE。
Feast 的 get_historical_features 靠 event_timestamp + ttl 自动规避这个坑。4.2 Tecton
Tecton 是 Feast 的商业化版本(由 Feast 的核心贡献者创建),增加了:
- Stream Feature Views:内置流处理特征(如 Kafka -> 实时聚合)
- 特征监控:自动检测特征漂移和数据质量下降
- 按需特征转换:在线请求时动态计算特征(适用于需要上下文信息的特征)
4.3 自建轻量级 Feature Store
此处有展示代码展开 ▼
python
class LightweightFeatureStore:
"""
轻量级 Feature Store 实现(适用于中小型量化团队)。
理念:在已有的数据基础设施(Parquet + Redis)上
增加一层 Feature 抽象,避免引入重型框架。
"""
def __init__(self,
offline_store_path: str,
redis_config: Dict = None):
"""
参数:
offline_store_path: 离线特征存储路径(Parquet 文件)
redis_config: Redis 连接配置(用于在线服务)
"""
self.offline_path = offline_store_path
self.redis_client = None
if redis_config:
try:
import redis
self.redis_client = redis.Redis(**redis_config)
except ImportError:
print("redis 库未安装,在线特征服务不可用")
def get_offline_features(self,
feature_names: List[str],
symbols: List[str],
start_date: str,
end_date: str) -> pd.DataFrame:
"""
从离线存储获取历史特征值(用于回测)。
自动处理 point-in-time join。
"""
frames = []
for feature in feature_names:
feature_path = f"{self.offline_path}/{feature}/"
try:
feature_data = pd.read_parquet(
feature_path,
filters=[
('symbol', 'in', symbols),
('date', '>=', start_date),
('date', '<=', end_date)
]
)
feature_data = feature_data.set_index(['date', 'symbol'])
frames.append(feature_data[[feature]])
except FileNotFoundError:
print(f"警告: 特征 '{feature}' 的数据未找到")
continue
if not frames:
return pd.DataFrame()
# 横向拼接
result = pd.concat(frames, axis=1).reset_index()
return result
def get_online_features(self,
symbol: str,
feature_names: List[str]) -> Dict[str, float]:
"""
从 Redis 在线存储获取最新的特征值(用于实时交易)。
延迟目标:< 1ms(本地 Redis)或 < 5ms(网络 Redis)。
"""
if self.redis_client is None:
raise RuntimeError("Redis 未配置,无法获取在线特征")
features = {}
pipeline = self.redis_client.pipeline()
for feat in feature_names:
key = f"feature:{symbol}:{feat}"
pipeline.get(key)
try:
values = pipeline.execute()
except Exception as e:
print(f"Redis pipeline 执行失败: {e}")
return {}
for feat, val in zip(feature_names, values):
features[feat] = float(val) if val is not None else None
return features
def materialize_features(self,
feature_names: List[str],
date: str):
"""
将当天的离线特征计算结果同步到在线 Redis 存储。
应在每日批处理完成后调用。
"""
if self.redis_client is None:
return
offline = self.get_offline_features(
feature_names,
symbols=[], # 所有 symbol
start_date=date,
end_date=date
)
pipeline = self.redis_client.pipeline()
for _, row in offline.iterrows():
symbol = row['symbol']
for feat in feature_names:
if feat in row and not pd.isna(row[feat]):
key = f"feature:{symbol}:{feat}"
pipeline.set(key, float(row[feat]))
pipeline.execute()点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `LightweightFeatureStore`(轻量级 Feature Store 实现(适用于中小型量化团队))。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
Feature Store 将量化策略的"特征工程"从临时性的脚本提升为有组织的系统工程。它确保了回测中使用的特征与实盘中使用的特征完全一致——这是量化策略从研究走向生产的最后也是最重要的一步。