news 2026/9/17 13:04:52

Feast Snowflake 离线存储(Snowflake Offline Store)配置与实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Feast Snowflake 离线存储(Snowflake Offline Store)配置与实战指南

Feast Snowflake 离线存储(Snowflake Offline Store)配置与实战指南

【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast

本文围绕 Feast 开源特征平台中面向 Snowflake 的离线存储实现展开,系统讲解其核心能力、feature_store.yaml完整配置、SnowflakeSource数据源定义、Point-in-Time 正确性连接(ASOF JOIN)的执行原理、单引号限制规避,以及SnowflakeRetrievalJob的各类结果导出能力。读完本文,你将掌握如何在 Feast 中完成基于 Snowflake 的离线特征仓库搭建、历史特征检索与结果导出,并能结合仓库源码理解底层实现细节。

Snowflake 离线存储能做什么

Feast 的SnowflakeOfflineStore是官方核心离线存储实现之一,负责在训练/离线阶段读取 SnowflakeSources(Snowflake 表或视图、或一段 SQL 查询),并完成特征检索工作。它的设计有几个关键特点:

  • 所有 Join 都在 Snowflake 内部完成:无论是特征表之间的关联,还是特征表与实体表(entity dataframe)之间的 Point-in-Time 连接,均由 Snowflake 的查询引擎执行,不需要把特征数据搬运到本地。
  • 实体数据(entity dataframe)有两种提供方式
    • 直接传入Pandas DataFrame:Feast 会把它上传到 Snowflake 临时表,以完成后续 join 操作;
    • 直接传入SQL 查询字符串:该查询的结果会被当作实体数据集参与连接。

这一点在 snowflake.py 的_upload_entity_df函数中有直接体现:Pandas DataFrame 走write_pandas(..., create_temp_table=True)写入临时表;SQL 字符串则执行CREATE TEMPORARY TABLE ... AS (查询)创建临时表。两种路径最终都把实体数据物化在 Snowflake 内部。

安装与项目初始化

使用 Snowflake 离线存储前,需要安装带 Snowflake 依赖的 Feast SDK:

pip install 'feast[snowflake]'

如果使用的是文件型注册表(file based registry,如本地registry.db或 S3/GCS 上的 registry 文件),还需要同时安装对应的云厂商扩展:

pip install 'feast[snowflake, aws]' # AWS pip install 'feast[snowflake, gcp]' # GCP pip install 'feast[snowflake, azure]' # Azure

其中CLOUD可选值为awsgcpazure

然后通过 Feast CLI 以 Snowflake 模板初始化特征仓库:

feast init -t snowflake

初始化后,仓库中会生成一个 feature_store.yaml 模板,其中已预置offline_store(类型为snowflake.offline)、batch_engine(类型为snowflake.engine,用于 Snowflake 上的批式物化计算)和online_store(类型为snowflake.online,用于在线特征存储)三段配置,占位符形如SNOWFLAKE_DEPLOYMENT_URLSNOWFLAKE_USER,替换为真实值即可。

feature_store.yaml 完整配置

官方文档给出的最小可用配置如下:

project: my_feature_repo registry: data/registry.db provider: local offline_store: type: snowflake.offline account: snowflake_deployment.us-east-1 user: user_login password: user_password role: SYSADMIN warehouse: COMPUTE_WH database: FEAST schema: PUBLIC

其中account是 Snowflake 部署标识符,例如snowflake_deployment.us-east-1不要.snowflakecomputing.com后缀;schema默认值为PUBLIC

除上述字段外,SnowflakeOfflineStoreConfig(定义于 snowflake.py 的SnowflakeOfflineStoreConfig类)还支持以下完整配置项:

