Feast Feature Transformation 全解析:三大引擎、UDF 与聚合/过滤/关联 API 实战指南
发布时间:2026/9/16 21:20:28来源:尧图网络
Feast Feature Transformation 全解析三大引擎、UDF 与聚合/过滤/关联 API 实战指南【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast导读本文是 FeastThe Open Source Feature Store for AI/ML官方架构文档 feature-transformation.md 的深度展开系统讲解特征变换Feature Transformation的核心概念、三种变换引擎的选型与权衡、以及feature_transformation/udf、Aggregation、ttl过滤、多 FeatureView 关联Join等核心 API 的完整用法。读完本文你将掌握如何在 Feast 中定义可序列化、可跨引擎执行的变换逻辑理解变换引擎与数据写入通信模式的耦合关系并能结合源码层面sdk/python/feast/transformation写出真正可运行、可下发的特征变换代码。什么是 Feature Transformation在 Feast 中特征变换Feature Transformation是接收一组输入数据、返回一组输出数据的函数。它的输入可以是原始数据raw data也可以是已派生的数据derived data例如对原始日志中的文本字段做清洗去多余空格、标准化大小写对交易金额字段做聚合1 小时内的 sum / avg对多个 FeatureView 做关联拼接形成复合特征视图。变换可以是批量的作用于整张 DataFrame也可以是单条的作用于单实体请求即 On Demand 场景。Feast 将变换逻辑与特征定义解耦使同一份特征配方能够在不同的执行后端上复用。Feature Transformation Engines三大变换引擎官方文档明确指出特征变换可以由三类变换引擎执行Feast Feature ServerPython 实现的在线特征服务端负责在请求时On Demand或写入时transform_on_write执行变换见 python-feature-server.mdOffline Store如 Snowflake、BigQuery、DuckDB、Spark 等在批量训练数据集生成get_historical_features或物化materialize阶段由离线存储自身执行 SQL 化的变换Compute EngineFeast 可插拔的特征管线执行抽象负责在 Spark、PyFlink、Ray、本地 ArrowPandas 等后端上执行变换、聚合、关联与物化见 compute-engine/README.md。变换引擎与写入通信模式的耦合Feast 采用 Push 模型 向在线存储写入特征因此三大变换引擎与写入通信模式write-patterns.md是耦合的。这意味着同样的特征变换代码在不同引擎下可能以完全不同的方式被调度执行。官方文档特别强调理解何时用哪种引擎 哪种通信模式是实现成功的关键通常建议根据数据生产方、特征/模型使用方式以及整体产品形态来选择最合适的引擎与网络调用方式。从 overview.md 可以看到 Feast 官方的定位当前对 On Demand 与 Streaming 数据源提供特征变换支持Batch 变换由离线存储或 Compute Engine 承接在线场景下预计算precompute是最优路径——把特征服务降级为一次轻量的数据库查询是延迟最低的理想模式。写入通信模式的权衡选型参考来自 write-patterns.md 的决策表可帮助针对不同应用选择数据写方式 × 特征计算方式的组合数据写方式特征计算典型场景推荐做法异步On Demand可容忍数据滞后的大数据量应用异步写 按需计算平衡负载与资源使用异步预计算高吞吐、非关键数据处理异步批作业 预计算变换兼顾效率与扩展性同步On Demand高利害决策如贷款审批同步写 按需计算确保数据新鲜度与正确性同步预计算需要快速反馈的用户型应用同步写 预计算特征降低延迟同步混合预计算 On Demand有延迟约束的高利害决策尽可能预计算配合少量按需计算其中transform_on_write参数源码见 feature_store.py控制写入在线存储时是否应用变换对于已预处理过的数据可以跳过变换同时保留 API 调用时的变换能力实现混合模式。API 详解feature_transformation / udffeature_transformation与udf是 Feast 定义特征变换的核心 API允许用户在物化materialization或检索retrieval期间对数据应用自定义逻辑。官方文档给出了三种等价的书写方式下面逐一展开并补充源码级说明。方式一SparkTransformation模式为 TransformationMode.SPARKdef remove_extra_spaces(df: DataFrame) - DataFrame: df[name] df[name].str.replace(\s, ) return df spark_transformation SparkTransformation( modeTransformationMode.SPARK, udfremove_extra_spaces, udf_stringremove extra spaces, ) feature_view FeatureView( feature_transformationspark_transformation, ... )方式二Transformation模式为 TransformationMode.SPARK_SQLspark_transformation Transformation( modeTransformationMode.SPARK_SQL, udfremove_extra_spaces_sql, udf_stringremove extra spaces sql, ) feature_view FeatureView( feature_transformationspark_transformation, ... )方式三transformation 装饰器transformation(modeTransformationMode.SPARK) def remove_extra_spaces_udf(df: pd.DataFrame) - pd.DataFrame: return df.assign(namedf[name].str.replace(\s, )) feature_view FeatureView( feature_transformationremove_extra_spaces_udf, ... )源码级说明Transformation 与 TransformationModeTransformationModemode.py是一个枚举定义了全部可用的变换模式PYTHON、PANDAS、SPARK_SQL、SPARK、FLINK、RAY、SQL、SUBSTRAIT。Transformation 基类base.py是抽象基类其__new__会根据传入的mode自动路由到对应的子类通过 factory.py 中的映射表python→PythonTransformation、pandas→PandasTransformation、spark/spark_sql→SparkTransformation、flink→FlinkTransformation、ray→RayTransformation、sql→SQLTransformation、substrait→SubstraitTransformation。若传入非法 mode会抛出ValueError。三个核心参数mode必填变换模式取值见上述枚举udf必填用户自定义变换函数udf_string必填UDF 的源码字符串。源码注释明确指出dill 获取源码并不适用于所有场景最好将源码作为字符串显式传入——这正是transformation装饰器内部使用dill.source.getsource(user_function)自动提取源码base.py的原因。可选参数name默认取udf.__name__、tags、description、owner。序列化机制Transformation.to_proto()base.py使用dill.dumps(self.udf, recurseTrue)序列化函数体并将udf_string作为body_text一并写入UserDefinedFunctionV2proto。这意味着变换逻辑会随特征注册进入 Registry并支持后续的变换生命周期管理。装饰器中的mainify还会把函数模块名改写为__main__保证 dill 序列化不受原始文件命名影响。执行入口子类各自实现transform/transform_arrow/transform_singleton。例如PythonTransformation.transformpython_transformation.py将输入 dict 传入 UDF 并把输出与输入合并返回SparkTransformation.transformspark_transformation.py在SPARK_SQL模式下会为每个输入 DataFrame 创建临时视图feast_transformation_temp_view_{name}再执行self.spark_session.sql(...)在SPARK模式下则直接调用 UDF其构造函数还会通过get_or_create_new_spark_session(spark_config)复用或新建 SparkSessionPandasTransformation会校验 UDF 返回类型必须是pd.DataFrame否则抛TypeErrorRayTransformationray_transformation.py专为分布式密集计算设计如 RAG 场景的 Embedding 生成、图像/视频处理UDF 接收并返回ray.data.Dataset推荐用有状态类配合map_batches避免每个 batch 重复加载模型FlinkTransformation的 UDF 接收并返回 PyFlink Table。特征推断PythonTransformation/PandasTransformation等均实现了infer_features用随机样例数据运行 UDF 后通过python_type_to_feast_value_type自动推断输出字段类型推断失败时如 UDF 返回空会给出明确的TypeError提示此时可以在 BatchFeatureView 中显式声明schema。注册链路在 feature_view.py 中FeatureView 的to_proto会调用transformation_to_proto将feature_transformation序列化进 FeatureView spec反向的from_proto则通过udf_rehydrate.resolve_udf依据body_text/body重建可调用对象。Aggregation内置的窗口聚合 APIAggregation是 Feast 内置的批式/流式聚合API用于描述在时间窗口上如何聚合数据如某特征在指定周期内的平均值或总和。from feast import Aggregation feature_view FeatureView( aggregations[ Aggregation( columnamount, functionsum ), Aggregation( columnamount, functionavg, time_window1h ), ], ... )源码级说明Aggregation 参数Aggregation 的完整参数如下参数类型说明columnstr参与聚合的特征列名functionstr内置聚合函数sum、max、min、count、mean、count_distinctcount_distinct会被别名映射为nunique见 aggregation/init.pytime_windowtimedelta聚合的时间窗口可选slide_intervaltimedelta滑动窗口步长可选默认等于time_windownamestr输出特征名覆盖默认推导为{function}_{column}带窗口时追加_{seconds}s见resolved_name此外类中的 docstring 明确提示Feast 自动处理的聚合目前仍在演进该类用于注册用户自定义聚合。to_proto/from_proto会将time_window、slide_interval与 protobuf 的Duration相互转换确保聚合定义可持久化到 Registry。在 Compute Engine 的 DAG 管线docs/reference/compute-engine/README.md中聚合对应AggregationNode管线顺序为SourceReadNode → TransformationNode若定义了 feature_transformation→ FilterNode始终存在应用 TTL 或用户过滤→ AggregationNode若定义了 aggregations→ DeduplicationNode → ValidationNode → Output。FilterTTL 时间过滤ttl决定特征可被物化或检索的时间长度只有实体行的时间戳大于当前时间减去 ttl的行才会被用于特征计算从而确保特征计算只使用近期数据。feature_view FeatureView( ttl1d, # Features will be available for 1 day ... )源码级说明ttl 的底层行为在 feature_view.py 中ttl是Optional[timedelta]类型默认值为timedelta(days0)注释明确说明ttl 为 0 表示特征永久有效而较大的 ttl或 0会带来存储与查询成本的上升需谨慎选择。get_ttl_duration()feature_view.py负责将timedelta转换为 protobufDuration写入 FeatureView spec反序列化时纳秒为 0 的 ttl 会被还原为timedelta(days0)feature_view.py。在 Compute Engine 的 DAG 中TTL 过滤由始终存在的FilterNode实现在 FeatureView 的校验逻辑中ttl、online、offline等属于跳过配置比较的字段feature_view.py即它们不影响特征定义的 diff 判定但会参与物化与检索时的行级过滤。Join多 FeatureView 关联Feast 可以将多个 FeatureView 关联起来创建复合特征视图composite feature view从而把来自不同数据源或不同视图的特征合并到同一个视图中feature_view FeatureView( namecomposite_feature_view, entities[entity_id], source[ FeatureView( namefeature_view_1, features[feature_1, feature_2], ... ), FeatureView( namefeature_view_2, features[feature_3, feature_4], ... ) ], ... )官方文档明确指出Join 的底层实现默认是 inner joinjoin key 就是实体 ID。从源码结构看这与 Compute Engine 的FeatureBuilder行为一致当 FeatureView 的source是多个 FeatureView 列表时FeatureResolver会按依赖关系对 DAG 节点做拓扑排序并通过JoinNode完成默认的关联操作见 docs/reference/compute-engine/README.md。关于 join key 的进一步细节在 feature_view.py 中FeatureView 的join_keys由关联实体的entity.join_key汇总而来若出现重复 join key 会直接报错实体 join key 的映射关系with_join_key_map可以结合 ADR-0004-entity-join-key-mapping.md 进一步了解该 ADR 讨论了实体 join key 映射的语义与查询期行为。实战建议与引擎选型小结综合官方文档与源码选择特征变换方案时建议遵循以下原则在线低延迟优先预计算把变换提前在批作业/流作业中完成让在线检索退化为一次在线存储查询这是 Feast 官方推荐的理想模式参见 overview.md 与 write-patterns.md。On Demand 变换用于个性化/时间敏感特征如距上次购买的时间这类混合模式特征适合在请求时即时计算。按数据规模选引擎小数据量本地开发用PANDAS/PYTHON本地 Compute Engine 基于 Arrow Pandas/Polars大规模批处理用SPARK/SPARK_SQL/FLINK/RAY全 SQL 场景可直接由离线存储Snowflake、BigQuery 等执行。务必显式提供udf_string由于 dill 提取源码存在局限显式传入源码字符串可保证变换可被序列化、可被from_proto重建这也是 Web UI 中展示变换源码的基础。把引擎选型与写入通信模式一起决策同步/异步写入 × 预计算/On Demand/混合计算会直接影响数据新鲜度、正确性、服务耦合度与应用延迟请结合前文决策表权衡。延伸阅读Write Patterns 与通信模式权衡Push 模型与 Pull 模型对比Compute Engine 与 DAG 执行管线Python Feature Server 与 push 端点Push DataSource 写入说明Feast 架构总览变换实现源码 与 聚合实现源码【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网