配置与实战指南)
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是官方核心离线存储实现之一负责在训练/离线阶段读取 SnowflakeSourcesSnowflake 表或视图、或一段 SQL 查询并完成特征检索工作。它的设计有几个关键特点所有 Join 都在 Snowflake 内部完成无论是特征表之间的关联还是特征表与实体表entity dataframe之间的 Point-in-Time 连接均由 Snowflake 的查询引擎执行不需要把特征数据搬运到本地。实体数据entity dataframe有两种提供方式直接传入Pandas DataFrameFeast 会把它上传到 Snowflake 临时表以完成后续 join 操作直接传入SQL 查询字符串该查询的结果会被当作实体数据集参与连接。这一点在 snowflake.py 的_upload_entity_df函数中有直接体现Pandas DataFrame 走write_pandas(..., create_temp_tableTrue)写入临时表SQL 字符串则执行CREATE TEMPORARY TABLE ... AS (查询)创建临时表。两种路径最终都把实体数据物化在 Snowflake 内部。安装与项目初始化使用 Snowflake 离线存储前需要安装带 Snowflake 依赖的 Feast SDKpip 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类还支持以下完整配置项配置项类型默认值说明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 Blobconvert_timestamp_columnsboolNone导出时是否将时间戳列转换为 Parquet 兼容格式max_file_sizeint16777216卸载时每个文件的大小上限字节需要说明的是schema在配置文件中以schema为别名解析而源码内部使用schema_字段承接Field(PUBLIC, aliasschema)同时model_config设置了extraallow允许配置文件携带额外的连接参数。从实现看database是唯一硬性必填项其余均可选——但实际连接 Snowflake 通常还需要提供account、user/password或私钥等认证信息。连接建立的底层逻辑所有 Snowflake 连接统一由GetSnowflakeConnectionsnowflake_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 数据源SnowflakeSourceSnowflakeSourcesnowflake_source.py用于声明“从哪张表/哪个查询读取特征”。它有两种指定方式二者必须且只能提供其一源码中会校验两者都缺失或同时提供都会抛ValueError。方式一表引用from feast import SnowflakeSource my_snowflake_source SnowflakeSource( databaseFEAST, schemaPUBLIC, tableFEATURE_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 分别加双引号并支持“只给 schematable”“只给 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凭据引用。当只提供databasetable而未提供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_dfPandas 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 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 valuevalue中的单引号会导致在 Snowflake 中执行失败。官方推荐改用成对的美元符$$value$$包裹字符串常量即 Snowflake 文档中的 dollar-quoted string constantsSELECT some_column FROM some_table WHERE other_column $$value$$在实体数据集、特征源或任何以查询字符串形式传入 Feast 的 SQL 中都应优先使用这种写法避免字符串常量中出现裸单引号。功能矩阵离线存储与检索任务能力一览Feast 定义了OfflineStore接口的五个核心方法详见 overview.mdSnowflake 离线存储对它们的支持情况如下方法说明Snowflakeget_historical_featuresPoint-in-Time 正确性连接检索历史特征yespull_latest_from_table_or_query检索最新特征值用于物化到在线存储yespull_all_from_table_or_query检索已保存的数据集saved datasetyesoffline_write_batch将 DataFrame 持久化到离线存储主要用于 push sourceyeswrite_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 导出能力与源码解析SnowflakeRetrievalJobsnowflake.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_tableTrue)导出为批次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()的批次逐批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分块导出为 Parquetuse_deprecated_int96_timestampsTrue保证时间戳精度上传到自动创建的临时 stage再通过infer_schema推断列类型自动建表auto_create_tableTrue时最后以COPY INTO ... FILE_FORMAT(TYPEPARQUET)加载。其中压缩算法支持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使用绝对路径、大数据集导出优先使用批次 APIarrow/pandas batches或数据湖卸载配置storage_integration_nameblob_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),仅供参考