配置项类型默认值说明
typestringsnowflake.offline离线存储类型选择器,必须固定为此值
config_pathstring~/.snowsql/configSnowflake snowsql 配置文件路径,必须是绝对路径(源码注释明确要求不能使用~
connection_namestringNoneSnowflake 连接名,通常在~/.snowflake/connections.toml中定义
accountstringNoneSnowflake 部署标识符(去掉.snowflakecomputing.com后缀)
userstringNoneSnowflake 用户名
passwordstringNoneSnowflake 密码
rolestringNoneSnowflake 角色名
warehousestringNoneSnowflake 仓库(warehouse)名
authenticatorstringNoneSnowflake 认证器名称
private_keystringNoneSnowflake 私钥文件路径
private_key_contentbytesNone以字节形式存储的 Snowflake 私钥
private_key_passphrasestringNoneSnowflake 私钥文件口令
databasestring必填Snowflake 数据库名(StrictStr,不可省略)
schemastringPUBLICSnowflake schema 名
storage_integration_namestringNoneSnowflake 存储集成名,用于把结果导出到云对象存储
blob_export_locationstringNone数据卸载位置(S3、Google Storage 或 Azure Blob)
convert_timestamp_columnsboolNone导出时是否将时间戳列转换为 Parquet 兼容格式
max_file_sizeint16777216卸载时每个文件的大小上限(字节)

需要说明的是,schema在配置文件中以schema为别名解析,而源码内部使用schema_字段承接(Field("PUBLIC", alias="schema"));同时model_config设置了extra="allow",允许配置文件携带额外的连接参数。从实现看,database是唯一硬性必填项,其余均可选——但实际连接 Snowflake 通常还需要提供accountuser/password或私钥等认证信息。

连接建立的底层逻辑

所有 Snowflake 连接统一由GetSnowflakeConnection(snowflake_utils.py)管理,其要点包括:

  • 按类型缓存连接offline_storetype值(如snowflake.offline)对应不同的连接缓存键(snowflake.registrysnowflake.offlinesnowflake.enginesnowflake.online分别缓存),同一配置的重复获取不会重建连接;
  • 支持从 snowsql 配置文件读取参数:若config_path指向的配置存在connections.feast_offline_store段落,会先读取该段落,再用 feature_store.yaml 中的显式配置覆盖(配置项中v is not None的键才会覆盖);
  • 建立连接时统一设置时区:连接建立后会执行ALTER SESSION SET TIMEZONE = 'UTC',保证时间语义一致;
  • 支持私钥认证:配置private_key/private_key_content时,会通过parse_private_key_path读取 PEM 私钥并转为 PKCS8 DER 字节(支持口令private_key_passphrase),详见 snowflake_utils.py。

另外,snowflake.py 的_get_snowflake_conn支持通过数据源的connection_ref覆盖连接参数(account、user、password、database、warehouse、role、schema、authenticator、private_key 等),实现“数据源级”的凭据覆盖。

定义 Snowflake 数据源(SnowflakeSource)

SnowflakeSource(snowflake_source.py)用于声明“从哪张表/哪个查询读取特征”。它有两种指定方式,二者必须且只能提供其一(源码中会校验:两者都缺失或同时提供都会抛ValueError)。

方式一:表引用

from feast import SnowflakeSource my_snowflake_source = SnowflakeSource( database="FEAST", schema="PUBLIC", table="FEATURE_TABLE", )

方式二:SQL 查询

from feast import SnowflakeSource my_snowflake_source = SnowflakeSource( query=""" SELECT timestamp_column AS "ts", "created", "f1", "f2" FROM `FEAST.PUBLIC.FEATURE_TABLE` """, )

使用时注意 Snowflake 对表名、列名的大小写与引号处理规则(quote identifiers):例如列名默认会被转换为大写,需要用小写列名时须加双引号。源码中get_table_query_string生成引用时会为 database/schema/table 分别加双引号,并支持“只给 schema+table”“只给 table”或“只给 query”三种缺省组合,缺省部分会回退到离线存储配置中的database/schema(见pull_latest_from_table_or_query_qualify_snowflake_from_expression的实现)。

SnowflakeSource还支持以下参数:timestamp_field(事件时间戳字段,用于 Point-in-Time join)、created_timestamp_column(行创建时间戳,用于去重)、field_mapping(源列名到特征表列名的映射)、descriptiontagsowner以及connection_ref(凭据引用)。当只提供database+table而未提供schema时,默认使用PUBLIC

在类型支持方面,Snowflake 数据源支持全部八种 Feast 基础类型;数组(Array)类型也支持,但不支持类型推断。源码中的snowflake_type_code_map展示了 Snowflake 类型码到 Feast 类型的映射(NUMBER、DOUBLE、VARCHAR、DATE、TIMESTAMP 系列、ARRAY、BINARY、BOOLEAN),而 VARIANT、OBJECT、TIME 类型会被视为不支持并抛错(提示转换为 VARCHAR)。

历史特征检索与 Point-in-Time 连接

离线特征检索的核心入口是get_historical_features,它返回一个懒执行的SnowflakeRetrievalJob。整个流程(见 snowflake.py 的get_historical_features)大致如下:

  1. 根据entity_df(Pandas DataFrame 或 SQL 字符串)推断实体 schema 与事件时间戳列、时间范围;
  2. 将实体数据上传/物化为临时表;
  3. 校验实体表中是否包含所有必需的 join key 列;
  4. 基于offline_utils生成查询上下文,套用MULTIPLE_FEATURE_VIEW_POINT_IN_TIME_JOINSQL 模板生成最终查询。

该 SQL 模板的关键步骤包括:

  • 为每个特征视图构造entity_dataframeCTE,并生成确定性哈希作为分组键(entity_row_unique_id);
  • 为每个特征视图生成subquery,选取时间戳、实体列与特征列(支持field_mappingfull_feature_names命名);
  • 若配置了created_timestamp_column,先按(entity, event_timestamp)MAX(created_timestamp)去重,避免 ASOF JOIN 在并列时间戳时产生不稳定结果(模板注释引用了 Snowflake 官方文档关于 ASOF JOIN ties 行为的说明);
  • 使用 Snowflake 原生ASOF JOINMATCH_CONDITION (e."entity_timestamp" >= v."event_timestamp"))完成 Point-in-Time 正确性连接;
  • 若特征视图配置了 TTL,则通过TIMESTAMPADD(second, -ttl, ...)过滤过期特征;
  • 最终将各特征视图的 TTL 结果按entity_row_unique_idLEFT JOIN 回实体数据集。

