Feast DuckDB 离线存储(Offline Store)配置与实战指南
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
DuckDB 离线存储(DuckDB OfflineStore)是 Feast 中面向文件型数据源的轻量级离线存储实现:它可以直接读取 Parquet 与 Delta 格式的FileSource,并在底层借助 Ibis 框架把离线存储操作翻译为 DuckDB 查询,从而在本地即可完成 Point-in-Time 正确性联结(PIT join)、特征物化(materialization)与训练数据集导出。读完本文,你将掌握 DuckDB 离线存储的安装、feature_store.yaml配置、FileSource定义方法、完整功能矩阵,以及其源码实现原理与测试验证方式。
一、DuckDB 离线存储是什么
根据官方文档 DuckDB offline store 的说明,DuckDB 离线存储的核心定位是:
- 面向文件数据源:提供对
FileSource的读取支持,可读取Parquet与Delta两种格式(文件可位于本地磁盘或 S3 上,详见 File source)。 - 基于 Ibis 构建:DuckDB 离线存储使用 Ibis 框架(Python 的通用 DataFrame / 数据库交互库)作为表达式引擎,将 Feast 的离线存储操作统一转换为 DuckDB SQL 查询执行,因此其核心逻辑与 Ibis 离线存储家族共享同一套实现(见下文源码分析)。
- 实体 DataFrame 灵活输入:
get_historical_features的实体数据(entity dataframe)可以直接以 Pandas DataFrame 形式传入,无需先落盘。
从实现上看,duckdb.py 中的DuckDBOfflineStore类继承自 Feast 的OfflineStore抽象接口,并声明supports_filter_by_created_timestamp = True,即支持按创建时间戳(created timestamp)过滤数据。
二、快速开始:安装与依赖
启用 DuckDB 离线存储只需安装 Feast 的 duckdb 扩展:
pip install 'feast[duckdb]'从当前仓库的 CI 依赖锁定文件 py3.10-ci-requirements.txt 可以看到,该扩展的核心依赖版本为:
duckdb==1.5.5ibis-framework[duckdb, mssql, oracle]==12.0.0(Ibis 框架,DuckDB 后端)deltalake==0.25.5(Delta Lake 读写)pyarrow(Arrow 数据交换)
这些依赖确保了 Parquet / Delta 读取、Arrow 导出以及 Ibis 表达式执行的能力。若你的数据源只用到本地 Parquet 文件,安装后即可开箱使用;若涉及 S3 上的 Delta 表,则需要额外的云存储凭证配置(见下文)。
三、feature_store.yaml 配置示例
在原文档中,DuckDB 离线存储的feature_store.yaml配置如下:
project: my_project registry: data/registry.db provider: local offline_store: type: duckdb online_store: path: data/online_store.db这份配置的要点:
| 配置项 | 取值 | 说明 |
|---|---|---|
project | my_project | 特征仓库所属项目名 |
registry | data/registry.db | 注册表(Registry)存储位置,此处为本地 SQLite 文件 |
provider | local | 本地 Provider,适合本地开发与测试 |
offline_store.type | duckdb | 指定离线存储类型为 DuckDB |
online_store.path | data/online_store.db | 在线存储为本地 SQLite,物化后的特征写入该文件 |
进阶配置项:staging_location
在基础示例之外,DuckDBOfflineStoreConfig还支持两个可选参数(见 duckdb.py):
staging_location(可选):远程暂存位置(如s3://bucket/staging)。当配置了该参数后,IbisRetrievalJob才支持"导出到远程存储"(export to remote storage)能力,结果会以 Delta 格式写入该位置。staging_location_endpoint_override(可选):用于覆盖 S3 端点地址(例如指向 MinIO 等 S3 兼容服务)。
一个带 staging 配置的示例(参考测试配置 duckdb_repo_configuration.py):
project: my_project registry: data/registry.db provider: local offline_store: type: duckdb staging_location: s3://my-bucket/staging staging_location_endpoint_override: http://minio:9000 online_store: path: data/online_store.db注意:type: duckdb的解析逻辑由DuckDBOfflineStoreConfig中的type: StrictStr = "duckdb"默认值承接,pydantic 会校验该字段类型。
四、定义文件数据源(FileSource)
DuckDB 离线存储读取的数据源类型为FileSource,支持Parquet和Delta两种格式。以 Parquet 为例,数据源定义如下(见 File source):
from feast import FileSource from feast.data_format import ParquetFormat parquet_file_source = FileSource( file_format=ParquetFormat(), path="file:///feast/customer.parquet", )Delta 格式同理,只需将ParquetFormat换成DeltaFormat。文件数据源支持全部八种基础类型及其对应的数组类型(array types),可用于定义 Batch Feature View 的原始数据。
源码视角:数据源如何被读取
在 duckdb.py 的_read_data_source函数中,读取逻辑按格式分支:
- Parquet(
ParquetFormat,或未指定格式但路径以.parquet结尾):调用ibis.read_parquet(data_source.path); - Delta(
DeltaFormat):调用ibis.read_delta(data_source.path); - S3 凭证注入:若路径以
s3://开头且配置了凭证(resolve_credentials()),会构造storage_options注入AWS_ACCESS_KEY_ID、AWS_SECRET_ACCESS_KEY,并支持s3_endpoint_override(AWS_ENDPOINT_URL)与AWS_SESSION_TOKEN; - 另外,源码中还支持通过 PyIceberg 读取 Iceberg Catalog 数据源(
IcebergSource),读取后以ibis.memtable载入内存。
也就是说,只要你的FileSource指向合法的 Parquet / Delta 文件(本地或 S3),DuckDB 离线存储即可在get_historical_features、pull_latest_from_table_or_query等操作中直接完成读取。
五、功能矩阵:支持什么、不支持什么
OfflineStore 接口方法支持情况
Feast 通过OfflineStore接口暴露 5 个核心方法(各方法含义详见 offline-stores overview),DuckDB 全部支持:
| 方法 | 用途 | DuckDB |
|---|---|---|
get_historical_features | Point-in-Time 正确性联结,取历史特征 | yes |
pull_latest_from_table_or_query | 取最新特征值,用于物化到在线存储 | yes |
pull_all_from_table_or_query | 取回已保存的数据集(saved dataset) | yes |
offline_write_batch | 将 DataFrame 持久化到离线存储(主要用于 push sources) | yes |
write_logged_features | 将记录的特征(logged features)持久化到离线存储 | yes |
IbisRetrievalJob 能力矩阵
上述前三个方法返回的 Retrieval Job 在 DuckDB 场景下为IbisRetrievalJob(共享自 Ibis 离线存储家族),其能力矩阵如下:
| 能力 | DuckDB |
|---|---|
| export to dataframe | yes |
| export to arrow table | yes |
| export to arrow batches | no |
| export to SQL | no |
| export to data lake (S3, GCS, etc.) | no |
| export to data warehouse | no |
| export as Spark dataframe | no |
| local execution of Python-based on-demand transforms | yes |
| remote execution of Python-based on-demand transforms | no |
| persist results in the offline store | yes |
| preview the query plan before execution | no |
| read partitioned data | yes |
与其它离线存储的完整对比,可查阅 功能矩阵。
要点解读:
- 本地强、远程弱:DuckDB 离线存储擅长本地执行(dataframe / arrow 导出、本地 on-demand transforms、分区读取、结果持久化),但不支持导出到 SQL、数据湖、数据仓库或 Spark DataFrame,也没有查询计划预览(preview)能力;
- 远程执行受限:Python on-demand transforms 仅在本地执行;远程存储导出(
to_remote_storage)只有在配置了staging_location时才可用(见下文源码分析); - Arrow batches 不支持:大内存数据集需分批导出时,应考虑其它存储(如 Redshift 支持)。
六、源码深度:DuckDBOfflineStore 实现原理
类结构与方法委托
duckdb.py 中DuckDBOfflineStore的五个核心方法并非各自实现 SQL,而是委托给 Ibis 共享实现:
pull_latest_from_table_or_query→pull_latest_from_table_or_query_ibisget_historical_features→get_historical_features_ibispull_all_from_table_or_query→pull_all_from_table_or_query_ibisoffline_write_batch→offline_write_batch_ibiswrite_logged_features→write_logged_features_ibis
这些共享函数位于 ibis.py,通过参数data_source_reader=_read_data_source、data_source_writer=_write_data_source将"如何读写文件"的具体逻辑注入进来,实现存储差异与查询逻辑的解耦。
写入逻辑:append 与 overwrite
_write_data_source 支持两种写入模式:
- overwrite:将表转为 PyArrow 后,若目标是
.parquet文件则write_table,否则视为目录用write_to_dataset;当目标已存在且allow_overwrite=False时,抛出SavedDatasetLocationAlreadyExists错误,避免误覆盖已有数据集; - append:读取现有 Parquet 并做 schema 对齐(cast 到旧 schema)后
concat_tables再写回;Delta 格式则调用table.to_delta(..., mode="append"/"overwrite"),通过deltalake库完成。
Retrieval Job 的导出能力
IbisRetrievalJob(见 ibis.py)关键实现:
_to_df_internal:self.table.execute(),即让 Ibis 把表达式交给 DuckDB 执行并返回 Pandas DataFrame;_to_arrow_internal:self.table.to_pyarrow(),返回 Arrow Table;persist:调用data_source_writer(即_write_data_source)以overwrite模式把结果写回离线存储,用于保存训练数据集(saved dataset);supports_remote_storage_export:仅当staging_location非空时为 True;to_remote_storage:把结果以 Delta 格式写入staging_location/{uuid}目录,并返回已写入的 S3 文件列表。
这解释了功能矩阵中"export to data lake = no / persist results = yes"的组合:DuckDB 场景下,数据写入目标由FileSource决定(本地或 S3 上的文件),而staging_location只开启"导出到远程存储"这一条独立路径。
数据质量监控(文档之外的源码能力)
值得补充的是,当前仓库中的DuckDBOfflineStore还实现了完整的数据质量监控接口(compute_monitoring_metrics、get_monitoring_max_timestamp、ensure_monitoring_tables、save_monitoring_metrics、query_monitoring_metrics、clear_monitoring_baseline,见 duckdb.py)。监控计算通过 DuckDB 原生 SQL 完成:
- 数值型特征统计:
COUNT、AVG、STDDEV_SAMP、MIN/MAX、PERCENTILE_CONT(p50/p75/p90/p95/p99)及直方图分箱; - 类别型特征统计:Top-N 取值频率、去重计数、空值率;
- 监控结果以 Parquet 文件形式按主键做 upsert 落盘,位于仓库根目录的
monitoring/目录下,支持按项目、特征视图、特征名过滤查询。
这意味着 DuckDB 离线存储不仅是"读取训练数据"的通道,也可作为特征监控(feature monitoring)的本地计算底座。
七、测试与验证
仓库的通用测试配置 duckdb_repo_configuration.py 提供了三种数据源变体来验证 DuckDB 离线存储:
DuckDBDataSourceCreator:本地 Parquet 文件数据源 +DuckDBOfflineStoreConfig();DuckDBDeltaDataSourceCreator:本地 Delta 文件数据源;DuckDBDeltaS3DataSourceCreator:S3 上的 Delta 数据源,并配置staging_location="s3://test/staging"与staging_location_endpoint_override(指向本地模拟 S3 端点),用于验证远程存储导出与 Delta S3 读取路径。
这些 creator 会被 Feast 的 universal 测试框架用来跑历史特征检索、物化、数据集保存等端到端用例,是验证 DuckDB 离线存储行为的最直接参考。你也可以在本地特性仓库(feature repo)中通过feast apply与feast materialize命令,配合 quickstart 中的标准流程验证本文配置。
八、适用场景与注意事项
推荐场景:
- 本地 / 单机开发与 CI 测试:无需外部集群即可完成 PIT join 与训练数据集构建;
- 中小规模 Parquet / Delta 数据集的离线特征检索;
- 与 SQLite 在线存储组合,形成"零外部依赖"的端到端特征存储链路;
- 在本地运行数据质量监控计算。
注意事项与限制:
- 功能矩阵中的
no项(Arrow batches、SQL/数据湖/数仓/Spark 导出、远程 on-demand transforms、查询计划预览)为设计限制,如需这些能力请评估 overview 中其它离线存储(如 BigQuery、Snowflake、Redshift); - 读取 S3 上的 Parquet / Delta 时需正确配置凭证与
s3_endpoint_override; staging_location仅影响"导出到远程存储"能力,不影响常规的 DataFrame 返回;- 版本能力以当前仓库为准:依赖锁定版本为
duckdb==1.5.5、ibis-framework==12.0.0、deltalake==0.25.5。
综上,DuckDB 离线存储是 Feast 生态中"本地优先、文件优先"的轻量离线方案:配置简单(一行type: duckdb)、对 Parquet/Delta 原生友好、完整支持 Point-in-Time 联结与物化取数,并通过 Ibis 共享实现保证了与其它离线存储一致的接口语义,是本地原型验证与中小规模特征工程的实用选择。
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考