尧图精选

Feast 基础设施公共抽象层 feast.infra.common 源码级解析:物化任务、历史检索任务与序列化机制

🕒 发布时间:2026/9/18 6:57:20 📁 来源:尧图网络
Feast 基础设施公共抽象层 feast.infra.common 源码级解析物化任务、历史检索任务与序列化机制【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast导读本文聚焦 Feast开源特征平台Python SDK 中负责承载基础设施通用抽象的核心包feast.infra.common。它位于 Provider 与各计算引擎本地、Spark、Ray、Flink、Snowflake、AWS Lambda、Kubernetes之间的公共边界以三个模块分别定义了物化任务MaterializationTask / MaterializationJob、历史特征检索任务HistoricalRetrievalTask与不可序列化构件的传输方案SerializedArtifacts。读完本文你将理解 Feast 如何用统一的任务描述与状态机抽象屏蔽不同离线/在线存储与计算引擎的差异并掌握这些类在实际代码调用链中的位置与用法。一、feast.infra.common 包的整体定位在 Feast 的 Python SDK 中sdk/python/docs/source/feast.infra.common.rst对应的包位于sdk/python/feast/infra/common/目录下仅包含 4 个文件sdk/python/feast/infra/common/ ├── __init__.py # 空文件仅用于声明 Python 包 ├── materialization_job.py # 物化任务与物化任务状态/抽象类 ├── retrieval_task.py # 历史特征检索任务 └── serde.py # 序列化辅助SerializedArtifacts从源码结构看该包是一个纯抽象与数据定义层它不依赖任何具体的离线存储、在线存储或计算引擎实现而是被 ComputeEngine 抽象基类 和 PassthroughProvider 等上层组件引用作为它们与各引擎实现之间传递数据的契约。该包被广泛引用涉及feast/infra/compute_engines/下几乎所有引擎local、spark、ray、flink、snowflake、aws_lambda、kubernetes以及feast/feature_store.py的 CLI 入口层可见它是 Feast 物化与历史检索两大核心流程的公共语言。二、materialization_job 模块物化任务的定义、状态机与作业抽象materialization_job.py是整个包的灵魂定义了三个关键构件分别回答要物化什么物化进行到哪一步如何管理一次物化作业三个问题。2.1 MaterializationTask一次物化的最小工作单元MaterializationTask是一个dataclass文档注释明确指出它代表需要从离线存储物化到在线存储的一份数据单元dataclass class MaterializationTask: project: str feature_view: Union[BatchFeatureView, StreamFeatureView, FeatureView] start_time: datetime end_time: datetime only_latest: bool True tqdm_builder: Union[None, Callable[[int], tqdm]] None disable_event_timestamp: bool False各字段含义与使用要点字段类型说明projectstr任务所属的 Feast 项目名用于在注册表Registry中定位对象feature_viewUnion[BatchFeatureView, StreamFeatureView, FeatureView]待物化的特征视图决定数据来源与转换逻辑start_time/end_timedatetime物化的时间窗口边界离线存储会按该窗口拉取数据only_latestbool默认True是否只物化每个实体键的最新特征值False时保留窗口内全部历史点tqdm_builderOptional[Callable[[int], tqdm]]进度条构建回调供 CLI 场景展示物化进度disable_event_timestampbool默认False是否在物化时禁用事件时间戳过滤逻辑单元测试 test_materialization_job.py 印证了这些默认行为only_latest与disable_event_timestamp在仅传必填参数时分别默认为True与False且任务创建后可原样读取project、feature_view、start_time、end_time字段。2.2 MaterializationJobStatus物化作业的九态状态机作业状态是一个enum.Enum共 9 个取值覆盖从排队到成功/取消/失败的完整生命周期class MaterializationJobStatus(enum.Enum): WAITING 1 # 排队等待调度 RUNNING 2 # 正在执行 AVAILABLE 3 # 结果已就绪可消费 ERROR 4 # 执行出错 CANCELLING 5 # 取消中 CANCELLED 6 # 已取消 SUCCEEDED 7 # 成功完成 PAUSED 8 # 已暂停 RETRYING 9 # 重试中该枚举的完整性同样被test_materialization_job.py中的test_all_statuses_defined用例锁定防止后续迭代中意外丢失状态值。2.3 MaterializationJob物化作业的抽象接口MaterializationJob继承ABC定义了所有计算引擎物化作业必须实现的 5 个抽象方法class MaterializationJob(ABC): task: MaterializationTask # 该作业对应的任务描述 abstractmethod def status(self) - MaterializationJobStatus: ... abstractmethod def error(self) - Optional[BaseException]: ... abstractmethod def should_be_retried(self) - bool: ... abstractmethod def job_id(self) - str: ... abstractmethod def url(self) - Optional[str]: ...这一接口设计使得上层调用方如PassthroughProvider无需关心作业跑在本地进程还是远端 Spark/Snowflake 集群只需通过统一的status()/error()/should_be_retried()等方法来轮询与判定结果。三、retrieval_task 模块历史特征检索任务的载体HistoricalRetrievalTask是历史特征检索训练数据生成即get_historical_features流程中从 Provider 传递到计算引擎/离线存储的任务描述dataclass class HistoricalRetrievalTask: project: str entity_df: Union[pd.DataFrame, str] # 实体 DataFrame 或其 SQL 字符串 feature_view: Union[BatchFeatureView, StreamFeatureView] full_feature_name: bool # 是否返回带视图名前缀的完整特征名 registry: Registry # 特征注册表用于解析实体/特征定义 start_time: Optional[datetime] None # 可选时间窗口下界 end_time: Optional[datetime] None # 可选时间窗口上界值得注意的两个设计点entity_df支持两种形态既可以是内存中的pandas.DataFrame也可以是对离线存储如 BigQuery/Snowflake的 SQL 字符串这为不同引擎的下推执行留下了空间registry直接内嵌在任务中因为历史检索需要实时解析实体Entity与特征视图FeatureView的定义将 Registry 随任务传递可避免再次从外部读取注册表。在 ComputeEngine.get_execution_context 中可以看到它的消费方式引擎从task.feature_view.entities逐个registry.get_entity(...)拉取实体定义并透传task.entity_df共同构建 DAG 执行上下文ExecutionContext。四、serde 模块跨进程/跨引擎传递不可序列化构件serde.py解决一个非常实际的问题当物化/检索作业需要被提交到独立进程或远端计算引擎如 Spark Application、Ray 集群时Python 对象必须可序列化但FeatureView与RepoConfig往往携带无法直接 pickled 的元素如 UDF 函数对象。SerializedArtifacts提供了一套标准化的换乘方案dataclass class SerializedArtifacts: Class to assist with serializing unpicklable artifacts to be passed to the compute engine. feature_view_proto: str # FeatureView 序列化为 protobuf 字符串 repo_config_byte: str # RepoConfig 通过 dill 序列化为字节串 classmethod def serialize(cls, feature_view, repo_config): feature_view_proto feature_view.to_proto().SerializeToString() repo_config_byte dill.dumps(repo_config) return SerializedArtifacts( feature_view_protofeature_view_proto, repo_config_byterepo_config_byte ) def unserialize(self): proto FeatureViewProto() proto.ParseFromString(self.feature_view_proto) # skip_udfTrue写入节点只需 schema 与实体元数据 feature_view FeatureView.from_proto(proto, skip_udfTrue) repo_config dill.loads(self.repo_config_byte) provider PassthroughProvider(repo_config) online_store provider.online_store offline_store provider.offline_store return feature_view, online_store, offline_store, repo_config其序列化策略非常讲究FeatureView 走 protobuf 路径利用FeatureView.to_proto()转成 protobuf 二进制字符串反序列化时用FeatureView.from_proto(proto, skip_udfTrue)跳过 UDF 还原——因为下游写入节点只需要 schema 与实体元数据无需真正的转换函数RepoConfig 走 dill 路径RepoConfig包含大量嵌套配置对象dill比标准pickle覆盖更广的对象类型反序列化时重建完整 Provider拿到repo_config后直接构造PassthroughProvider并暴露online_store/offline_store让远端进程无需任何外部上下文即可初始化存储连接。这解释了为何该模块同时 import 了feast.protos.feast.core.FeatureView_pb2与PassthroughProvider——它是为远端引擎独立重建执行环境这一场景量身定做的。五、调用链全景从 FeatureStore 到公共抽象再到引擎实现将上述三个模块串联起来的核心枢纽是 PassthroughProvider 与 ComputeEngine。5.1 物化路径materializePassthroughProvider.materialize_single_feature_viewpassthrough_provider.py是标准的逐视图物化入口校验feature_view类型BatchFeatureView/StreamFeatureView/FeatureView均可OnDemandFeatureView需开启write_to_online_store将参数组装为MaterializationTask(project, feature_view, start_time, end_time, tqdm_builder, disable_event_timestamp)交给self.batch_engine.materialize(registry, task)返回作业列表断言len(jobs) 1若jobs[0].status() MaterializationJobStatus.ERROR则抛出jobs[0].error()。而在 ComputeEngine.materialize 中单个任务会被规范化为列表逐一调用引擎自实现的_materialize_one(registry, task)最终返回List[MaterializationJob]。以本地引擎为例local/compute.py 的_materialize_one会用f{task.feature_view.name}-{task.start_time}-{task.end_time}生成job_id构建LocalFeatureBuilder生成执行计划并执行成功返回SUCCEEDED状态、失败返回ERROR状态并携带异常——其返回的 LocalMaterializationJob 就是一个完整实现MaterializationJob接口status/error/should_be_retried/job_id/url的具体类。5.2 历史检索路径get_historical_featuresPassthroughProvider.get_historical_featurespassthrough_provider.py在未指定计算引擎时直接委托offline_store.get_historical_features(...)而引擎路径则由ComputeEngine.get_historical_features(registry, task)接收HistoricalRetrievalTask并返回RetrievalJob或pa.Table。本地引擎的 LocalRetrievalJob 采用惰性执行策略构造时只构建执行计划直到调用to_df()/to_arrow()时才真正执行并以error字段承载构建阶段的失败。六、测试佐证公共抽象的契约被测试锁定仓库中与feast.infra.common直接相关的单元测试位于 test_materialization_job.py覆盖两层契约状态机完整性MaterializationJobStatus必须恰好包含 9 个成员WAITING、RUNNING、AVAILABLE、ERROR、CANCELLING、CANCELLED、SUCCEEDED、PAUSED、RETRYING任务默认行为MaterializationTask的only_latest默认为True、disable_event_timestamp默认为False。此外sdk/python/tests/component/下的spark/test_compute*.py、ray/test_compute.py以及unit/infra/compute_engines/下的test_local_job.py、test_spark_application.py、flink/test_flink_compute_engine.py等用例均在各自引擎实现层面间接验证了这些公共抽象被正确消费。七、实战视角何时会与这些抽象打交道对大多数 Feast 使用者而言这些类不会直接出现在业务代码中但理解它们有助于排查三类问题物化失败定位MaterializationJobStatus的ERROR/RETRYING/CANCELLED等状态直接对应feast materialize命令的输出与日志通过job.error()可拿到底层异常历史检索性能分析HistoricalRetrievalTask中的entity_df形态DataFrame vs SQL决定了查询是否可下推到离线存储执行是调优get_historical_features的关键切入点自研计算引擎/离线存储接入如果要参照 adding-a-new-offline-store.md 与 creating-a-custom-compute-engine.md 扩展 Feast只需实现ComputeEngine接口并正确构造/消费MaterializationTask、HistoricalRetrievalTask、MaterializationJob即可无缝接入现有 Provider 与 CLI 流程而SerializedArtifacts则是把任务投递给远端引擎时的标准序列化方案。结语feast.infra.common虽然是一个仅有三个模块的小包却是 Feast 基础设施层面向接口编程的核心体现materialization_job定义了物化任务与作业状态契约retrieval_task封装了历史检索的输入描述serde解决了远端引擎间的对象传输问题。三者共同构成了 Provider、计算引擎与存储层之间稳定、可测试的公共边界让 Feast 得以在本地、Spark、Ray、Flink、Snowflake、AWS Lambda 与 Kubernetes 等异构环境下保持统一的物化与检索语义。延伸阅读仓库内路径公共抽象定义materialization_job.py、retrieval_task.py、serde.py抽象消费方base.py、passthrough_provider.py本地引擎实现local/compute.py、local/job.py单元测试test_materialization_job.py相关扩展指南adding-a-new-offline-store.md、creating-a-custom-compute-engine.md【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联 返回资讯列表 →