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.py | FlinkComputeEngine与FlinkComputeEngineConfig,引擎入口与配置模型 |
feature_builder.py | FlinkFeatureBuilder,把 FeatureView 解析为 Flink 专属 DAG 节点 |
nodes.py | 各DAGNode的 Flink 实现(Source/Join/Filter/Agg/Dedup/Transform/Validation/Output) |
job.py | FlinkDAGRetrievalJob与FlinkMaterializationJob |
utils.py | PyFlink 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"]需要特别说明两点版本约束,它们直接来自仓库配置:
- PyArrow 版本冲突:
flinkextra 要求pyarrow<21,而 Feast 默认安装保持pyarrow>=21。Feast 的 uv lock 将 Flink extra 解析为独立依赖分支([tool.uv.conflicts]声明了flink与ge、ci、dev、docsextra 互斥),因此正常 Feast 安装不会被降级 Arrow。 - 互斥关系:由于上述冲突声明,
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_config、parallelism、execution_mode写入 PyFlinkConfiguration,并据此创建TableEnvironment。
配置选项详解
| Option | Type | Default | Description |
|---|---|---|---|
type | string | flink.engine | 必须为flink.engine,引擎类型选择器。 |
execution_mode | string | batch | PyFlink 执行模式:batch或streaming。 |
parallelism | integer | null | 引擎创建作业的默认 Flink 并行度。 |
table_config | map | null | 额外的 PyFlink table 配置项(键值对)。 |
pandas_split_num | integer | 1 | 把 pandas entity DataFrame 转成 Flink Table 时的 Arrow source 分片数。 |
各选项对应的源码实现(flink/compute.py 中FlinkComputeEngineConfig):
type:Literal["flink.engine"],固定为flink.engine,用于repo_config的引擎注册表查找;execution_mode:Literal["batch", "streaming"],默认batch。在 flink/utils.py 的create_flink_table_environment中,streaming走in_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_num由pandas_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 → 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_id,ORDER 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)。
六、当前限制与规避方式
原文档明确列出两条限制,本文补充源码级佐证:
- 窗口聚合尚未实现:
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.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),仅供参考