十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Feast Flink 计算引擎(FlinkComputeEngine)实战指南:用 PyFlink 分布式执行特征物化与历史检索

Feast Flink 计算引擎(FlinkComputeEngine)实战指南:用 PyFlink 分布式执行特征物化与历史检索 Feast Flink 计算引擎FlinkComputeEngine实战指南用 PyFlink 分布式执行特征物化与历史检索【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feastApache Flink 计算引擎是 Feast 中基于 PyFlink Table API 的分布式计算后端它实现 Feast 统一的ComputeEngine接口可承担批量物化materialize/materialize-incremental与历史特征检索get_historical_features两类核心任务。本文以 docs/reference/compute-engine/flink.md 为骨架结合仓库源码与测试用例系统讲解其安装配置、feature_store.yaml参数、modeflink特征变换写法、DAG 各节点在 Flink SQL 层面的实现以及当前已知限制帮助你把它接入自己的特征仓库并理解底层工作原理。一、Flink 计算引擎是什么在 Feast 的ComputeEngine抽象体系见 compute-engine 总览中计算引擎负责“执行特征流水线”——包括变换transformation、聚合aggregation、关联join以及物化/历史检索。Flink 引擎是其中面向大规模分布式执行的选项通过PyFlink Table API构建分布式执行计划数据读取走配置好的 Feast offline store离线存储若原生暴露to_flink_table(table_env)检索任务则直接把 Flink Table 交给引擎否则引擎将标准 Arrow 路径的结果转换为 Flink Table源码见 flink/nodes.py 中_retrieval_job_to_flink_tablejoin、filter、aggregate、dedupe、projection 等步骤全部由 Flink Table/SQL 算子完成物化结果写入配置的 online 与/或 offline store。从源码看引擎的入口类是FlinkComputeEngine其_materialize_one()与get_historical_features()分别对应物化与历史检索两条路径flink/compute.py执行逻辑统一由FlinkFeatureBuilder构建 DAG 后交给ExecutionPlan按拓扑顺序执行。整个引擎位于sdk/python/feast/infra/compute_engines/flink/目录文件职责compute.pyFlinkComputeEngine与FlinkComputeEngineConfig引擎入口与配置模型feature_builder.pyFlinkFeatureBuilder把 FeatureView 解析为 Flink 专属 DAG 节点nodes.py各DAGNode的 Flink 实现Source/Join/Filter/Agg/Dedup/Transform/Validation/Outputjob.pyFlinkDAGRetrievalJob与FlinkMaterializationJobutils.pyPyFlink TableEnvironment 创建、pandas↔Flink 转换、临时视图管理二、安装与依赖flinkextraFlink 引擎依赖 PyFlinkFeast 通过flinkextra 提供安装入口。在 Feast 源码检出目录下用uv安装uv sync --extra flink --no-dev关于依赖约束仓库根目录 pyproject.toml 中flinkextra 定义为flink [apache-flink2.2.1,3, pyarrow21.0.0]需要特别说明两点版本约束它们直接来自仓库配置PyArrow 版本冲突flinkextra 要求pyarrow21而 Feast 默认安装保持pyarrow21。Feast 的 uv lock 将 Flink extra 解析为独立依赖分支[tool.uv.conflicts]声明了flink与ge、ci、dev、docsextra 互斥因此正常 Feast 安装不会被降级 Arrow。互斥关系由于上述冲突声明flinkextra 不应与其他列出的 extrage、ci、dev、docs同时启用这也是安装命令使用--no-dev的原因之一。若运行环境缺少 PyFlink引擎会抛出明确的ImportError提示安装flinkextra见 flink/utils.py 中create_flink_table_environment。三、在 feature_store.yaml 中配置引擎在feature_store.yaml中通过batch_engine配置段启用 Flink 引擎。原文档给出的完整示例project: my_project registry: data/registry.db provider: local offline_store: type: file online_store: type: sqlite path: data/online_store.db batch_engine: type: flink.engine execution_mode: batch parallelism: 4 table_config: pipeline.name: Feast Flink Compute Engine pandas_split_num: 4配置解析链路为RepoConfig.batch_engine属性读取batch_engine字段通过get_batch_engine_config_from_type()依据type查找注册表BATCH_ENGINE_CLASS_FOR_TYPE其中flink.engine映射到feast.infra.compute_engines.flink.compute.FlinkComputeEngine见 repo_config.py。引擎初始化时会把table_config、parallelism、execution_mode写入 PyFlinkConfiguration并据此创建TableEnvironment。配置选项详解OptionTypeDefaultDescriptiontypestringflink.engine必须为flink.engine引擎类型选择器。execution_modestringbatchPyFlink 执行模式batch或streaming。parallelismintegernull引擎创建作业的默认 Flink 并行度。table_configmapnull额外的 PyFlink table 配置项键值对。pandas_split_numinteger1把 pandas entity DataFrame 转成 Flink Table 时的 Arrow source 分片数。各选项对应的源码实现flink/compute.py 中FlinkComputeEngineConfigtypeLiteral[flink.engine]固定为flink.engine用于repo_config的引擎注册表查找execution_modeLiteral[batch, streaming]默认batch。在 flink/utils.py 的create_flink_table_environment中streaming走in_streaming_mode()其余走in_batch_mode()parallelism可选整数非空时写入 Flink 配置键parallelism.defaulttable_config字典逐项调用flink_conf.set_string(key, value)注入 PyFlinkConfiguration示例中的pipeline.name即作业名pandas_split_num整数默认1。它作为splits_num传入table_env.from_pandas(df, splits_num...)见pandas_to_flink_table用于控制 pandas entity DataFrame 转换时的并行分片数。四、用 modeflink 编写 Flink 特征变换当BatchFeatureView的变换函数需要接收并返回 PyFlink Table 对象时使用modeflinkfrom feast import BatchFeatureView, Field from feast.types import Float32 def double_rates(table): # 生产环境中可以在这里使用 PyFlink Table API 操作并返回一个 table。 return table driver_stats BatchFeatureView( namedriver_stats, entities[driver], modeflink, udfdouble_rates, schema[Field(nameconv_rate, dtypeFloat32)], sourcedriver_stats_source, onlineTrue, )约束modeflink的变换函数必须返回 PyFlink Table 对象返回 pandas DataFrame 的 UDF 不被 Flink 计算引擎接受。源码依据batch_feature_view.py 的get_feature_transformation()把modeflink归入支持的变换模式列表flink/nodes.py 中FlinkTransformationNode.execute()直接调用self.transformation_fn(*input_tables)并把返回值当作 Flink Table若返回对象无法取得 schema 则抛出TypeErrorfeature_builder.py 也将flink列为需要变换节点的模式之一。在 Flink 引擎的 DAG 构建流程中变换节点位于 Source 读取之后FeatureBuilder._build()先构建 source 节点若视图定义有变换则接TransformationNode历史检索场景再追加 entity join 节点见 flink/feature_builder.py。五、DAG 各节点在 Flink 中的实现Flink 引擎以 Flink 专属节点实现 Feast 计算 DAG每个节点类型在 flink/nodes.py 中都有对应类。构建顺序与节点职责如下可对照 compute-engine 总览 中的 Feature Builder Flow1. Source 读取节点FlinkSourceReadNode通过create_offline_store_retrieval_job()从 Feast offline store 创建检索任务优先原生 Flink 表若检索任务提供to_flink_table(table_env)直接使用其返回的 Flink Table否则走 Arrow 回退调用to_arrow()拿到pyarrow.Table再经pandas_to_flink_table()转为 Flink Tablesplits_num由pandas_split_num控制若存在字段映射field_mapping会注册临时视图并用 SQLSELECT ... AS重命名列field_mapping处理逻辑同样在 Source 节点内实现。2. 变换节点FlinkTransformationNode把输入 Flink Table 直接传给modeflink的 UDF保留 UDF 返回的原生 Flink Table 输出不经过 pandas 中转。3. 关联节点FlinkJoinNode特征关联与实体关联都通过Flink SQL 临时视图实现先把各输入表注册为临时视图视图名形如__feast_join_uuid再生成多表LEFT JOIN查询ON 条件为 join keys 等值匹配对于历史检索若存在entity_df还会把实体表注册为视图与特征视图做LEFT JOIN保证每个实体行对齐对应特征行。4. 过滤节点FlinkFilterNode将 point-in-timefeature_timestamp entity_ts、TTL 下界、自定义过滤表达式合并为WHERE条件TTL 由_subtract_flink_intervals()生成 FlinkINTERVAL N DAY/HOUR/MINUTE/SECOND字面量表达式无任何条件时直接透传输入避免多余 SQL 包装。5. 聚合节点FlinkAggregationNode支持非窗口Feast 聚合使用 Flink SQL 聚合函数Feast 聚合算子到 SQL 函数的映射见_execute_sql_aggregation()mean/avg → AVG、sum → SUM、min → MIN、max → MAX、count → COUNT、nunique → COUNT(DISTINCT ...)、std → STDDEV_SAMP、var → VAR_SAMP窗口聚合time-windowed会通过aggregation_specs_to_agg_ops()的time_window_unsupported_error_message抛出明确错误。6. 去重节点FlinkDedupNode使用ROW_NUMBER()窗口函数PARTITION BY实体键或内部__feast_entity_row_idORDER BY按 timestamp/created_timestamp 降序取ROW_NUMBER() 1从而在历史检索时为每个实体行保留一条最新特征行。7. 校验节点FlinkValidationNode检查输出是否包含期望的特征列expected_columns缺失即抛ValueErrorJSON 值校验必须在上游 Flink SQL 中处理因为引擎不会把中间数据收集出 Flink 再校验若存在 JSON 类型列需要校验会抛NotImplementedError提示在 Flink SQL 中预先校验或关闭该 FeatureView 的 JSON 校验。8. 输出节点FlinkOutputNode仅物化任务执行写入历史检索只读物化时通过flink_table_to_arrow_batches()把结果按批次online_write_batch_size默认批大小常量10_000见 flink/utils.py流式转成 Arrow再分别写入 online storeonline_write_batch与 offline storeoffline_write_batch由feature_view.online/offline开关控制写入前会剔除内部列ENTITY_ROW_ID_drop_internal_columns。历史检索的 entity_df 支持历史检索get_historical_features接受两类实体数据pandas DataFrame转换时补全 join keys、推断事件时间戳列并追加内部行号列ENTITY_ROW_IDSQL 字符串被解释为针对当前 TableEnvironment/catalog 的Flink SQL 查询结果必须包含event_timestamp列否则抛ValueError引擎用ROW_NUMBER() OVER (ORDER BY ...) - 1生成内部行号并注册临时视图参与后续关联。FlinkDAGRetrievalJobflink/job.py在首次调用to_df()/to_arrow()时惰性执行 DAG结果收集为 Arrow 表后清理所有临时视图其persist()、to_remote_storage()、to_sql()目前均未实现会抛NotImplementedError。六、当前限制与规避方式原文档明确列出两条限制本文补充源码级佐证窗口聚合尚未实现FlinkAggregationNode对带时间窗口的聚合直接报错错误信息见 flink/nodes.py。规避方式使用 Feast 非窗口聚合或在上游 Flink 中预先做窗口化处理。JSON 值校验未实现FlinkValidationNode遇到 JSON 列需要校验时抛NotImplementedError。原因在于引擎不在 Flink 之外收集中间数据做校验。规避方式在 Flink SQL 上游完成 JSON 校验或关闭该 FeatureView 的 JSON 校验enable_validation。另外从实现角度补充两点可预期的行为物化任务在_materialize_one()执行失败时返回MaterializationJobStatus.ERROR并携带异常对象成功则返回SUCCEEDED每次执行结束含异常路径都会调用cleanup_flink_temporary_views()清理临时视图flink/compute.pyupdate()与teardown_infra()为空实现——Flink 引擎不管理 Feast 的共享基础设施仅负责计算执行。七、测试与进一步阅读仓库为 Flink 引擎提供了完整的单元测试sdk/python/tests/unit/infra/compute_engines/flink/test_flink_compute_engine.py约 1200 行覆盖FlinkComputeEngineConfig、各 DAG 节点Source/Join/Filter/Agg/Dedup/Transform/Validation/Output、pandas 与 SQL 两种 entity_df 输入、批式结果收集execute().collect()路径等。测试使用FakeTableEnvironment/FakeFlinkTable模拟 PyFlink 环境无需真实集群即可验证 DAG 构建与 SQL 生成逻辑可作为二次开发或排错的参考。想进一步了解 Feast 计算引擎整体架构、其他后端Spark、Ray、Snowflake、Local 等与自定义计算引擎的接入方式可继续阅读 compute-engine 总览其中也包含ComputeEngine接口签名与FeatureBuilder各 build 方法的扩展模板对理解 Flink 引擎在 Feast 抽象体系中的位置很有帮助。适用前提提示以上配置与行为以当前仓库代码为准。使用 Flink 引擎需要可用的 PyFlink 运行环境含其 Java 依赖引擎面向批式物化与历史检索若需要窗口聚合或 JSON 值校验等能力请按第六节的限制评估是否适合你的场景。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表