ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Feast Snowpark Python UDF 模块解析:Snowflake 特征物化的类型序列化内核

Feast Snowpark Python UDF 模块解析:Snowflake 特征物化的类型序列化内核 Feast Snowpark Python UDF 模块解析Snowflake 特征物化的类型序列化内核【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feastFeast 的feast.infra.utils.snowflake.snowpark包是 Feast 在 Snowflake 上执行批量特征物化materialization的核心支撑它以一组注册在 Snowflake 中的 Python UDFSnowpark Python UDF为桥梁把 Snowflake 表里的原生 SQL 类型BINARY、VARCHAR、NUMBER、ARRAY 等逐一转换成 Feast 的 protobuf 序列化形式并负责把实体键entity key编码为可供在线存储查询的二进制键。本文以该模块的 API 文档feast.infra.utils.snowflake.snowpark.rst为骨架结合其真实实现、部署 SQL 模板与计算引擎调用链讲清每支 UDF 的作用、类型映射规则、部署与清理机制以及运行时版本的自定义方式。模块定位一篇 RST 指向的实质代码该 RST 文档是 Sphinxautomodule风格的 API 文档条目它声明的文档对象有两个层次feast.infra.utils.snowflake.snowpark.snowflake_udfs子模块真正承载全部 UDF 函数实现feast.infra.utils.snowflake.snowpark包本身作为模块容器导出。对应的源码位于 snowflake_udfs.py同目录还包含两份部署模板 snowflake_python_udfs_creation.sql 与 snowflake_python_udfs_deletion.sql。也就是说本文讨论的文档本质上是对一段支撑生产物化链路的 Python 代码的官方 API 说明内容充实、可直接对照源码验证。为什么需要一组 Snowpark Python UDFFeast 的离线存储使用 Snowflake 时物化任务需要把 Snowflake 查询结果中的原始列值压缩编码为 Feast 在线存储如 Redis、DynamoDB 或 Snowflake 自身的 online store可识别的 protobuf 二进制值。这一步如果拉到 Python 端逐行转换会产生巨大的网络与计算开销Feast 的做法是把序列化逻辑下沉到 Snowflake 内部执行——以 Python UDF 形式部署在 Snowflake 仓库上让转换发生在查询引擎内。从源码结构看这套 UDF 由SnowflakeComputeEngine统一部署和调用见 snowflake_engine.py而连接管理、zip 打包等辅助能力位于 snowflake_utils.py。UDF 全清单与类型映射snowflake_udfs.py中共有 16 支函数全部以vectorized(inputpandas.DataFrame)修饰Snowflake 的批处理 Python UDF 模式输入输出为 Pandas DataFrame返回值统一为BINARY即序列化后的 protobuf 字节。下表按标量类型与列表类型分类汇总其中ValueType枚举定义于 value_type.py函数名Snowflake 输入类型对应 Feast ValueType说明feast_snowflake_binary_to_bytes_protoBINARYBYTES字节值转 protobuffeast_snowflake_varchar_to_string_protoVARCHARSTRING字符串值转 protobuffeast_snowflake_number_to_int32_protoNUMBERINT3232 位整数feast_snowflake_number_to_int64_protoNUMBERINT6464 位整数feast_snowflake_float_to_double_protoDOUBLEFLOAT/DOUBLE浮点统一按 double 存储Snowflake 浮点类型即 DOUBLEfeast_snowflake_boolean_to_bool_boolean_protoBOOLEANBOOL布尔值feast_snowflake_timestamp_to_unix_timestamp_protoNUMBERUNIX_TIMESTAMP纳秒时间戳转 Unix 时间戳 protobuffeast_snowflake_array_bytes_to_list_bytes_protoARRAYBYTES_LIST字节数组feast_snowflake_array_varchar_to_list_string_protoARRAYSTRING_LIST及 UUID/DECIMAL 系列字符串数组feast_snowflake_array_number_to_list_int32_protoARRAYINT32_LIST整数数组feast_snowflake_array_number_to_list_int64_protoARRAYINT64_LIST长整数数组feast_snowflake_array_float_to_list_double_protoARRAYDOUBLE_LIST浮点数组feast_snowflake_array_boolean_to_list_bool_protoARRAYBOOL_LIST布尔数组feast_snowflake_array_timestamp_to_list_unix_timestamp_protoARRAYUNIX_TIMESTAMP_LIST时间戳数组feast_serialize_entity_keysARRAY×3——将 1n 个实体键序列化为单个二进制键lookup 用feast_entity_key_proto_to_stringARRAY×3——将实体键编码为EntityKeyProto的序列化字符串其中_convert_value_name_to_snowflake_udftype_map.py维护了 ValueType 名称到 UDF 名的完整映射包括UUID、TIME_UUID、DECIMAL复用varchar_to_string_proto以及UUID_LIST、DECIMAL_LIST、DECIMAL_SET等复用array_varchar_to_list_string_proto的兼容分支读者可在该函数中核对全部 24 个 ValueType 分支的对应关系。序列化实现细节每支标量/列表 UDF 的实现模式高度一致以feast_snowflake_binary_to_bytes_proto为例vectorized(inputpandas.DataFrame) def feast_snowflake_binary_to_bytes_proto(df): sys._xoptions[snowflake_partner_attribution].append(feast) df list( map( ValueProto.SerializeToString, python_values_to_proto_values(df[0].to_numpy(), ValueType.BYTES), ) ) return df要点有三partner attribution每支 UDF 首先执行sys._xoptions[snowflake_partner_attribution].append(feast)用于向 Snowflake 标记调用来源为 Feast 生态类型转换借助 python_values_to_proto_values 把 numpy 数组按指定ValueType转成ValueProto字节输出通过ValueProto.SerializeToString输出二进制供下游直接落库。部分 UDF 有特殊处理列表类函数中array_bytes_to_list_bytes_proto与array_float_to_list_double_proto会先用np.asarray(...).astype(bytes)/.astype(float)做显式类型收敛避免 Snowflake ARRAY 元素类型与 Python 端不一致array_timestamp_to_list_unix_timestamp_proto使用np.datetime64统一时间表示timestamp_to_unix_timestamp_proto用pandas.to_datetime(df[0], unitns)把纳秒级数值还原为时间戳再序列化。实体键序列化在线查找的关键路径模块中两支实体键 UDF 签名均为(names ARRAY, data ARRAY, types ARRAY)分别对应实体 join key 名、值、类型三个数组由调用方用ARRAY_CONSTRUCT在 SQL 中构造。feast_serialize_entity_keys的内部流程用create_entity_dict(names, types)通过_convert_value_type_str_to_value_type把类型字符串还原为ValueType建立列名 → 类型字典将data数组重组为以列为单位的 DataFrame逐列调用python_values_to_proto_values生成ValueProto其中BYTES类型因为 Snowflake 会把二进制转成十六进制字符串需要先用unhexlify还原对每一行构造EntityKeyProto(join_keys..., entity_values[...])并调用serialize_entity_key(..., entity_key_serialization_version3)见 key_encoding_utils.py输出最终的二进制实体键。feast_entity_key_proto_to_string与前者几乎相同区别在于它只做EntityKeyProto.SerializeToString()不经过版本化的serialize_entity_key。从 snowflake_engine.py 的调用逻辑可以推断当在线存储是 Snowflake 自身snowflake.online时使用serialize_entity_keys其他在线存储则使用entity_key_proto_to_string即两支 UDF 服务于不同的在线存储编码约定。部署机制SQL 模板与占位符UDF 的部署不是手工执行而是由SnowflakeComputeEngine.update()在feast apply阶段自动完成流程如下snowflake_engine.py计算 stage 路径database.schema.feast_projectCREATE STAGE IF NOT EXISTS创建内部 stage调用package_snowpark_zip(project)把 UDF 运行所需代码打成feast.zip并PUT上传到 stage读取 snowflake_python_udfs_creation.sql逐条替换占位符后执行STAGE_HOLDER→ 实际 stage 路径对应IMPORTS (stage/feast.zip)PROJECT_NAME→ 当前 Feast 项目名UDF 名带feast_project_前缀实现多项目隔离RUNTIME_VERSION_HOLDER→ 配置的python_udf_runtime_version。需要特别说明的是snowflake_udfs.py文件头部那些CREATE OR REPLACE FUNCTION ...块只是便于阅读的参考文档并非 Feast 实际执行的语句真正执行的 SQL 是上述模板填充后的结果。模板中每支函数固定声明LANGUAGE PYTHON、PACKAGES (protobuf, pandas)并显式指定HANDLER为feast.infra.utils.snowflake.snowpark.snowflake_udfs.函数名。feast.zip 的打包内容package_snowpark_zipsnowflake_utils.py会把以下文件从安装好的 Feast 包中复制并打成 zip作为feast包上传到 Snowflake stageinfra/utils/snowflake/snowpark/snowflake_udfs.pyUDF 实现本体infra/key_encoding_utils.py实体键序列化type_map.py、value_type.py类型转换与枚举protos/feast/types/Value_pb2.py、protos/feast/types/EntityKey_pb2.pyprotobuf 消息这正是每支 UDF 的IMPORTS (stage/feast.zip)所引用的运行时依赖也是 UDF 能完整复现 Feast 侧序列化语义的前提。Python UDF 运行时版本可配置项与默认值SnowflakeComputeEngineConfig.python_udf_runtime_versionsnowflake_engine.py控制部署时写入RUNTIME_VERSION的 Python 版本默认值为3.10与 Feast 自身pyproject.toml中声明的requires-python最低版本一致校验规则要求符合\d\.\d(\.\d)?格式例如3.10或3.11.2Snowflake 会周期性下线旧版 Python UDF 运行时例如 3.9 曾因此被替换为避免被某个硬编码版本卡住该字段开放给用户在feature_store.yaml的batch_engine段覆盖例如batch_engine: type: snowflake.engine account: my_account user: my_user password: my_password database: my_db schema: PUBLIC python_udf_runtime_version: 3.11此外SnowflakeComputeEngine.update()每次部署都使用CREATE OR REPLACE FUNCTION无条件重建 UDF而非存在即跳过——这一设计源码注释中明确提及确保修改python_udf_runtime_version后重新feast apply即可让新版本真正生效避免旧运行时被悄悄沿用。物化查询中的实际调用在物化阶段generate_snowflake_materialization_querysnowflake_engine.py为每个 feature 按feature.dtype.to_value_type().name从_convert_value_name_to_snowflake_udf查出对应 UDF 名构造形如UDF(feature_name) AS feature_name的 SQL 片段其中UNIX_TIMESTAMP类型会先做DATE_PART(EPOCH_NANOSECOND, col::TIMESTAMP_LTZ)转为纳秒数再喂给 UDFDOUBLE类型显式::DOUBLE转换实体键列则用ARRAY_CONSTRUCT组装names/data/types三个数组交给上述实体键 UDF 输出entity_key列。随后materialize_to_snowflake_online_store或materialize_to_external_online_store消费这批 protobuf 二进制列前者在 Snowflake 内直接MERGE INTO在线表后者通过fetch_pandas_batches()拉回 Python 端后调用在线存储的online_write_batch写入。清理与删除对应的 snowflake_python_udfs_deletion.sql 按签名精确删除全部 16 支函数DROP FUNCTION IF EXISTS ...由SnowflakeComputeEngine.teardown_infra()在销毁项目基础设施时执行同时还会DROP STAGE IF EXISTS清理上传的 stage 与feast.zip。相关阅读入口若需深入这套序列化链路的更多上下文可继续阅读仓库中的以下文件snowflake_udfs.py全部 UDF 的实现与头部参考 SQLsnowflake_python_udfs_creation.sql 与 snowflake_python_udfs_deletion.sql部署/删除模板snowflake_engine.pyUDF 的部署、物化 SQL 生成与在线写入snowflake_utils.pySnowflake 连接管理与feast.zip打包type_map.pyValueType 到 UDF 名的完整映射表配套的 feast.infra.utils.snowflake.rst 与 feast.infra.compute_engines.snowflake.rst 文档页可对照查看连接工具与计算引擎的 API 说明。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表