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可选值为aws、gcp、azure。
然后通过 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_URL、SNOWFLAKE_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类)还支持以下完整配置项:
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
type | string | snowflake.offline | 离线存储类型选择器,必须固定为此值 |
config_path | string | ~/.snowsql/config | Snowflake snowsql 配置文件路径,必须是绝对路径(源码注释明确要求不能使用~) |
connection_name | string | None | Snowflake 连接名,通常在~/.snowflake/connections.toml中定义 |
account | string | None | Snowflake 部署标识符(去掉.snowflakecomputing.com后缀) |
user | string | None | Snowflake 用户名 |
password | string | None | Snowflake 密码 |
role | string | None | Snowflake 角色名 |
warehouse | string | None | Snowflake 仓库(warehouse)名 |
authenticator | string | None | Snowflake 认证器名称 |
private_key | string | None | Snowflake 私钥文件路径 |
private_key_content | bytes | None | 以字节形式存储的 Snowflake 私钥 |
private_key_passphrase | string | None | Snowflake 私钥文件口令 |
database | string | 必填 | Snowflake 数据库名(StrictStr,不可省略) |
schema | string | PUBLIC | Snowflake schema 名 |
storage_integration_name | string | None | Snowflake 存储集成名,用于把结果导出到云对象存储 |
blob_export_location | string | None | 数据卸载位置(S3、Google Storage 或 Azure Blob) |
convert_timestamp_columns | bool | None | 导出时是否将时间戳列转换为 Parquet 兼容格式 |
max_file_size | int | 16777216 | 卸载时每个文件的大小上限(字节) |
需要说明的是,schema在配置文件中以schema为别名解析,而源码内部使用schema_字段承接(Field("PUBLIC", alias="schema"));同时model_config设置了extra="allow",允许配置文件携带额外的连接参数。从实现看,database是唯一硬性必填项,其余均可选——但实际连接 Snowflake 通常还需要提供account、user/password或私钥等认证信息。
连接建立的底层逻辑
所有 Snowflake 连接统一由GetSnowflakeConnection(snowflake_utils.py)管理,其要点包括:
- 按类型缓存连接:
offline_store的type值(如snowflake.offline)对应不同的连接缓存键(snowflake.registry、snowflake.offline、snowflake.engine、snowflake.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(源列名到特征表列名的映射)、description、tags、owner以及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)大致如下:
- 根据
entity_df(Pandas DataFrame 或 SQL 字符串)推断实体 schema 与事件时间戳列、时间范围; - 将实体数据上传/物化为临时表;
- 校验实体表中是否包含所有必需的 join key 列;
- 基于
offline_utils生成查询上下文,套用MULTIPLE_FEATURE_VIEW_POINT_IN_TIME_JOINSQL 模板生成最终查询。
该 SQL 模板的关键步骤包括:
- 为每个特征视图构造
entity_dataframeCTE,并生成确定性哈希作为分组键(entity_row_unique_id); - 为每个特征视图生成
subquery,选取时间戳、实体列与特征列(支持field_mapping与full_feature_names命名); - 若配置了
created_timestamp_column,先按(entity, event_timestamp)取MAX(created_timestamp)去重,避免 ASOF JOIN 在并列时间戳时产生不稳定结果(模板注释引用了 Snowflake 官方文档关于 ASOF JOIN ties 行为的说明); - 使用 Snowflake 原生ASOF JOIN(
MATCH_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_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 source) | yes |
write_logged_features | 将日志特征持久化到离线存储(特征日志) | yes |
其中前三个方法都返回该离线存储特有的RetrievalJob(如SnowflakeRetrievalJob)。SnowflakeRetrievalJob支持的功能矩阵如下:
| 功能 | Snowflake |
|---|---|
| 导出为 dataframe | yes |
| 导出为 arrow table | yes |
| 导出为 arrow batches | yes |
| 导出为 SQL | yes |
| 导出到数据湖(S3、GCS 等) | yes |
| 导出到数据仓库 | yes |
| 导出为 Spark dataframe | yes |
| 本地执行 Python 型 on-demand transforms | yes |
| 远程执行 Python 型 on-demand transforms | no |
| 将结果持久化到离线存储 | yes |
| 执行前预览查询计划 | yes |
| 读取分区数据 | yes |
与其他离线存储(Dask、BigQuery、Redshift 等)的完整对比可参考 功能矩阵。可以看到 Snowflake 离线存储在导出维度上覆盖非常全面,是少数同时支持“导出为 Spark dataframe”与“导出到数据湖”的实现之一。
SnowflakeRetrievalJob 导出能力与源码解析
SnowflakeRetrievalJob(snowflake.py)封装了检索结果的各种消费方式。它是懒执行的:查询字符串保存在内部,直到调用导出方法才真正下发执行,因此支持“执行前预览查询计划”(调用to_sql()即可拿到将要执行的 SQL,无需真正运行)。
各导出方法的底层实现要点:
- 导出为 DataFrame:
to_df()内部执行execute_snowflake_statement(...).fetch_pandas_all();若特征包含 Array 类型,会把 Snowflake 返回的 JSON 字符串json.loads还原为列表(_to_df_internal中逐列处理 Array(String/Bytes/Int32/Int64/UnixTimestamp/Float64/Float32/Bool)); - 导出为 Arrow:
to_arrow()使用fetch_arrow_all(force_return_table=True); - 导出为批次:
to_arrow_batches()/to_pandas_batches()先把查询结果落为临时表(CREATE TEMPORARY TABLE),再通过fetch_arrow_batches()/fetch_pandas_batches()分批拉取,适合内存放不下的大数据集; - 导出为 SQL:
to_sql()直接返回查询字符串; - 导出到数据仓库:
to_snowflake(table_name, ...)支持allow_overwrite(CREATE OR REPLACE)与temporary(TEMPORARY)选项;若带 on-demand 特征视图则先本地转换再write_pandas写入; - 导出为 Spark DataFrame:
to_spark_df(spark_session)开启 PySpark Arrow 支持后,把to_pandas_batches()的批次逐批createDataFrame再unionAll合并; - 导出到数据湖:
to_remote_storage()需要同时配置storage_integration_name与blob_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_into从COPY 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)加载。其中压缩算法支持gzip与snappy,上传并行度默认 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),仅供参考