这也解释了为何“所有 Join 都发生在 Snowflake 内部”:从实体表上传到最终的 ASOF JOIN,全部由 Snowflake 执行。

单引号限制与规避方法

Feast 在使用 Snowflake 时对 SQL 查询字符串有一个已知限制:尽量规避 SQL 查询字符串中的单引号'。例如以下查询会失败:

SELECT some_column FROM some_table WHERE other_column = 'value'

'value'中的单引号会导致在 Snowflake 中执行失败。官方推荐改用成对的美元符$$value$$包裹字符串常量(即 Snowflake 文档中的 dollar-quoted string constants):

SELECT some_column FROM some_table WHERE other_column = $$value$$

在实体数据集、特征源或任何以查询字符串形式传入 Feast 的 SQL 中,都应优先使用这种写法,避免字符串常量中出现裸单引号。

功能矩阵:离线存储与检索任务能力一览

Feast 定义了OfflineStore接口的五个核心方法(详见 overview.md),Snowflake 离线存储对它们的支持情况如下:

方法说明Snowflake
get_historical_featuresPoint-in-Time 正确性连接,检索历史特征yes
pull_latest_from_table_or_query检索最新特征值(用于物化到在线存储)yes
pull_all_from_table_or_query检索已保存的数据集(saved dataset)yes
offline_write_batch将 DataFrame 持久化到离线存储(主要用于 push source)yes
write_logged_features将日志特征持久化到离线存储(特征日志)yes

其中前三个方法都返回该离线存储特有的RetrievalJob(如SnowflakeRetrievalJob)。SnowflakeRetrievalJob支持的功能矩阵如下:

功能Snowflake
导出为 dataframeyes
导出为 arrow tableyes
导出为 arrow batchesyes
导出为 SQLyes
导出到数据湖(S3、GCS 等)yes
导出到数据仓库yes
导出为 Spark dataframeyes
本地执行 Python 型 on-demand transformsyes
远程执行 Python 型 on-demand transformsno
将结果持久化到离线存储yes
执行前预览查询计划yes
读取分区数据yes

与其他离线存储(Dask、BigQuery、Redshift 等)的完整对比可参考 功能矩阵。可以看到 Snowflake 离线存储在导出维度上覆盖非常全面,是少数同时支持“导出为 Spark dataframe”与“导出到数据湖”的实现之一。

SnowflakeRetrievalJob 导出能力与源码解析

SnowflakeRetrievalJob(snowflake.py)封装了检索结果的各种消费方式。它是懒执行的:查询字符串保存在内部,直到调用导出方法才真正下发执行,因此支持“执行前预览查询计划”(调用to_sql()即可拿到将要执行的 SQL,无需真正运行)。

各导出方法的底层实现要点:

  • 导出为 DataFrameto_df()内部执行execute_snowflake_statement(...).fetch_pandas_all();若特征包含 Array 类型,会把 Snowflake 返回的 JSON 字符串json.loads还原为列表(_to_df_internal中逐列处理 Array(String/Bytes/Int32/Int64/UnixTimestamp/Float64/Float32/Bool));
  • 导出为 Arrowto_arrow()使用fetch_arrow_all(force_return_table=True)
  • 导出为批次to_arrow_batches()/to_pandas_batches()先把查询结果落为临时表(CREATE TEMPORARY TABLE),再通过fetch_arrow_batches()/fetch_pandas_batches()分批拉取,适合内存放不下的大数据集;
  • 导出为 SQLto_sql()直接返回查询字符串;
  • 导出到数据仓库to_snowflake(table_name, ...)支持allow_overwriteCREATE OR REPLACE)与temporaryTEMPORARY)选项;若带 on-demand 特征视图则先本地转换再write_pandas写入;
  • 导出为 Spark DataFrameto_spark_df(spark_session)开启 PySpark Arrow 支持后,把to_pandas_batches()的批次逐批createDataFrameunionAll合并;
  • 导出到数据湖to_remote_storage()需要同时配置storage_integration_nameblob_export_location,实现上先建临时表,再用COPY INTO '<export_path>/<table>' ... STORAGE_INTEGRATION = ... FILE_FORMAT = (TYPE = PARQUET)把数据以 Parquet 形式卸载到 S3/GCS/Azure 等外部位置,并可通过max_file_size控制单文件大小;针对 AWS GovCloud 还会把s3gov://前缀规范化为s3://
  • 持久化到离线存储persist(storage)把结果写为SavedDatasetSnowflakeStorage指向的 Snowflake 表,供后续pull_all_from_table_or_query直接读取。

