news 2026/9/17 23:28:38

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Feast Flink 计算引擎(FlinkComputeEngine)实战指南:用 PyFlink 分布式执行特征物化与历史检索

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

【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast

Apache Flink 计算引擎是 Feast 中基于 PyFlink Table API 的分布式计算后端,它实现 Feast 统一的ComputeEngine接口,可承担批量物化(materialize/materialize-incremental)与历史特征检索(get_historical_features)两类核心任务。本文以 docs/reference/compute-engine/flink.md 为骨架,结合仓库源码与测试用例,系统讲解其安装配置、feature_store.yaml参数、mode="flink"特征变换写法、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_table);
  • join、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.pyFlinkComputeEngineFlinkComputeEngineConfig,引擎入口与配置模型
feature_builder.pyFlinkFeatureBuilder,把 FeatureView 解析为 Flink 专属 DAG 节点
nodes.pyDAGNode的 Flink 实现(Source/Join/Filter/Agg/Dedup/Transform/Validation/Output)
job.pyFlinkDAGRetrievalJobFlinkMaterializationJob
utils.pyPyFlink TableEnvironment 创建、pandas↔Flink 转换、临时视图管理

二、安装与依赖:flinkextra

Flink 引擎依赖 PyFlink,Feast 通过flinkextra 提供安装入口。在 Feast 源码检出目录下用uv安装:

uv sync --extra flink --no-dev

关于依赖约束,仓库根目录 pyproject.toml 中flinkextra 定义为:

flink = ["apache-flink>=2.2.1,<3", "pyarrow<21.0.0"]

需要特别说明两点版本约束,它们直接来自仓库配置:

  1. PyArrow 版本冲突flinkextra 要求pyarrow<21,而 Feast 默认安装保持pyarrow>=21。Feast 的 uv lock 将 Flink extra 解析为独立依赖分支([tool.uv.conflicts]声明了flinkgecidevdocsextra 互斥),因此正常 Feast 安装不会被降级 Arrow。
  2. 互斥关系:由于上述冲突声明,flinkextra 不应与其他列出的 extra(ge、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_configparallelismexecution_mode写入 PyFlinkConfiguration,并据此创建TableEnvironment

配置选项详解

OptionTypeDefaultDescription
typestringflink.engine必须为flink.engine,引擎类型选择器。
execution_modestringbatchPyFlink 执行模式:batchstreaming
parallelismintegernull引擎创建作业的默认 Flink 并行度。
table_configmapnull额外的 PyFlink table 配置项(键值对)。
pandas_split_numinteger1把 pandas entity DataFrame 转成 Flink Table 时的 Arrow source 分片数。

各选项对应的源码实现(flink/compute.py 中FlinkComputeEngineConfig):

  • typeLiteral["flink.engine"],固定为flink.engine,用于repo_config的引擎注册表查找;
  • execution_modeLiteral["batch", "streaming"],默认batch。在 flink/utils.py 的create_flink_table_environment中,streamingin_streaming_mode(),其余走in_batch_mode()
  • parallelism:可选整数,非空时写入 Flink 配置键parallelism.default
  • table_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 转换时的并行分片数。

四、用 mode="flink" 编写 Flink 特征变换

BatchFeatureView的变换函数需要接收并返回 PyFlink Table 对象时,使用mode="flink"

from feast import BatchFeatureView, Field from feast.types import Float32 def double_rates(table): # 生产环境中可以在这里使用 PyFlink Table API 操作并返回一个 table。 return table driver_stats = BatchFeatureView( name="driver_stats", entities=[driver], mode="flink", udf=double_rates, schema=[Field(name="conv_rate", dtype=Float32)], source=driver_stats_source, online=True, )

约束mode="flink"的变换函数必须返回 PyFlink Table 对象;返回 pandas DataFrame 的 UDF 不被 Flink 计算引擎接受。源码依据:

  • batch_feature_view.py 的get_feature_transformation()mode="flink"归入支持的变换模式列表;
  • flink/nodes.py 中FlinkTransformationNode.execute()直接调用self.transformation_fn(*input_tables)并把返回值当作 Flink Table,若返回对象无法取得 schema 则抛出TypeError
  • feature_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 Flow):

1. 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 Table(splits_numpandas_split_num控制);
  • 若存在字段映射(field_mapping),会注册临时视图并用 SQLSELECT ... AS重命名列(field_mapping处理逻辑同样在 Source 节点内实现)。

2. 变换节点(FlinkTransformationNode)

  • 把输入 Flink Table 直接传给mode="flink"的 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-time(feature_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 → AVGsum → SUMmin → MINmax → MAXcount → COUNTnunique → COUNT(DISTINCT ...)std → STDDEV_SAMPvar → 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),缺失即抛ValueError
  • JSON 值校验必须在上游 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 store(online_write_batch)与 offline store(offline_write_batch),由feature_view.online/offline开关控制;
  • 写入前会剔除内部列ENTITY_ROW_ID_drop_internal_columns)。

历史检索的 entity_df 支持

历史检索(get_historical_features)接受两类实体数据:

  • pandas DataFrame:转换时补全 join keys、推断事件时间戳列,并追加内部行号列ENTITY_ROW_ID
  • SQL 字符串:被解释为针对当前 TableEnvironment/catalog 的Flink SQL 查询,结果必须包含event_timestamp列(否则抛ValueError);引擎用ROW_NUMBER() OVER (ORDER BY ...) - 1生成内部行号,并注册临时视图参与后续关联。

FlinkDAGRetrievalJob(flink/job.py)在首次调用to_df()/to_arrow()时惰性执行 DAG,结果收集为 Arrow 表后清理所有临时视图;其persist()to_remote_storage()to_sql()目前均未实现(会抛NotImplementedError)。

六、当前限制与规避方式

原文档明确列出两条限制,本文补充源码级佐证:

  1. 窗口聚合尚未实现FlinkAggregationNode对带时间窗口的聚合直接报错(错误信息见 flink/nodes.py)。规避方式:使用 Feast 非窗口聚合,或在上游 Flink 中预先做窗口化处理。
  2. JSON 值校验未实现FlinkValidationNode遇到 JSON 列需要校验时抛NotImplementedError。原因在于引擎不在 Flink 之外收集中间数据做校验。规避方式:在 Flink SQL 上游完成 JSON 校验,或关闭该 FeatureView 的 JSON 校验(enable_validation)。

另外从实现角度补充两点可预期的行为:

  • 物化任务在_materialize_one()执行失败时返回MaterializationJobStatus.ERROR并携带异常对象,成功则返回SUCCEEDED;每次执行结束(含异常路径)都会调用cleanup_flink_temporary_views()清理临时视图(flink/compute.py);
  • update()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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 23:26:38

IDC运维工程师面试:Linux、MySQL、Redis、Docker排障

简介&#xff1a;面向 IDC 机房运维岗位求职者与初级运维工程师的面试备考资料&#xff0c;以一份 PDF 问答文档形式呈现&#xff0c;覆盖 Windows、Linux 与网络基础三大知识板块。内容按基础技能测试题组织&#xff0c;逐条给出参考答案&#xff0c;涉及远程登录工具与端口辨…

作者头像 李华