ARTICLE DETAIL

资讯详情

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

Feast Compute Engine 深度解析:统一计算引擎架构、DAG 执行模型与批量/流式引擎配置实战

Feast Compute Engine 深度解析:统一计算引擎架构、DAG 执行模型与批量/流式引擎配置实战 Feast Compute Engine 深度解析统一计算引擎架构、DAG 执行模型与批量/流式引擎配置实战【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feastFeast 的 Compute Engine计算引擎是负责特征物化materialization与历史特征检索historical retrieval的核心组件它统一了本地、Spark、Snowflake、AWS Lambda、Flink、Ray 等多种执行后端。本文将以 docs/getting-started/components/compute-engine.md 为骨架结合sdk/python/feast/infra/compute_engines/下的真实源码实现系统讲解 Compute Engine 的抽象接口、feature_store.yaml与 FeatureView 级引擎配置、基于 DAG 的执行计划构建流程以及如何扩展自定义引擎。Compute Engine 是什么在 Feast 中Compute Engine 是一个负责**物化materialization和历史检索historical retrieval**任务的组件。它负责执行 FeatureView 定义中声明的逻辑包括聚合Aggregation变换Transformation自定义用户函数UDFJoin、过滤等其他特征生成操作从源码角度所有计算引擎都实现同一个抽象基类ComputeEngine定义于 sdk/python/feast/infra/compute_engines/base.py。该类的 docstring 明确说明了它的定位The interface that Feast uses to control the compute system that handles materialization and get_historical_features. Each engine must implement: materialize() to generate and persist features; get_historical_features() to perform historical retrieval of features.每个引擎必须实现两类核心能力抽象方法职责materialize()生成并持久化特征写在线/离线存储get_historical_features()执行历史特征检索同时基类还定义了update()在feast apply时准备物化所需的云资源与teardown_infra()在feast teardown时清理资源两个生命周期方法以及supports_batch、applies_materialization等能力属性。值得注意的是base.py中还内置了_get_feature_view_engine_config()方法负责合并仓库级默认引擎配置与 FeatureView 级覆盖配置这为下文要讲的batch_engine/stream_engine优先级机制提供了底层支撑。一个关键的设计抽象是MaterializationTask物化任务它把具体使用哪种技术/框架来物化数据这件事与任务本身解耦。用户既可以使用纯本地的序列化执行方式默认的LocalComputeEngine也可以把物化委托给独立的组件例如 AWS Lambda 实现的LambdaComputeEngine。支持的 Compute Engine 一览原文档给出了 Feast 当前支持的内置计算引擎清单结合仓库源码 sdk/python/feast/repo_config.py 中BATCH_ENGINE_CLASS_FOR_TYPE的注册表可以得到更完整的对应关系Compute Engine说明支持feature_store.yaml中的 typeLocalComputeEngine基于 Arrow Pandas/Polars/Dask 等面向轻量级变换✅localSparkComputeEngine基于 Apache Spark面向大规模分布式特征生成✅spark.engineSnowflakeComputeEngine基于 Snowflake使用 Snowflake SQL 做可扩展特征生成✅snowflake.engineLambdaComputeEngine基于 AWS Lambda面向无服务器特征生成✅lambdaFlinkComputeEngine基于 Apache Flink通过 PyFlink Table API 做分布式特征生成✅flink.engineRayComputeEngine基于 Ray面向分布式特征生成与机器学习负载✅ray.engine此外仓库中还注册了KubernetesComputeEnginetype 为k8s实现位于 sdk/python/feast/infra/compute_engines/kubernetes/k8s_engine.py与SparkApplicationComputeEnginetype 为spark_application实现位于 sdk/python/feast/infra/compute_engines/spark_application/compute.py可用于需要独立 Driver/Application 承载物化任务的场景。如果内置引擎无法满足需求可以创建自定义物化引擎具体参见自定义计算引擎指南。引擎的全局配置入口feature_store.yaml引擎在feature_store.yaml中配置。该文件的完整字段说明见 feature_store.yaml 参考文档。从源码看仓库配置中的batch_engine字段被解析到RepoConfig.batch_engine_config见 sdk/python/feast/repo_config.py并且当feature_store.yaml中完全未声明batch_engine时Feast 会默认使用本地进程内物化引擎batch_engine_config local见 sdk/python/feast/repo_config.py这也是大多数快速上手场景的默认行为。一个典型的feature_store.yaml引擎配置如下project: my_project registry: registry.db batch_engine: type: spark.engine config: spark_master: local[*] spark_app_name: Feast Batch Engine spark_conf: spark.sql.shuffle.partitions: 100 spark.executor.memory: 4g online_store: type: sqlite path: online_store.db offline_store: type: file其中type必须是上表列出的引擎类型标识或自定义引擎的模块路径config部分则透传给对应引擎的配置类。例如spark.engine会通过get_batch_engine_config_from_type()sdk/python/feast/repo_config.py解析为SparkComputeEngineConfig。Batch Engine批量物化与历史检索的默认引擎batch_engine配置在feature_store.yaml中声明后将作为所有物化任务与历史检索任务的默认配置。它可以在两个层级设置层级一feature_store.yaml 全局默认batch_engine: type: spark.engine config: spark_master: local[*] spark_app_name: Feast Batch Engine spark_conf: spark.sql.shuffle.partitions: 100 spark.executor.memory: 4g层级二BatchFeatureView 级覆盖在 Python 定义 FeatureView 时通过BatchFeatureView.batch_engine参数为单个视图覆盖全局默认配置from feast import BatchFeatureView fv BatchFeatureView( batch_engine{ spark_conf: { spark.sql.shuffle.partitions: 200, spark.executor.memory: 8g }, } )当物化该 FeatureView 时Feast 将使用 FeatureView 中指定的batch_engine配置——即spark.sql.shuffle.partitions 200、spark.executor.memory 8g覆盖仓库级默认的 100 与 4g。两层配置的合并优先级从实现上看这一优先级逻辑由 sdk/python/feast/infra/compute_engines/base.py 的_get_feature_view_engine_config()保证以仓库配置self.repo_config.batch_engine_config作为基线baseline若 FeatureView 提供了运行时覆盖BatchFeatureView 的batch_engine/ StreamFeatureView 的stream_engine则以浅合并方式覆盖基线{**default_conf, **runtime_conf}若运行时配置不是 dict会抛出TypeError提示引擎配置必须是字典类型。这意味着你可以为绝大多数视图保留轻量默认值只对计算密集型的少数视图单独调大 Spark 资源配置。Stream Engine流式物化的引擎配置与 Batch Engine 对应stream_engine配置在feature_store.yaml中作为所有流式物化与流式历史检索任务的默认配置。它可以在StreamFeatureView中通过stream_engine参数按视图覆盖。全局默认feature_store.yamlstream_engine: type: spark.engine config: spark_master: local[*] spark_app_name: Feast Stream Engine spark_conf: spark.sql.shuffle.partitions: 100 spark.executor.memory: 4gStreamFeatureView 级覆盖from feast import StreamFeatureView fv StreamFeatureView( stream_engine{ spark_conf: { spark.sql.shuffle.partitions: 200, spark.executor.memory: 8g }, } )物化该 FeatureView 时Feast 会使用视图内声明的stream_engine配置shuffle partitions 为 200、executor memory 为 8g其合并逻辑与 Batch Engine 完全相同均由基类的_get_feature_view_engine_config()统一处理。计算引擎的 API 与执行计划构建Compute Engine 以FeatureBuilder的形式将执行计划构建为 DAG 格式。它从 FeatureView 定义中推导出特征生成所需的全部操作Transformation通过 Transformation APIAggregation通过 Aggregation APIJoin与实体数据集 join、自定义 JOIN、或与另一个 FeatureView joinFilterpoint-in-time 过滤、TTL 过滤、按自定义表达式过滤其他必要的衍生操作在源码层面这一过程由两条关键链路实现sdk/python/feast/infra/compute_engines/feature_builder.py 中的FeatureBuilder抽象基类定义了build_source_node/build_transformation_node/build_join_node/build_filter_node/build_aggregation_node/build_dedup_node/build_validation_node/build_output_nodes等构建方法并提供了_build()编排方法按读取源 → 变换/join → 过滤 → 聚合或去重 → 校验 → 输出的固定顺序组织节点每种引擎提供自己的FeatureBuilder实现例如LocalFeatureBuildersdk/python/feast/infra/compute_engines/local/feature_builder.py、Spark、Flink、Ray 各自的 builder分别位于 spark/、flink/、ray/ 目录。FeatureBuilder.build()的完整流程feature_builder.py用FeatureResolver把 FeatureView 之间的依赖关系解析为逻辑 DAG并做拓扑排序按拓扑序为每个 FeatureView 构建对应的物理执行 DAGNode为根 FeatureView 构建输出节点对最终 DAG 再做一次拓扑排序返回ExecutionPlansdk/python/feast/infra/compute_engines/dag/plan.py由引擎按序执行。ExecutionPlan.execute()会逐个执行 DAG 节点并把每个节点的输出DAGValue缓存在ExecutionContext.node_outputs中使下游节点可以直接复用上游结果而无需重复计算。Compute Engine 的核心组件Compute Engine 负责执行 FeatureView 中定义的物化与检索任务它构建一张描述特征生成所需操作的有向无环图DAG。核心组件有三个Feature Builder负责从 FeatureView 中解析特征并执行 DAG 中定义的操作处理变换、聚合、join 与过滤的执行。每个引擎都有专属的 Builder 实现用于把统一的 FeatureView 定义翻译成该引擎可执行的计算图。Feature Resolver计算引擎的核心组件负责为特征生成构造执行计划。它接收 FeatureView 定义构建一张需要执行的操作 DAG。源码实现在 sdk/python/feast/infra/compute_engines/feature_resolver.pyFeatureResolver.resolve()递归遍历 FeatureView 的source_views依赖链用_node_cache缓存已解析节点、用_resolution_path检测循环依赖一旦发现环会抛出Cycle detected in FeatureView DAG异常。它还提供topological_sort()与debug_dag()便于调试依赖图。DAGDAG 是需要执行的特征生成操作的有向无环图包含变换、聚合、join、过滤等操作节点。DAG 由 Feature Resolver 构建由 Feature Builder 执行。节点抽象定义于 sdk/python/feast/infra/compute_engines/dag/node.py每个DAGNode持有name、inputs、outputs通过add_input()建立边关系通过execute(context)执行自身逻辑并通过get_input_values()/get_single_input_value()从执行上下文读取上游输出。DAG 节点管线从数据源到输出的完整链路DAG 节点按如下流水线组织该图为原文档的核心架构图完整保留--------------------- | SourceReadNode | - Read data from offline store (e.g. Snowflake, BigQuery, etc. or custom source) --------------------- | v -------------------------------------- | TransformationNode / JoinNode (*) | - Merge data sources, custom transformations by user, or default join -------------------------------------- | v --------------------- | FilterNode | - used for point-in-time filtering --------------------- | v --------------------- | AggregationNode (*) | - only if aggregations are defined --------------------- | v --------------------- | DeduplicationNode | - used if no aggregation and for history --------------------- retrieval | v --------------------- | ValidationNode (*) | - optional validation checks --------------------- | v ---------- | Output | ---------- / \ v v ---------------- ---------------- | OnlineStoreWrite| OfflineStoreWrite| ---------------- ----------------各节点的执行语义与源码依据如下SourceReadNode从离线存储Snowflake、BigQuery 等或自定义 source读取数据并按start_time/end_time限定时间范围同时通过get_column_info()确定需要读取的 join key、特征列、时间戳列等见 local/feature_builder.pyTransformationNode / JoinNode合并数据源、执行用户自定义变换或在无变换时执行默认 join。用户定义了feature_transformation的视图走变换分支否则走默认 join 分支见 feature_builder.pyFilterNode执行 point-in-time 过滤也处理视图的filter表达式与ttl见 local/feature_builder.pyAggregationNode仅当视图定义了aggregations时插入通过aggregation_specs_to_agg_ops()把聚合规格转换为聚合算子并以视图实体作为 group by 键local/feature_builder.py。注意本地引擎不支持时间窗口聚合会抛出明确的错误提示DeduplicationNode当没有聚合且任务属于历史检索HistoricalRetrievalTask或only_latest时使用保证只保留每个实体键的最新记录feature_builder.pyValidationNode可选节点仅在视图开启enable_validation时插入可依据视图特征声明自动推导 PyArrow 类型做列类型校验并对 JSON 类型列做内容校验local/feature_builder.pyOutput最终输出分流到 OnlineStoreWrite 与 OfflineStoreWrite分别写在线与离线存储。从FeatureBuilder._build()的编排逻辑feature_builder.py可以看出这套节点序列与上图完全一致先读源/变换/join再过滤然后二选一执行聚合或去重最后可选校验并输出。本地引擎的轻量实现LocalComputeEngineLocalComputeEnginesdk/python/feast/infra/compute_engines/local/compute.py是默认引擎其配置类LocalComputeEngineConfig支持两个字段class LocalComputeEngineConfig(FeastConfigBaseModel): type: Literal[local] local backend: Optional[str] None # 例如 pandas、polarstype固定为localbackend指定 DataFrame 操作后端如pandas、polars。若未显式指定backendLocalComputeEngine._get_backend()会根据运行时entity_df的类型自动推断pd.DataFrame、pyarrow.Table或空值默认使用 Pandas 后端polars.DataFrame自动切换为 Polars 后端推断失败则抛出ValueError推断逻辑见 sdk/python/feast/infra/compute_engines/backends/factory.py 的BackendFactory。本地物化的执行路径是_materialize_one()构造LocalFeatureBuilder→builder.build()生成ExecutionPlan→plan.execute(context)串行执行 DAG最终返回携带SUCCEEDED/ERROR状态的LocalMaterializationJoblocal/compute.py。历史检索则返回懒执行的LocalRetrievalJob持有 plan 与上下文供后续拉取结果。自定义 Compute Engine 扩展指南当内置引擎不足时可以开发自己的计算引擎。核心要求是实现ComputeEngine接口契约sdk/python/feast/infra/compute_engines/base.py完整教程见创建自定义计算引擎这里提炼关键步骤Step 1定义引擎类。继承ComputeEngine至少实现update()与materialize()如需历史检索还要实现get_historical_features()from feast.infra.compute_engines.base import ComputeEngine class MyCustomEngine(ComputeEngine): def update(self, project, views_to_delete, views_to_keep, entities_to_delete, entities_to_keep): print(Creating new infrastructure is easy here!) pass def materialize(self, registry, tasks): print(Launching custom batch jobs is pretty easy...) return [self._materialize_one(registry, task, **{}) for task in tasks] def get_historical_features(self, task): raise NotImplementedErrorStep 2在feature_store.yaml中指向引擎类。batch_engine字段可以填模块路径模块名.类名需以Engine结尾例如project: repo registry: registry.db batch_engine: feast_custom_engine.MyCustomEngine online_store: type: sqlite path: online_store.db offline_store: type: fileStep 3运行feast apply生效。若引擎模块不在默认搜索路径需要通过PYTHONPATH指定PYTHONPATH$PYTHONPATH:/home/my_user/my_custom_engine feast applyfeast apply时会在Deploying infrastructure阶段触发你自定义的update()逻辑若引擎需要独立进程/集群承载也可以选择重写materialize()以提交自定义批量作业Spark、Beam、AWS Lambda 等并在feast teardown时通过teardown_infra()回收资源。总结Feast 的 Compute Engine 通过统一接口 多后端实现的架构把特征物化与历史检索从具体技术栈中解放出来LocalComputeEngine满足轻量场景SparkComputeEngine/SnowflakeComputeEngine/FlinkComputeEngine/RayComputeEngine/LambdaComputeEngine等满足分布式与无服务器场景batch_engine与stream_engine支持从feature_store.yaml全局配置到单个 FeatureView 按需覆盖的两级配置模型FeatureResolverFeatureBuilder DAG 的执行计划机制保证了不同引擎共享同一套特征生成语义。深入阅读 compute_engines 目录下的dag/、backends/与各引擎子目录可以进一步掌握每个节点的实现细节。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表