Feast MongoDB Offline Store:基于单集合聚合管道的离线特征存储实战指南

发布时间:2026/9/17 21:25:42

Feast MongoDB Offline Store:基于单集合聚合管道的离线特征存储实战指南 Feast MongoDB Offline Store基于单集合聚合管道的离线特征存储实战指南【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feastFeast 的 MongoDB Offline Store 是一个将 MongoDB 作为离线特征存储Offline Store的贡献型contrib实现允许你直接基于 MongoDB 中的数据训练模型、运行批量打分batch scoring并在读取路径上采用 MongoDB 聚合管道与复合索引将单实体查询成本控制在 O(log n) 量级。读完本文你将掌握该离线存储的单集合数据模型、复合索引设计、feature_store.yaml配置方式、pull_latest/get_historical_features两条读取路径的实现原理评分路径与服务端去重、训练路径与merge_asof、strict_pit语义以及写入与内存行为并了解其配套单元测试对关键行为的验证方式。设计定位单集合、共享 schema 的离线特征存储该离线存储的完整实现位于 mongodb.py模块文档对其设计做了三点概括单集合single-collectionschema所有 Feature View 共享同一个 MongoDB collection默认名为feature_history通过文档中的feature_view字段区分彼此服务端去重scoring path当entity_df中实体 ID 唯一时聚合管道追加$group阶段每个(entity_id, feature_view)最多返回一条文档传输量由 O(N×P×K) 降为 O(N×K)N 为实体数P 为每实体观测数K 为 Feature View 数复合索引支撑整条管道索引(entity_id ASC, feature_view ASC, event_timestamp DESC, created_at DESC)使单实体查询成本从 O(P) 降为 O(log P)。数据模型所有 Feature View 共享的feature_history集合所有 Feature View 的观测数据写入同一个 collection默认feature_history由feature_view字段做判别discriminator。README 给出了文档的 JSON 形态// Collection: feature_history { entity_id: Binary(...), // Serialized entity key (bytes) feature_view: driver_stats, // Discriminator features: { // Nested subdocument trips_today: 5, rating: 4.8 }, event_timestamp: ISODate(2024-01-15T10:00:00Z), created_at: ISODate(2024-01-15T10:00:01Z) }各字段的含义与实现对应关系如下字段类型说明源码依据entity_idBinarybytes序列化后的实体键entity key由 Feast 的serialize_entity_key生成mongodb.py中_serialize_entity_key_from_row与_serfeature_viewstring判别字段即数据源MongoDBSource的nameMongoDBSource.feature_view_name返回self.namefeatures嵌套子文档特征名 → 特征值 的映射offline_write_batch逐列构造event_timestampISODate事件时间戳Point-in-Time 连接的核心字段默认timestamp_fieldevent_timestampcreated_atISODate写入时间戳用于同事件时间下的冲突裁决取最新默认created_timestamp_columncreated_at注意entity_id在 MongoDB 中是序列化后的字节串不直接存可读的 join key 列这正是读取时需要_expand_entity_id_column反序列化、把 join key 还原成独立列的原因。反序列化逻辑见mongodb.py的_expand_entity_id_column它调用deserialize_entity_key将字节展开为各 join key 列后再输出。复合索引一次懒创建支撑全部查询存储层在首次使用时懒创建一个复合索引源码中_ensure_indexes使用create_index(..., nameentity_fv_ts_idx, backgroundTrue)并通过模块级缓存_indexes_ensured集合f{conn_str}/{db}/{collection}去重避免每次调用都重复建索引。索引定义如下db.feature_history.createIndex({ entity_id: 1, feature_view: 1, event_timestamp: -1, created_at: -1 })从源码结构可以推断该索引的每一列都服务于具体管道阶段entity_id升序支撑$match: {entity_id: {$in: [...]}}的实体点查feature_view升序配合$match中的feature_view判别过滤event_timestamp降序支撑$lte: max_ts的时间窗口过滤以及$sort中的时间降序created_at降序支撑$sort/$group $first对同一事件时间下“最新写入”的选择。配置接入feature_store.yaml在 Feast 的feature_store.yaml中通过offline_store段启用该离线存储offline_store: type: feast.infra.offline_stores.contrib.mongodb_offline_store.mongodb.MongoDBOfflineStore connection_string: mongodb://localhost:27017 database: feast collection: feature_history # optional, default: feature_history对应源码中的MongoDBOfflineStoreConfigFeastConfigBaseModel子类三个可配置项及默认值如下配置项默认值说明typefeast.infra.offline_stores.contrib.mongodb_offline_store.mongodb.MongoDBOfflineStore离线存储实现类路径connection_stringmongodb://localhost:27017MongoDB 连接 URI支持带认证与多副本的 URI 形式databasefeastMongoDB 数据库名collectionfeature_history所有 Feature View 共享的集合名在 Feature View 定义中batch source 需要使用MongoDBSource其name即文档判别字段from feast.infra.offline_stores.contrib.mongodb_offline_store.mongodb import MongoDBSource driver_source MongoDBSource( namedriver_stats, timestamp_fieldevent_timestamp, created_timestamp_columncreated_at, )MongoDBSource继承自 Feast 的DataSourcesource_type()返回CUSTOM_SOURCE其 proto 序列化将{feature_view: self.name}写入custom_options.configuration见_to_proto_impl。读取路径一pull_latest_from_table_or_querypull_latest_from_table_or_query返回时间窗口内每个实体最新一条观测管道形如$match → $sort → $group($first) → $project见mongodb.py$match按feature_view判别 event_timestamp的$gte/$lte时间窗口$sortentity_id升序、event_timestamp降序、created_at降序——保证“最新写入”排最前$group按entity_id分组取$first每个实体仅返回一条$project展开features子文档的各特征列并保留event_timestamp可选created_at。与之配套的pull_all_from_table_or_query则只做$match $project返回窗口内全部原始行、不去重通常服务于离线特征物料training data的原始抽取。读取路径二get_historical_features的评分路径与训练路径get_historical_features是训练/批量打分的主入口源码根据entity_df的形态在两条路径间按 Feature View 逐条自动切换Scoring path评分路径——当entity_df中实体 ID 唯一且strict_pitTrue时所有实体请求时间戳相同管道为$match $sort $group在 MongoDB 服务端完成去重每个(entity_id, feature_view)至多返回一条文档复合索引使单实体成本为 O(log P)且避免了把每个实体的全部历史观测拉到 Python 侧Python 侧随后做一次向量化 left join并施加future_mask严格 PIT 时把晚于请求时间的文档置NULL。Training path训练路径——当entity_df存在重复实体 ID 且位于不同时间戳典型的 PIT 训练数据形态省略$group阶段将候选文档与实体表按实体键分组后在 Python 中执行pandas.merge_asofdirectionbackward做逐行 Point-in-Time 连接该操作由 pandas 底层的 C 实现优化为正确处理“同一事件时间、不同写入时间”的冲突fv_df会先按[event_timestamp, created_at]排序保证merge_asof命中created_at最新的文档对应测试test_training_path_created_at_tiebreaker。路径选择的核心判定代码位于get_historical_features的_run_single中unique_entities result[eid_col].nunique() len(result) scoring_path unique_entities and ( not strict_pit or result[event_timestamp_col].nunique() 1 )即实体键唯一 请求时间戳一致或strict_pitFalse时走服务端$group否则回退merge_asof。这保证了不同 Feature View 可以混合使用不同路径见测试test_mixed_join_key_cardinalityuser_id维度的 FV 走merge_asof而(user_id, device_id)维度的 FV 仍可走评分路径。strict_pit参数语义get_historical_features接受strict_pit关键字参数默认Truestrict_pitTrue默认训练/评估安全文档时间戳严格晚于实体请求时间戳的观测被返回为NULL避免“未来泄漏”future leakagestrict_pitFalse用于实时推理real-time inference场景始终返回该实体最新的观测即使其时间戳晚于名义请求时间。对应的测试用例覆盖了三种场景见测试文件test_mongodb.pytest_scoring_path_nulls_future_doc验证未来文档被置NULLtest_scoring_path_nulls_future_doc_chunk_size_1验证在分块_CHUNK_SIZE50_000默认值边界下行为一致test_strict_pit_false_returns_future_doc验证strict_pitFalse时未来文档会被返回。源码在$match阶段用ts_filter {$lte: max_ts} if strict_pit else {}做服务端过滤Python 侧再用future_mask兜底置空。查询折叠Query-collapse从 K 次往返降为一次README 强调的核心优化是Query-collapse共享相同 join key 集合join key signature的多个 Feature View会被分组成一次 MongoDB 聚合往返而不是每个 Feature View 一次。往返次数从 KFeature View 数降为“唯一 join key 签名数”——在常见场景下即为 1。从源码看get_historical_features先构造fv_by_proj按投影名name_to_use()索引 Feature View、fv_mongo_name投影名 → MongoDB 判别值、fv_mapped_join_keys投影名 → 映射后的 join key等映射表再对每个投影逐一遍历执行管道同一 join key 集合的 Feature Views 共享一次$match {entity_id: {$in: batch_ids}}的实体 ID 批量查询每个批次的实体 ID 上限由_MONGO_BATCH_SIZE 10_000控制。测试test_k_collapse_multiple_feature_views验证了driver_stats_k与vehicle_stats_k两个共享driver_id的 Feature View 在同一次检索中被正确解析。此外当entity_df行数超过_CHUNK_SIZE 50_000时数据会被分块处理_chunk_dataframe各块结果按原始行号_row_idx排序拼接保证输出顺序与输入一致。写入数据offline_write_batch与feast materialize写入侧使用offline_write_batch它由feast materialize自动调用。README 给出直接调用方式store.write_to_offline_store(feature_view_name, df)写入语义为纯追加append-only不做 upsert冲突在读取时裁决——pull_latest与评分路径均通过$sort created_at DESC → $group $first或 merge_asof 前的created_at排序选择created_at最高的文档。从源码看offline_write_batch的处理流程使用原始未映射join key 名序列化实体键确保与get_historical_features的序列化字节一致源码注释明确说明该点测试test_offline_write_batch_round_trip和test_int32_entity_key均验证了“写入字节 读取字节”时间戳列统一规范化为带 UTC 时区的datetime特征值逐行提取NaN 跳过、numpy 标量经.item()转为 Python 原生类型、非标量list/dict保留created_at缺省时取当前 UTC 时间按每批 10,000 条insert_many(..., orderedFalse)写入progress回调在每批后上报行数。文档判别值取自feature_view.batch_source.feature_view_name因此通过 push /write_to_offline_store写入的数据与初始 ingest 落在同一集合分区。内存行为$match先行而非全量加载README 明确指出存储的内存占用特性存储按实体键在$match阶段过滤而不是把整个集合加载到内存。因此内存占用上界为“唯一实体 ID 数 × 每实体文档数”与集合总体积无关。这一特性对海量历史数据尤其重要——$match配合复合索引在 MongoDB 服务端完成裁剪Python 侧只接收与请求实体相关的候选文档。结果持久化与类型推断结果落盘MongoDBRetrievalJob.persist()支持将检索结果写为 Parquet 文件需配合SavedDatasetFileStorage目标文件已存在且未设置allow_overwriteTrue时抛出SavedDatasetLocationAlreadyExists对应测试test_persist_writes_parquet、test_persist_raises_if_file_exists、test_persist_allow_overwrite。列类型推断MongoDBSource.get_table_column_names_and_types通过读取feature_view对应的一条样本文档来推断列名与类型——event_timestamp/created_at映射为datetimefeatures子文档中的值按 Python 类型映射为bool/int64/float64/string/list/dict/object。类型字符串到 FeastValueType的映射见 type_map.py 的mongodb_to_feast_value_type如int64→INT64、float64→DOUBLE、list[int]→INT64_LIST无法识别的类型映射为UNKNOWN。测试验证与使用前提该存储的单元测试位于 test_mongodb.py使用testcontainers.mongodb.MongoDbContainermongo:latest起真实 MongoDB 实例Docker 不可用时相关用例被_requires_docker跳过。测试覆盖了本文涉及的几乎所有关键行为pull_latest每实体最新行、训练路径逐行 PIT 连接、created_at平局裁决、评分路径服务端去重、pull_all全量窗口、TTL 过期置NULL、K-collapse 多 Feature View 合并、混合 join key 基数、异构时间戳回退训练路径、重叠特征名full_feature_namesTrue时的fv__feature列命名、复合 join key、entity_df 额外标签列不污染实体键序列化、INT32 实体键字节一致性、persist 行为、offline_write_batch写读往返以及strict_pit三态语义。使用前请注意以下前提与限制均以当前仓库源码为准依赖pymongo未安装时会抛出FeastExtrasDependencyImportError(pymongo, mongodb)get_historical_features不支持 SQL 字符串形式的entity_df传入字符串会抛出ValueError请使用 pandas DataFrameentity_key_serialization_version必须与写入侧一致否则字节不匹配会导致查询结果为空测试中统一使用 version 3TTLFeature View 的ttl在读取侧执行评分路径与训练路径都会把超过 TTL 的过期观测置NULL对应测试test_ttl_excludes_stale_features该离线存储位于contrib社区贡献目录属于自定义离线存储实现可参考 adding-a-new-offline-store.md 了解 Feast 离线存储扩展点OfflineStore抽象类定义了pull_latest_from_table_or_query、get_historical_features、offline_write_batch等接口。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/17 21:25:42