单元测试 test_snowflake.py 验证了to_remote_storage的完整行为:它会先调用to_snowflake创建临时表,再调用_get_file_names_from_copy_intoCOPY INTO ... DETAILED_OUTPUT = TRUE的结果中提取FILE_NAME列,最终返回形如s3://.../file的文件路径列表。

写入与特征日志

offline_write_batch用于把 PyArrow 表写入离线存储(push source 场景):它会先校验输入表的列与顺序必须与批数据源的 schema 一致,类型不符时通过cast_arrow_table_to_schema转换,最后调用write_pandas(自动建表)写入feature_view.batch_source.table指向的表。

write_logged_features用于特征日志持久化:数据若为本地 Parquet 路径则走write_parquet,否则走write_pandas,均写入SnowflakeLoggingDestination指定的table_name,并支持自动建表。

在数据写入路径上,snowflake_utils.py 的write_pandas采用“先落本地 Parquet → PUT 上传临时 stage → COPY INTO 目标表”的经典管道:DataFrame 按chunk_size分块导出为 Parquet(use_deprecated_int96_timestamps=True保证时间戳精度),上传到自动创建的临时 stage,再通过infer_schema推断列类型自动建表(auto_create_table=True时),最后以COPY INTO ... FILE_FORMAT=(TYPE=PARQUET)加载。其中压缩算法支持gzipsnappy,上传并行度默认 4。

小结

Feast 的 Snowflake 离线存储是一个能力相当完整的核心实现:它把 Point-in-Time 正确性连接的复杂逻辑(ASOF JOIN、created_timestamp 去重、TTL 过滤)全部下推到 Snowflake 执行,实体数据既可用 Pandas DataFrame(自动上传临时表)也可用 SQL 查询提供;SnowflakeRetrievalJob覆盖了从 DataFrame、Arrow、SQL、数据湖到 Spark DataFrame 的几乎所有主流导出形态,同时支持特征日志、批量写入与离线数据监控。

实际落地时需重点注意:SQL 查询字符串中避免裸单引号(改用$$...$$)、account不要带域名后缀、config_path使用绝对路径、大数据集导出优先使用批次 API(arrow/pandas batches)或数据湖卸载(配置storage_integration_name+blob_export_location)以控制内存与文件大小。相关实现细节可继续阅读 snowflake.py、snowflake_source.py、snowflake_utils.py 与模板配置 feature_store.yaml。

【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 13:04:04

Win10越用越慢?13种优化方法从启动项到电源策略

Windows 10用久了变慢几乎是躲不掉的事&#xff0c;但十有八九不是硬件真不行&#xff0c;而是系统里堆了太多默认开着、默认运行、默认保留的东西。我手上这台办公本用了四年多&#xff0c;中间只加过一条内存&#xff0c;系统一直没重装&#xff0c;靠的就是隔一段时间按固定…

作者头像 李华
网站建设 2026/9/17 13:02:55

UML用例图从入门到实战:参与者、用例关系与软考考点解析

1. 用例图到底解决什么问题UML设计系列写到这里&#xff0c;前面几篇分别聊了类图、对象图这些偏静态结构的图&#xff0c;这一篇轮到用例图。说实话&#xff0c;用例图经常被当成UML里最“简单”的一张图&#xff0c;很多同学画的时候也就拖几个椭圆、拉几根线&#xff0c;感觉…

作者头像 李华
网站建设 2026/9/17 12:59:20

传奇C++源码深度解析:从服务端主循环到客户端渲染

如果你手里也有这么一份名为“传奇源代码cpp版本”的压缩包&#xff0c;那这篇文章刚好可以拿去参考。我最早是因为想搞清楚“一个网游玩起来到底在收发包、跑逻辑、画画面时做了什么”&#xff0c;才去翻的这套代码。当时没人带&#xff0c;全靠对着源码一行行猜&#xff0c;再…

作者头像 李华