
Feast 离线存储技术详解OfflineStore 接口、RetrievalJob 执行模型与各实现的功能矩阵【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast本文基于 Feast 官方文档 离线存储总览 与配套源码系统讲解离线存储Offline Store的五大核心接口方法、RetrievalJob延迟执行模型以及 Dask、BigQuery、Snowflake、Redshift、Postgres、Spark、Trino、DuckDB、Couchbase、Ray 等实现的功能支持矩阵。读完本文后你可以理解 Feast 离线存储的抽象设计与各后端的实际能力边界并在feature_store.yaml中做出有依据的后端选型。离线存储在 Feast 中的定位离线存储是 Feast 架构中负责离线特征Offline Features的存储与计算系统。它承载特征视图Feature View所引用的历史数据是特征训练、特征导出、批量回写的落点实时场景下的在线存储Online Store则承担低延迟点查。关于离线存储的概念性讲解可参阅 离线存储组件文档。从源码结构看离线存储的抽象定义在 offline_store.py 中由两个抽象基类构成OfflineStore见 OfflineStore 基类定义 Feast 与离线特征的存储和计算系统交互的接口。每个实现只与其对应的数据源DataSource搭配工作例如SnowflakeOfflineStore只能处理SnowflakeSource而不能处理FileSourceRetrievalJob见 RetrievalJob 基类管理从离线存储取数这一查询任务的执行不同后端有自己的具体实现如SnowflakeRetrievalJob。OfflineStore 接口的五个核心方法官方文档列出了OfflineStore接口暴露的五个核心方法及其支持的核心功能结合源码签名逐一说明如下。get_historical_features基于时间点point-in-time的历史特征查询该方法执行时间点正确的 joinpoint-in-time correct join来检索历史特征是特征训练数据集构建的入口。源码签名见 get_historical_features揭示了几个关键细节entity_df可以是 pandas DataFrame也可以是一条 SQL 查询若为None则按指定的时间范围start_date/end_date直接取特征full_feature_names为True时特征列名会从feature变为feature_view__feature形式例如daily_transactions变为customer_fv__daily_transactions避免多个特征视图存在同名字段时冲突关键字参数filter_by_created_timestamp若源声明了supports_filter_by_created_timestamp则只取创建时间戳不晚于实体行事件时间戳的特征值进一步收紧时间正确性语义。pull_latest_from_table_or_query抽取最新行用于在线物化该方法从指定数据源中抽取时间范围内的最新实体行join key 列 特征列 时间戳列的组合主要服务于离线数据物化到在线存储的链路feast materialize。源码 docstring 特别强调传入的所有列名必须是数据源中真实存在的列列名映射mapping必须在此前已经完成。签名中还包括created_timestamp_column用于在时间戳相同时打破平局等参数见 pull_latest_from_table_or_query。pull_all_from_table_or_query抽取全量行以读取保存的数据集与pull_latest相对该方法抽取指定时间范围内数据源中的全部实体行不做每个实体只取最新一行的去重用于读取已保存的数据集Saved Dataset等场景。start_date/end_date均为可选见 pull_all_from_table_or_query。offline_write_batch批量回写面向 Push 源将指定的 Arrow Table 写入某个特征视图底层数据源主要服务于 push 型数据源实时推送写入。实现中接受一个progress回调函数用于向调用方报告写入进度见 offline_write_batch。write_logged_features特征日志持久化将推理时记录的线上特征logged features写入离线存储的指定目标。其行为语义是追加写目标不存在则创建存在则追加因此可以反复调用同一目标、以分块方式刷入日志。入参data可以是 Arrow Table也可以是一个包含待写日志的 Parquet 目录路径见 write_logged_features。源码中接口还包含哪些扩展能力除文档列出的五个核心方法外当前版本源码中OfflineStore基类还定义了若干扩展接口未实现的后端保持NotImplementedError由调用方决定回退策略validate_data_source/get_table_column_names_and_types_from_data_source数据源校验与列类型读取见 validate_data_sourcecompute_monitoring_metrics、get_monitoring_max_timestamp直接用后端原生计算引擎计算监控指标统计量、分位数、直方图不支持的后端由监控服务回退到 Python 计算ensure_monitoring_tables、save_monitoring_metrics、query_monitoring_metrics、clear_monitoring_baseline在离线存储中以原生表持久化监控指标upsert 语义支持daily、weekly、biweekly、monthly、quarterly等粒度见 监控指标存储接口。RetrievalJob延迟执行的结果集抽象get_historical_features、pull_latest_from_table_or_query、pull_all_from_table_or_query三个方法返回的都是特定于离线存储后端的RetrievalJob如SnowflakeRetrievalJob见 SnowflakeRetrievalJob 定义。文档列出了RetrievalJob支持的完整功能清单导出为 DataFrameto_df导出为 Arrow Tableto_arrow导出为 Arrow 批次to_arrow_batches用于处理内存中的大数据集导出为 SQLto_sql导出到数据湖S3、GCS 等to_remote_storage导出到数据仓库导出为 Spark DataFrame在本地执行 Python 的按需转换on-demand transforms远程执行 Python 的按需转换将结果持久化回离线存储persist执行前预览查询计划RetrievalJob是延迟执行的——只有调用导出/持久化方法时才真正运行查询读取分区数据源码中的执行链路与细节从 RetrievalJob 基类实现 可以看到具体执行机制to_df并不是独立实现而是先调用to_arrow再转成 pandas DataFrame即 Arrow 是统一的中间表示。to_arrow内部调用子类的_to_arrow_internal真正执行查询按需转换On-Demand Feature View在to_arrow内完成基类遍历on_demand_feature_views调用odfv.transform_arrow(features_table, ...)将转换结果列按只保留显式请求的特征列的语义追加到结果表中见 按需转换执行段若传入了validation_reference实验性功能结果表在返回前会先经过数据集校验失败则抛出ValidationFailedto_remote_storage的契约是把结果导出为多个大小合适的 Parquet 文件并返回文件路径列表默认实现不支持supports_remote_storage_export返回False由具体后端覆盖如 BigQuery 的 to_remote_storage当前版本源码中RetrievalJob还提供了文档矩阵之外的便捷出口to_feast_df返回带引擎信息与元数据的FeastDataFrame、to_tensor直接转换为 torch 张量字典便于训练、to_ray_dataset转换为 Ray DatasetRay 后端可保持计算在集群内其他后端则先落到驱动端 Arrow 再转换见 to_ray_dataset 说明可观测性to_arrow会记录offline_store_request_total、offline_store_request_latency_seconds、offline_store_row_count等指标开启 audit logging 时还会按特征视图维度输出审计日志见 指标与审计段。各离线存储实现的功能矩阵文档指出目前有四个核心离线存储实现DaskOfflineStore、BigQueryOfflineStore、SnowflakeOfflineStore与RedshiftOfflineStore另有社区贡献的若干实现PostgreSQLOfflineStore、SparkOfflineStore、TrinoOfflineStore与RayOfflineStore不保证稳定也不保证与核心实现功能对齐。各存储的具体配置方式如如何在feature_store.yaml中配置见 离线存储文档索引。核心方法支持矩阵DaskBigQuerySnowflakeRedshiftPostgresSparkTrinoCouchbaseRayget_historical_featuresyesyesyesyesyesyesyesyesyespull_latest_from_table_or_queryyesyesyesyesyesyesyesyesyespull_all_from_table_or_queryyesyesyesyesyesyesyesyesyesoffline_write_batchyesyesyesyesnonononoyeswrite_logged_featuresyesyesyesyesnonononoyes可以看到查询类三方法为全量支持而写入类方法offline_write_batch、write_logged_features在 Postgres、Spark、Trino、Couchbase 上不支持——即在这些后端上push 源写入与特征日志落盘不可用。RetrievalJob 功能支持矩阵DaskBigQuerySnowflakeRedshiftPostgresSparkTrinoDuckDBCouchbaseRayexport to dataframeyesyesyesyesyesyesyesyesyesyesexport to arrow tableyesyesyesyesyesyesyesyesyesyesexport to arrow batchesnononoyesnonononononoexport to SQLnoyesyesyesyesnoyesnoyesnoexport to data lake (S3, GCS, etc.)nonoyesnoyesnononoyesyesexport to data warehousenoyesyesyesyesnononoyesnoexport as Spark dataframenonoyesnonoyesnonononolocal execution of Python-based on-demand transformsyesyesyesyesyesnoyesyesyesyesremote execution of Python-based on-demand transformsnonononononononononopersist results in the offline storeyesyesyesyesyesyesnoyesyesyespreview the query plan before executionyesyesyesyesyesyesyesnoyesyesread partitioned datayesyesyesyesyesyesyesyesyesyes矩阵中的几个要点Arrow 批次导出目前仅 Redshift 支持其to_arrow_batches实现见 Snowflake/Redshift 对比参考适合在内存受限环境下流式处理超大结果集to_sql在 Spark 上不支持、DuckDB 也不支持这与其执行引擎的查询表达方式有关按需转换ODFV所有列出的后端均为本地执行远程执行在所有后端中均未支持矩阵中全为 no这与 Feast 当前按需转换的运行模型一致persist在 Trino 上不支持即 Trino 后端的查询结果无法直接由RetrievalJob持久化回离线存储。核心实现与社区实现在仓库中的位置从源码结构看各后端实现分布如下可据此定位具体查询生成与执行逻辑核心实现位于 offline_stores 目录dask.py ——DaskOfflineStore/DaskRetrievalJob基于 Dask 的本地/分布式 Parquet 计算是最易上手的默认后端bigquery.py ——BigQueryRetrievalJob含to_sql、to_remote_storage实现snowflake.py ——SnowflakeRetrievalJob含to_sql、to_arrow_batches、to_remote_storageredshift.py ——RedshiftRetrievalJob唯一支持 Arrow 批次导出社区实现位于 contrib 目录postgres_offline_store/、spark_offline_store/、trino_offline_store/、ray_offline_store/以及athena_offline_store/、clickhouse_offline_store/、couchbase_offline_store/、mongodb_offline_store/、mssql_offline_store/、oracle_offline_store/等其他remote.py 中的RemoteRetrievalJob支撑远端离线存储通过离线服务代理访问ibis.py 中的IbisRetrievalJob基于 Ibis 提供通用 SQL 引擎后端DuckDB 即走该路径hybrid_offline_store.py 提供混合离线存储能力。选型与配置建议结合功能矩阵与源码事实可以做如下工程判断本地开发/快速起步Dask 后端五个核心方法全支持、RetrievalJob全功能支持且无需外部数仓适合本地 Parquet 数据file源开发调试需要结果导出到数据湖/数仓的大规模训练集优先选择 Snowflake、BigQuery 或 Redshift——它们同时支持export to SQL、export to data warehouse其中 Snowflake/BigQuery 还支持导出到 S3/GCS 数据湖Redshift 额外支持 Arrow 批次流式导出Push 源或特征日志仅 Dask、BigQuery、Snowflake、Redshift、Ray 支持offline_write_batch与write_logged_featuresPostgres/Spark/Trino/Couchbase 后端需另行规划写入链路配置方式离线存储通过feature_store.yaml的offline_store.type字段指定配合各后端专属参数各后端完整配置项请查阅 离线存储文档索引 下对应的 Dask、Snowflake、BigQuery、Redshift、Postgres、Spark、Trino、Ray 等文档稳定性预期社区贡献的 Postgres、Spark、Trino、Ray 等实现不保证稳定、也不保证与核心实现功能对齐生产使用前建议核对上表功能矩阵并以当前仓库源码为准若需跨环境复用后端能力还可评估 远程离线存储 与 混合离线存储 方案。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考