激光技术课件自动化:python-pptx、M²拟合与交付自检

简介:这份《专题一 激光技术.ppt》面向物理、光电信息、电子工程等专业的学生与初入激光领域的自学者,用于系统梳理激光原理与技术脉络。课件从爱因斯坦1916年提出受激辐射讲起,串联汤斯与肖洛的经典论文、梅曼的红宝石激光器、He-Ne气体激光…

2026/9/17 22:10:54

Spring Boot+Vue数字化物资管理系统设计与实现

1. 项目背景与核心价值在灾害救援场景中,物资管理效率直接关系到受灾群众的生命安全。去年参与某地洪灾救援时,我亲眼目睹了传统纸质台账导致的物资调配混乱:救援队需要3小时才能确认库存情况,而受灾群众等待帐篷和药品的时间超过…

2026/9/17 22:10:54

Docker部署ONLYOFFICE文档服务:Nginx反代与HTTPS配置全攻略

1. 为什么选Docker部署ONLYOFFICE,而不是裸机装1.1 官方镜像到底装了什么,一个容器等于一整套服务我第一次接触ONLYOFFICE时,下意识以为它和LibreOffice差不多,装个依赖、跑个服务、配置一下就能用。真正上手才发现,ON…

2026/9/17 22:10:54

Transformer架构核心组件与工业实践全解析

