Feast 第三方集成生态:插件标准、扩展机制与自定义 Store 开发指南
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
Feast 作为开源 Feature Store,其核心设计目标之一就是"融入你已有的技术栈"——无论是数据源(BigQuery、Snowflake、Redshift、Kafka),还是离线/在线存储(Redis、DynamoDB、DuckDB、向量数据库),Feast 都提供了官方内置或社区维护的插件化集成。本文以 Feast 仓库的第三方集成文档为主线,系统梳理集成清单与 Roadmap、插件被收录和合入主仓库必须满足的工程标准,并结合仓库源码深入讲解自定义离线存储(Offline Store)与在线存储(Online Store)的接口实现、配置加载与通用测试机制,帮助你判断"我的存储/数据源如何接入 Feast",以及"如何把一个插件从 contrib 提升为官方集成"。
一、Feast 的第三方集成生态概览
Feast 与一大批工具和技术集成,覆盖了数据源、离线存储、在线存储、向量检索、流式摄取、部署形态、特征服务与治理等完整链路。这些集成大部分以插件(plugin)的形式维护在 Feast 主仓库内(sdk/python/feast/infra/下),少部分由社区在独立仓库维护。官方对集成的完整功能清单与演进状态集中在 docs/roadmap.md 中,其中每一项都标注了完成状态([x]已完成 /[ ]规划中)以及对应的参考文档链接。
从 roadmap 的目录结构可以清晰地看到 Feast 集成面的广度:
| 分类 | 已完成的代表性集成 |
|---|---|
| 数据源 | Snowflake、Redshift、BigQuery、Parquet 文件、Kafka/Kinesis(经由 push 支持)、Postgres、Spark、Athena、ClickHouse、Oracle、MongoDB、Ray 等 |
| 离线存储 | Snowflake、Redshift、BigQuery、DuckDB、Dask、Remote、Hybrid,以及 MSSQL、Hive、Trino、Spark、Couchbase、ClickHouse、Ray、Oracle、MongoDB 等 contrib 插件 |
| 在线存储 | Snowflake、DynamoDB、Redis、Dragonfly、Datastore、Bigtable、SQLite、Remote、Postgres、HBase、Cassandra/AstraDB、ScyllaDB、MySQL、Hazelcast、Elasticsearch、SingleStore、Couchbase、MongoDB、Aerospike |
| 向量存储 | Qdrant、Milvus、Faiss(作为向量在线存储) |
| 流式 | 自定义流式摄取作业、基于 push 的流式数据写入在线/离线存储 |
| 部署 | AWS Lambda、Kubernetes(见 docs/how-to-guides/running-feast-in-production.md) |
| 特征服务 | Python 客户端、Python Feature Server、Java/Go Feature Server(alpha)、Offline Feature Server、Registry Server、Feast Operator |
| 治理与发现 | Python SDK / CLI 浏览 registry、Feature Service、Amundsen 与 DataHub 集成、Feast Web UI |
需要特别说明的是:这份清单不是封闭的。如果你在列表中没有找到自己正在使用的离线存储或在线存储,Feast 明确提供了两条"自己动手"的路径:
- 添加一个新的离线存储
- 添加一个新的在线存储
二、插件(Plugin)的定位与 contrib 机制
在深入自定义开发之前,先理解 Feast 对"插件"的等级划分。从 sdk/python/feast/infra/offline_stores/ 与 sdk/python/feast/infra/online_stores/ 的目录结构可以直观看到:官方内置实现(如bigquery.py、redshift.py、snowflake.py、redis.py、sqlite.py、dynamodb.py)直接位于父目录下,而社区贡献的实现统一放在contrib/子目录中,例如:
- 离线存储 contrib:
athena_offline_store、clickhouse_offline_store、mssql_offline_store、oracle_offline_store、postgres_offline_store、ray_offline_store、spark_offline_store、trino_offline_store、couchbase_offline_store、mongodb_offline_store - 在线存储 contrib(在线存储目录中按子包组织):
mysql_online_store、hbase_online_store、cassandra_online_store、hazelcast_online_store、elasticsearch_online_store、singlestore_online_store、qdrant_online_store、milvus_online_store、aerospike_online_store、scylladb_online_store等
contrib 插件的定位(官方明确声明):
- 不保证实现接口的全部方法;
- 不保证 API 稳定,后续可能发生变化;
- 实现中应通过
RuntimeWarning向用户提示"这是一个 contrib 插件,不由 Feast 维护者维护",源码中的标准警告文案为:"This offline store is an experimental feature in alpha development. Some functionality may still be unstable so functionality can change in the future."
从 contrib 提升为官方插件需要满足两个条件:
- CI(
make test-python-integration)已配置为对该插件运行全部测试并通过; - 至少两名贡献者拥有该插件的维护权(理想情况下记录在
OWNERS/CODEOWNERS文件中)。
三、集成被"收录"与"合入主仓库"的工程标准
第三方集成文档的核心价值在于定义了清晰的验收门槛,分为两个层级:
3.1 插件被收录(被官方列出)的要求
一个插件集成若要被官方收录并突出展示,必须满足以下三条:
- 必须有测试。理想情况下应复用 Feast 的 universal 测试套件(示例见 添加或复用测试),自定义测试也可以接受;
- 必须有基本的使用文档,说明如何使用该插件;
- 作者必须与维护者协作通过一次基础代码评审,确保实现与 Feast 核心实现大致对齐(例如遵循相同的数据模型与调用约定)。
3.2 插件合入主仓库的要求
若插件要合并进 Feast 主仓库,标准更严格:
- PR 必须通过全部集成测试,且必须更新 universal 测试(专门为自定义集成设计的测试)以覆盖该集成;
- 必须有文档和教程说明如何使用该集成;
- 作者(或其他贡献者)必须承诺拥有这些文件的维护权,并持续维护;
- 如果插件由组织而非个人贡献,该组织应为集成测试提供基础设施(或云额度)。
这套标准保证了第三方集成不仅仅是"能跑",而是可持续、可验证、有归属的工程资产。
四、自定义离线存储:接口、配置与调用链
离线存储是 Feast 中"读历史特征、做 point-in-time join、支撑物化"的存储与计算层。完整的自定义教程见 添加一个新的离线存储,其流程分为 8 步:定义OfflineStore类、定义OfflineStoreConfig类、定义RetrievalJob类、定义DataSource类、在feature_store.yaml中引用、测试、更新依赖、补充文档。
4.1 OfflineStore 接口
离线存储的抽象基类定义在 sdk/python/feast/infra/offline_stores/offline_store.py 中(OfflineStore类,OfflineStore(ABC))。从源码看,一个完整的离线存储实现需要覆盖以下核心方法:
| 方法 | 触发时机 | 说明 |
|---|---|---|
pull_latest_from_table_or_query | feast materialize/feast materialize-incremental命令,或FeatureStore.materialize() | 从离线存储拉取数据,由FeatureStore负责写入在线存储 |
get_historical_features | FeatureStore.get_historical_features() | 读取历史特征,典型用于训练 ML 模型,执行 point-in-time 正确性连接 |
pull_all_from_table_or_query(可选) | SavedDatasets 与特征质量监控的兜底计算路径 | 按起止日期拉取全部数据 |
write_logged_features(可选) | 内部用于 SavedDatasets | 将 pyarrow 表或 parquet 路径写入LoggingSource/LoggingConfig指定的源 |
offline_write_batch(可选) | Push API | 将 pyarrow 表直接推送到指定 feature view 的 batch source |
接口签名示例(节选自自定义教程):
def get_historical_features(self, config: RepoConfig, feature_views: List[FeatureView], feature_refs: List[str], entity_df: Union[pd.DataFrame, str], registry: Registry, project: str, full_feature_names: bool = False) -> RetrievalJob: """Perform point-in-time correct join of features onto an entity dataframe.""" # Implementation here. pass一个值得注意的实现细节:OfflineStore类名必须以OfflineStore后缀结尾(教程中以 hint 形式强调)。这一约定不是装饰性的——sdk/python/feast/repo_config.py 中的get_offline_store_type函数在加载时会校验:如果配置的 type 不在OFFLINE_STORE_CLASS_FOR_TYPE字典中,则要求字符串以OfflineStore结尾,否则抛出FeastOfflineStoreInvalidName。同时,Feast 通过OFFLINE_STORE_CLASS_FOR_TYPE字典(如"file"、"bigquery"、"redshift"、"duckdb"等)实现了"简短别名 → 全限定类名"的映射,第三方实现与第一方实现走同一条类加载代码路径。
另外,与 OnlineStore 不同,Feast 不为离线存储管理任何基础设施——离线存储本质上是"计算 + 存储"的查询引擎,其方法(如get_historical_features、pull_latest_from_table_or_query)均为静态方法,直接对底层引擎发起查询。
4.2 RetrievalJob:惰性执行的查询句柄
离线存储的方法不应急切执行读操作,而是返回一个RetrievalJob实例,代表对底层存储的实际查询执行。RetrievalJob抽象类同样定义在 sdk/python/feast/infra/offline_stores/offline_store.py 中,其核心能力包括:
to_df():同步执行查询并返回 pandas DataFrame(会执行 on-demand 变换与可选的验证);to_arrow():同步执行查询并返回 pyarrow Table(to_df内部实际委托给to_arrow().to_pandas());to_remote_storage()(可选):将结果以多个 parquet 文件导出到远端存储(如 S3/GCS),配合自定义 Materialization Engine 实现可扩展的批量物化——将物化数据分块并行处理。若未实现,Feast 默认走本地物化(把全部记录拉入内存);persist()、metadata、to_sql()等辅助能力。
自定义RetrievalJob的典型实现模式(以自定义文件离线存储为例):
class CustomFileRetrievalJob(RetrievalJob): def __init__(self, evaluation_function: Callable): """Initialize a lazy historical retrieval job""" # The evaluation function executes a stored procedure to compute a historical retrieval. self.evaluation_function = evaluation_function def to_df(self): # Only execute the evaluation function to build the final historical retrieval dataframe at the last moment. df = self.evaluation_function() return df def to_arrow(self): df = self.evaluation_function() return pyarrow.Table.from_pandas(df) def to_remote_storage(self): # Optional method to write to an offline storage location to support scalable batch materialization. pass4.3 OfflineStoreConfig:基于 pydantic 的配置模型
所有OfflineStore实现必须在同文件内定义对应的OfflineStoreConfig类,该类继承自 sdk/python/feast/repo_config.py 中的FeastConfigBaseModel(一个 pydantic 模型,负责把 YAML 配置解析为 Python 对象,并可通过 validator 校验配置正确性)。命名约束是:配置类名 = 离线存储类名 +Config后缀,且必须包含type字段,值为对应OfflineStore类的全限定类名。
class CustomFileOfflineStoreConfig(FeastConfigBaseModel): """ Custom offline store config for local (file-based) store """ type: Literal["feast_custom_offline_store.file.CustomFileOfflineStore"] \ = "feast_custom_offline_store.file.CustomFileOfflineStore" uri: str # URI for your offline store(in this case it would be a path)对应在feature_store.yaml中的配置方式:
project: my_project registry: data/registry.db provider: local offline_store: type: feast_custom_offline_store.file.CustomFileOfflineStore uri: <File URI> online_store: path: data/online_store.db该配置会以config.offline_store字段的形式注入到OfflineStore的每个方法参数config: RepoConfig中,实现中通过assert isinstance(offline_store_config, CustomFileOfflineStoreConfig)取回类型安全的配置对象。
4.4 DataSource 与类型映射
要让自定义离线存储作为 feature view 的 batch source 使用,还需定义DataSource的子类,实现from_proto与to_proto两个方法。对于不在主仓库实现的自定义离线存储,应使用custom_options字段存储数据源所需配置,并在to_proto中将其序列化为字节、在from_proto中读回(通常以 JSON 形式编码):
class CustomFileDataSource(FileSource): """Custom data source class for local files""" @staticmethod def from_proto(data_source: DataSourceProto): custom_source_options = str( data_source.custom_options.configuration, encoding="utf8" ) path = json.loads(custom_source_options)["path"] return CustomFileDataSource( field_mapping=dict(data_source.field_mapping), path=path, timestamp_field=data_source.timestamp_field, created_timestamp_column=data_source.created_timestamp_column, date_partition_column=data_source.date_partition_column, ) def to_proto(self) -> DataSourceProto: config_json = json.dumps({"path": self.path}) data_source_proto = DataSourceProto( type=DataSourceProto.CUSTOM_SOURCE, custom_options=DataSourceProto.CustomSourceOptions( configuration=bytes(config_json, encoding="utf8") ), ) data_source_proto.timestamp_field = self.timestamp_field data_source_proto.created_timestamp_column = self.created_timestamp_column data_source_proto.date_partition_column = self.date_partition_column return data_source_proto此外,大多数离线存储都需要实现自定义类型映射:在DataSource类中实现source_datatype_to_feast_value_type(将数据源类型转换为 Feast 值类型)与get_column_names_and_types(获取列名与对应类型)。类型转换的辅助函数可以放入sdk/python/feast/type_map.py。类型映射必须正确,否则 Feast 处理特征列时可能发生错误转换,导致信息丢失或数据错误——这是自定义 store 开发中出错率最高的环节之一。
4.5 在 feature repo 中使用自定义离线存储
只要OfflineStore类在 Python 环境中可用,Feast 会在运行时动态导入。feature_store.yaml中offline_store.type应指定为可导入的全限定类名:
project: test_custom registry: data/registry.db provider: local offline_store: # Make sure to specify the type as the fully qualified path that Feast can import. type: feast_custom_offline_store.file.CustomFileOfflineStore如果无需额外配置,可以省略其他字段,只写type:
offline_store: feast_custom_offline_store.file.CustomFileOfflineStore随后在 repo 定义文件中使用自定义数据源定义 FeatureView:
driver_hourly_stats = CustomFileDataSource( path="feature_repo/data/driver_stats.parquet", timestamp_field="event_timestamp", created_timestamp_column="created", ) driver_hourly_stats_view = FeatureView( source=driver_hourly_stats, ... )五、自定义在线存储:基础设施管理与读写方法
在线存储是服务阶段低延迟读取特征值的存储层。完整教程见 添加一个新的在线存储,流程为 6 步:定义OnlineStore类、定义OnlineStoreConfig类、在feature_store.yaml中引用、测试、更新依赖、补充文档。
5.1 OnlineStore 接口的两组方法
抽象基类定义在 sdk/python/feast/infra/online_stores/online_store.py(OnlineStore(ABC))。从源码结构看,其方法分为两组:
(1)基础设施管理方法——与离线存储不同,Feast会为在线存储管理基础设施:
| 方法 | 触发时机 | 职责 |
|---|---|---|
update | feast apply命令或FeatureStore.apply()SDK 方法 | 执行写/读数据前的必要操作,例如为新的 feature view 创建 MySQL 表与索引 |
teardown | feast teardown或FeatureStore.teardown() | 执行清理,例如删除被删 feature view 对应的表与索引 |
(2)读写方法:
| 方法 | 触发时机 | 说明 |
|---|---|---|
online_write_batch | feast materialize/feast materialize-incremental或FeatureStore.materialize() | 批量写入特征行,数据以(EntityKeyProto, Dict[str, ValueProto], event_ts, created_ts)四元组列表传入 |
online_read | FeatureStore.get_online_features() | 按实体键读取特征值,返回与entity_keys等长的(event_ts, {feature_name: ValueProto})列表 |
在线存储实现中,实体键需要经过serialize_entity_key(entity_key, entity_key_serialization_version=config.entity_key_serialization_version)序列化后写入存储;读取时通过ValueProto.ParseFromString()还原特征值(protobuf 编码)。此外,OnlineStore基类还提供get_online_features、plan、异步读写、预计算向量读写等扩展能力,自定义实现可根据需要覆盖。
5.2 OnlineStoreConfig 与 feature_store.yaml 引用
与离线存储相同的约定:OnlineStoreConfig类继承FeastConfigBaseModel、以Config后缀命名、必须含type字段。以 MySQL 为例:
class MySQLOnlineStoreConfig(FeastConfigBaseModel): type: Literal["feast_custom_online_store.mysql.MySQLOnlineStore"] = "feast_custom_online_store.mysql.MySQLOnlineStore" host: Optional[StrictStr] = None user: Optional[StrictStr] = None password: Optional[StrictStr] = None database: Optional[StrictStr] = Noneproject: test_custom registry: data/registry.db provider: local online_store: # Make sure to specify the type as the fully qualified path that Feast can import. type: feast_custom_online_store.mysql.MySQLOnlineStore user: foo password: bar同样,若无需额外配置可以简写为online_store: feast_custom_online_store.mysql.MySQLOnlineStore。配置在方法内通过config.online_store访问,例如host=online_store_config.host or "127.0.0.1"。类名必须以OnlineStore结尾,repo_config.py中的get_online_config_from_type会据此校验并动态加载。
六、插件的通用测试机制:FULL_REPO_CONFIGS 与 universal 测试
Feast 为第三方存储插件提供了复用官方测试套件的能力,这是插件质量门槛("必须有测试")的技术支撑。测试体系细节见 添加或复用测试。
6.1 测试套件结构
- 单元测试位于
sdk/python/tests/unit,集成测试位于sdk/python/tests/integration; - 集成测试按组件组织:
e2e(端到端)、offline_store(历史检索、push API、特征日志)、online_store(在线检索、push)、registration(registry、CLI、类型推断)、materialization、feature_repos(测试夹具与配置)等; - universal feature repo是一套可参数化的夹具(如
environment与universal_data_sources),可覆盖不同离线存储 × 在线存储 × provider 的组合,让同一测试代码跑遍所有存储后端; - 测试标记(marker):
@pytest.mark.integration标记集成测试,@pytest.mark.universal_offline_stores将测试参数化到所有 universal 离线存储(file、redshift、bigquery、snowflake),universal_online_stores(only=["redis"])可指定特定在线存储。
6.2 用 FULL_REPO_CONFIGS_MODULE 覆盖测试矩阵
核心机制在 sdk/python/tests/universal/feature_repos/repo_configuration.py:
- Feast 通过
FULL_REPO_CONFIGS变量参数化集成测试,该变量聚合AVAILABLE_OFFLINE_STORES(每个条目为(provider, DataSourceCreator))与AVAILABLE_ONLINE_STORES(每个条目为(store_config, OnlineStoreCreator))生成IntegrationTestRepoConfig; - 插件无需修改 Feast 仓库即可覆盖测试矩阵:创建自己的模块文件定义
AVAILABLE_OFFLINE_STORES/AVAILABLE_ONLINE_STORES/FULL_REPO_CONFIGS,通过环境变量FULL_REPO_CONFIGS_MODULE指向该模块,再用make test-python-universal运行。
仓库中真实的 contrib 插件配置示例 sdk/python/feast/infra/offline_stores/contrib/postgres_repo_configuration.py:
from feast.infra.offline_stores.contrib.postgres_offline_store.tests.data_source import ( PostgreSQLDataSourceCreator, ) from tests.universal.feature_repos.repo_configuration import REDIS_CONFIG from tests.universal.feature_repos.universal.online_store.redis import ( RedisOnlineStoreCreator, ) AVAILABLE_OFFLINE_STORES = [("local", PostgreSQLDataSourceCreator)] AVAILABLE_ONLINE_STORES = {"redis": (REDIS_CONFIG, RedisOnlineStoreCreator)}配套的 Makefile 目标(以 Spark 为例,摘自自定义离线存储教程):
test-python-universal-spark: PYTHONPATH='.' \ FULL_REPO_CONFIGS_MODULE=sdk.python.feast.infra.offline_stores.contrib.spark_repo_configuration \ PYTEST_PLUGINS=feast.infra.offline_stores.contrib.spark_offline_store.tests \ IS_TEST=True \ python -m pytest -n 8 --integration \ -k "not test_historical_retrieval_fails_on_validation and \ not test_historical_retrieval_with_validation and ..." \ sdk/python/tests其中PYTEST_PLUGINS环境变量让 pytest 加载DataSourceCreator/OnlineStoreCreator;-k选项可剔除与特定存储无关或暂不适用的测试。
6.3 DataSourceCreator 与 OnlineStoreCreator
DataSourceCreator.create_data_source(df, ...):把测试传入的 DataFrame 上传/写入对应离线存储,返回指向该位置的数据源对象(BigQueryDataSourceCreator是参考实现);OnlineStoreCreator:用于在容器中拉起在线存储实例。例如 Redis 的 creator 会启动容器并等待 "Ready to accept connections" 日志后再开始测试,这样其他开发者无需自建实例即可运行测试。
若集成测试失败,说明该存储插件的实现存在错误——这是最直接的信号。测试通过后,还应把数据源注册到repo_config.py的OFFLINE_STORE_CLASS_FOR_TYPE/ONLINE_STORE_CLASS_FOR_TYPE字典中(与spark、trino等一致),这样 Feast 才能从feature_store.yaml中加载对应类。
七、依赖、文档与合入流程的收尾
一个插件在代码与测试就绪后,还需完成以下工程化步骤:
- 依赖管理:在 sdk/python/pyproject.toml 的
[project.optional-dependencies]下为插件新增 extra(如<offline_store> = ["package1>=1.0", "package2"]),使其默认不随 Feast 安装;随后运行make lock-python-dependencies-all重新生成锁文件; - 文档:
- 离线存储:在
docs/reference/offline-stores/与docs/reference/data-sources/各新增一个 Markdown 文件,并在docs/reference/data-sources/README.md与docs/SUMMARY.md中登记索引; - 在线存储:在
docs/reference/online-stores/新增文档并在其README.md与docs/SUMMARY.md中登记; - 文档必须覆盖:如何创建数据源、
feature_store.yaml需要哪些配置、明确标注该数据源处于 alpha 开发阶段、该存储的数据模型说明; - 最后运行
make build-sphinx生成 Python API 文档;
- 离线存储:在
- 合入主仓库前的收尾:按前文"合入主仓库的四条要求"逐一核对——PR 通过全部集成测试并更新 universal 测试、提供文档与教程、明确文件所有权与后续维护承诺、组织贡献者提供测试基础设施或额度。
八、总结:从"接入"到"贡献"的完整路径
Feast 的第三方集成机制可以总结为一条清晰的路径:
- 查清单:在 docs/roadmap.md 确认你的技术栈是否已有集成,以及其状态(官方 / contrib / 规划中);
- 选路线:没有现成集成时,按 添加离线存储 或 添加在线存储 教程实现接口(注意类名后缀、Config 类命名、
type字段与类型映射等硬性约定); - 接入使用:在
feature_store.yaml中用全限定类名引用,Feast 运行时动态导入; - 过测试:通过
FULL_REPO_CONFIGS_MODULE+make test-python-universal复用 universal 测试套件,用DataSourceCreator/OnlineStoreCreator接入测试基础设施; - 工程化收尾:以 extra 形式管理依赖、补齐参考文档与教程、登记索引,最后对照"收录三条 + 合入四条"标准推进评审与合入。
理解这套标准与机制,既能让现有技术栈顺利接入 Feast,也能为贡献者提供一份可执行、可验收的插件开发路线图——这正是第三方集成文档的核心价值所在。
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考