1. Transformer架构全景解析2017年那篇《Attention Is All You Need》论文扔进NLP领域就像颗核弹,把传统RNN架构炸得七零八落。我在实际业务中部署Transformer模型时发现,相比LSTM那些需要顺序处理的"老古董",自注意力机制让长距离…

2026/9/17 22:10:54

2023上半年软考数据库系统工程师上午真题复盘与备考指南

“数据库系统工程师”这个中级资格在软考里一直是报考大户,2023年上半年那场考试更是特别——它是软考机考改革前一场大规模的传统笔试,上午的《基础知识》科目仍然沿用75道单选题、150分钟、45分及格的老规矩。我在考场上最大的感受是:这张卷…

2026/9/17 22:05:53

发那科机器人报警代码详解:从紧急停止到伺服与编码器排查

简介:面向FANUC发那科工业机器人维护与调试人员,这份中文故障代码与报警处理全集,系统梳理了伺服系统中最常见的紧急停止与报警类型,覆盖SRVO-001操作面板紧急停止、SRVO-002示教操作盘紧急停止、SRVO-003紧急时自动停机开关、SRV…

2026/9/16 12:52:37

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/17 0:03:13

WiFi密码安全测试:从原理到实战的字典暴力破解指南

1. 写在前面:我为什么要研究WiFi密码这件事先交代一下背景。我身边有不少朋友,家里的WiFi密码常年是"12345678"或者"88888888",问就是"好记"。直到有一次,隔壁邻居蹭网蹭到我家路由器后台都进不去&…

2026/9/17 0:03:13

redis-py服务控制与监控函数实战:从ping到slowlog的巡检指南

我用 redis-py 写了快五年的业务代码,坦白说,真正让我觉得这个客户端“像一个成熟工具箱”的,不是 get/set 那套基本操作,而是它那批专门做服务控制与状态监控的辅助函数。日常开发里,大家把redis.Redis(host..., deco…

2026/9/17 0:03:13

SpringBoot+Vue3实现中小企业设备管理系统开发实践

1. 项目概述与核心价值中小企业设备管理系统是制造业、服务业等领域的基础信息化工具。传统设备管理往往依赖Excel表格或纸质记录,存在数据孤岛、流程混乱、维护成本高等痛点。这套基于Java SpringBootVue3MyBatis的技术方案,通过前后端分离架构实现了设…

2026/9/16 22:55:57

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/16 22:56:09

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/16 22:56:16

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码