Feast Aerospike Online Store 接入指南:配置、数据模型与实现原理(Preview)

Feast Aerospike Online Store 接入指南:配置、数据模型与实现原理(Preview) Feast Aerospike Online Store 接入指南配置、数据模型与实现原理Preview【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast导读本文以 Feast 官方文档 docs/reference/online-stores/aerospike.md 为骨架结合仓库内 Aerospike 在线存储的完整实现源码 aerospike.py 与对应测试用例系统讲解如何在 Feast 项目中把 Aerospike 作为在线存储online store从基础配置、集群参数、认证与 TLS到按 Feature View 粒度的 namespace/set 覆盖、prewriting hook 扩展再到记录模型、TTL 语义、异步读写与功能矩阵。读完本文你将能够独立完成feature_store.yaml的 Aerospike 配置并理解其底层 Map CDT 数据布局与关键设计取舍。预览状态提示Aerospike online store 当前处于preview阶段。部分功能可能不稳定未来版本可能出现破坏性变更breaking changes。1. 功能总览Aerospike online store 负责将特征值物化materialize进 Aerospike 集群用于在线特征服务。其核心能力如下同步与异步读写双路径同时支持online_read/online_read_async、online_write_batch/online_write_batch_async。异步方法通过run_in_executor将阻塞式客户端调用放到线程池执行保持 feature-server 负载下事件循环event loop的响应性。基于 Aerospike Map CDT 的部分服务端 upsert写入某个 feature view 永远不会覆盖同一实体上其他 feature view 的数据。记录级 TTL由单一ttl_seconds配置项控制兼容 namespace 默认 TTL、永不过期哨兵值或显式秒数三种语义。按 Feature View 的 namespace 覆盖与 set 覆盖可将单个 feature view 固定到纯内存RAM-only或 SSD 支持的 namespace或将某个 view 隔离到独立 set而无需拆分项目。Prewriting hook一个通过 import 字符串解析的可配置回调应用于每次写批write batch用于 PII 脱敏、应用侧加密、值强转等横切关注点。认证与 TLS面向 Aerospike Enterprise Edition 的认证与 TLS 选项原样透传给 Aerospike Python 客户端。client_kwargs逃生舱AerospikeOnlineStoreConfig未暴露的任何高级客户端配置字段均可通过该参数传入。基线要求Aerospike Server≥ 6.0使用 batch-write / batch-operate API该存储基于 CE 8.x 开发验证。2. 快速开始安装 Aerospike extra同时安装所选离线存储的依赖pip install feast[aerospike]你可以从任意标准模板起步如feast init -t local或feast init -t aws然后按下文示例将 online store 切换为 Aerospike。2.1 基础配置 —— 本地 Aerospike CEfeature_store.yaml中最简配置project: my_feature_repo registry: data/registry.db provider: local online_store: type: aerospike hosts: - [127.0.0.1, 3000] namespace: feast2.2 多节点集群配置project: my_feature_repo registry: data/registry.db provider: local online_store: type: aerospike hosts: - [aerospike-1.internal, 3000] - [aerospike-2.internal, 3000] - [aerospike-3.internal, 3000] namespace: feast ttl_seconds: 86400 # 24h 记录级 TTL read_timeout_ms: 150 # 单记录 get 的硬性截止时间 write_timeout_ms: 300 # 单记录 put/operate 的硬性截止时间 batch_total_timeout_ms: 500 # online_read / online_write_batch 的硬性截止时间 batch_max_records: 1000 # batch_write / batch_operate 的分块大小 socket_timeout_ms: 50 # 单次尝试的截止时间使 max_retries 真正生效 max_retries: 2超时语义Timeout semanticsAerospike 客户端区分单次尝试per-attemptsocket_timeout与总截止时间totaltotal_timeout。*_timeout_ms系列参数映射到total_timeout——即包含重试在内的整体调用预算。必须同时设置socket_timeout_ms让每次尝试拥有自己更短的截止时间否则max_retries实际上永远不会触发因为第一次尝试就被允许消耗完整个总截止时间。从源码可以印证这一设计在 _get_client 中read、write、batch三套 policy 都同时写入total_timeout与max_retries且仅在配置了socket_timeout_ms时才将其注入 policyread_policy {total_timeout: store_cfg.read_timeout_ms, max_retries: store_cfg.max_retries} ... if store_cfg.socket_timeout_ms is not None: read_policy[socket_timeout] store_cfg.socket_timeout_ms批量分块Batch chunkingonline_read与online_write_batch会将大请求按batch_max_records默认1000分块。Aerospike 通过服务器端batch-max-requests设置历史上为5000强制每个节点的批处理上限。如果集群上限更严格请调低batch_max_records只有当服务器限制与客户端超时允许时才调高。源码中 _DEFAULT_BATCH_MAX_RECORDS 定义为1_000注释明确说明其目的是保持在服务器batch-max-requests之下避免物化与宽 feature-server 请求触发BatchMaxRequestError错误码 151。读写路径通过 _chunked 生成器切片执行。2.3 Aerospike Enterprise 认证配置需要 Aerospike Enterprise Edition。Community Edition 服务器没有内置用户/安全模型会拒绝这些配置键。project: my_feature_repo registry: data/registry.db provider: local online_store: type: aerospike hosts: - [aerospike.internal, 3000] namespace: feast user: feast_user password: ${AEROSPIKE_PASSWORD} # pragma: allowlist secret auth_mode: internal # internal | external | pkiauth_mode支持三种取值源码 _AUTH_MODE_TO_CONSTANT 将其映射为 Aerospike 客户端常量auth_mode含义internalCE/EE 的用户名/密码认证默认值externalLDAP/Kerberos 等外部认证pki基于证书的认证从源码看当设置了user但未设置password时客户端构造会直接抛出ValueErroruser is set but password is notpassword在配置模型中使用SecretStr类型声明aerospike_repo_configuration.py 对应的配置类并在构造客户端时通过get_secret_value()取用避免明文打印。2.4 Aerospike Enterprise TLS 配置需要 Aerospike Enterprise Edition。Community Edition 服务器未实现 TLS因此tls配置仅对 EE 集群生效。project: my_feature_repo registry: data/registry.db provider: local online_store: type: aerospike hosts: - [aerospike-1.internal, 4333, aerospike-tls] namespace: feast tls: enable: true cafile: /etc/aerospike/certs/ca.pem certfile: /etc/aerospike/certs/client.pem keyfile: /etc/aerospike/certs/client.key注意hosts中每个种子节点的元组形式变为(host, port, tls_name)三元素形式配置模型定义见 AerospikeOnlineStoreConfig.hosts。tls字典会被原样透传给 Aerospike Python 客户端的tlspolicy源码 aerospike.py 中client_config[tls] store_cfg.tls。3. 按 Feature View 的 namespace / set 覆盖两个Dict[str, str]配置字段——namespace_overrides和set_overrides——允许你把个别 feature view 放到不同的 Aerospike namespace 或 set 上而无需把项目拆分到多个存储。凡未在两个映射中列出的 feature view一律回落到存储级默认值namespace/set_name_template。常见的使用场景热点、低延迟的 view 放在纯内存RAM-onlynamespace宽表、冷数据 view 放在 SSD-backed namespace。同一项目不同存储层级。希望feast apply删除或truncate某个 feature view 时是 O(1) 操作、不必扫描其他 view 的记录——为该 view 分配独立 set。配置示例project: my_feature_repo registry: data/registry.db provider: local online_store: type: aerospike hosts: - [aerospike.internal, 3000] namespace: feast # 默认 namespace set_name_template: {project}_{collection_suffix} namespace_overrides: driver_realtime_stats: feast_ram # 内存 namespace driver_history_lookup: feast_ssd # 设备/SSD namespace set_overrides: isolated_view: my_feature_repo_isolated3.1 权衡Tradeoffsnamespace_overrides中列出的每个 namespace 必须已存在于集群——Aerospike 无法在运行时创建 namespace缺失的 namespace 会在第一次读或写时暴露为不透明的AEROSPIKE_ERR_PARAM错误。将 feature view 放到不同 set意味着同一实体的多 feature-view 读取会变成每个 set 一次 Aerospike 往返而不是总共一次往返。只有当下述运维隔离价值大于该成本时才启用。仅触达单个 feature view 的读取不受影响。管理操作会自动遵守覆盖规则update()由feast apply调用将待删除的 feature view 按解析后的(namespace, set)分组每个组发起一次后台扫描background scanteardown()会 truncate 项目可能写入过的每一个唯一(namespace, set)组合含存储级默认值。源码中 _set_name 与 _namespace_for_fv 分别实现 set 与 namespace 的解析逻辑优先查set_overrides/namespace_overrides否则回落默认值。set_name_template支持{project}与{collection_suffix}两个替换变量collection_suffix默认值为latest。update()的分组扫描逻辑见源码 aerospike.py按(ns, set_name)分组后对每个组构造map_remove_by_key操作列表分别移除features与event_ts两个 Map 中该 feature view 的槽位通过client.scan(...).execute_background()以单次服务端后台扫描完成清理。4. Prewriting hooks写前钩子prewriting_hook是一个可调用对象的 import 路径每次online_write_batch调用时被触发一次接收即将写入的行并返回真正落盘的行。它用于那些你不想在每次物化任务里重复粘贴的写侧横切关注点——PII 脱敏、应用侧加密、双写扇出dual-write fan-out、值强转等。Hook 通过 import 字符串而非 PythonCallable值引用这样配置可以在 YAML/JSON 序列化与远程 feature-server 传输中存活。解析后的可调用对象缓存在 store 实例上import 成本每个 store 生命周期只支付一次若配置的 import 字符串在调用间发生变化会在下一次写入时自动重新解析源码见 _resolve_prewriting_hook。4.1 Hook 签名def hook( config: RepoConfig, table: FeatureView, data: list[ tuple[ EntityKeyProto, dict[str, ValueProto], datetime, datetime | None, ] ], ) - list[ tuple[ EntityKeyProto, dict[str, ValueProto], datetime, datetime | None, ] ]: ...Hook必须返回与输入相同 schema 的行列表。返回[]会短路写入——与空输入路径相同不发起任何网络调用。抛出异常的 Hook 会使整个批次失败没有逐行回退机制。源码中PrewritingHook类型别名定义了这一契约aerospike.py并在 online_write_batch 中先解析 hook、再判断空批保证无行调用不支付 import 成本。4.2 实战示例PII 字符串哈希脱敏第 1 步在项目中放置 hook 函数。任何位于每个通过 Feast 写入的进程物化 worker、registry CLI 主机以及 feature server的PYTHONPATH中的模块均可my_feature_repo/hooks.pyPrewriting hooks for the Aerospike online store. from __future__ import annotations import hashlib import os from datetime import datetime from typing import Optional from feast import FeatureView from feast.protos.feast.types.EntityKey_pb2 import EntityKey as EntityKeyProto from feast.protos.feast.types.Value_pb2 import Value as ValueProto from feast.repo_config import RepoConfig # 绝不允许以明文进入在线存储的 feature 名。 # 按精确 feature 名匹配可按项目约定调整。 _SENSITIVE_FEATURES {email, phone_number, ssn} def hash_pii_string_features( config: RepoConfig, table: FeatureView, data: list[ tuple[ EntityKeyProto, dict[str, ValueProto], datetime, Optional[datetime], ] ], ) - list[ tuple[ EntityKeyProto, dict[str, ValueProto], datetime, Optional[datetime], ] ]: 将任何敏感字符串 feature 替换为加盐 SHA-256 十六进制摘要。 哈希是确定性的相同输入 → 相同摘要因此下游以相同方式哈希候选值 的查询仍然能命中。FEAST_PII_SALT 必须在每个物化特征的进程上设置 salt 未设置时直接抛出异常而不是静默回落为明文。 salt os.environ.get(FEAST_PII_SALT) if salt is None: raise RuntimeError( FEAST_PII_SALT is not set; refusing to write feature batches without a configured PII salt. ) salt_bytes salt.encode(utf-8) def _digest(plaintext: str) - str: h hashlib.sha256() h.update(salt_bytes) h.update(plaintext.encode(utf-8)) return h.hexdigest() transformed: list[ tuple[ EntityKeyProto, dict[str, ValueProto], datetime, Optional[datetime], ] ] [] for entity_key, values, event_ts, created_ts in data: new_values dict(values) for feature_name in _SENSITIVE_FEATURES.intersection(new_values): v new_values[feature_name] if v.HasField(string_val) and v.string_val: new_values[feature_name] ValueProto(string_val_digest(v.string_val)) transformed.append((entity_key, new_values, event_ts, created_ts)) return transformed第 2 步在feature_store.yaml中引用 hookproject: my_feature_repo registry: data/registry.db provider: local online_store: type: aerospike hosts: - [aerospike.internal, 3000] namespace: feast prewriting_hook: my_feature_repo.hooks.hash_pii_string_features4.3 运维注意事项Hook 只在写路径被调用读路径不受影响地直通 store。如果你的 hook 是单向的如哈希你必须在读取时对候选值自行应用同样的变换。Hook 与写入者运行在同一进程内——它们不是 RPC也不在沙箱中。它们可以读取环境变量、打开文件、调用 KMS 等。请将它们视为可信代码库的一部分。配置错误的prewriting_hookimport 路径错误、函数缺失、目标不可调用会在第一次online_write_batch调用时抛出ValueError/TypeError而不是在 store 构造时。建议在部署时添加一个写单行的冒烟测试让错误配置在真实批处理前暴露。源码中_resolve_prewriting_hook对 import 失败、属性缺失、非可调用对象分别抛出带具体消息的ValueError/TypeError。完整配置选项可查看AerospikeOnlineStoreConfig源码类定义见 aerospike.py。5. 数据模型Aerospike online store 采用**每个项目一个 set 实体键共置entity-key collocation**的布局。同一实体的多个 feature view 的特征值存储在同一条 Aerospike 记录上类似于 MongoDB online store 的每个实体一个文档布局。Aerospike 概念Feast 映射Namespaceonline_store.namespace必须已在集群上预配置通过online_store.namespace_overrides按 feature view 覆盖Setonline_store.set_name_template→ 默认{project}_{collection_suffix}通过online_store.set_overrides按 feature view 覆盖Keyserialize_entity_key(entity_key)作为bytearray用户键BinfeaturesMap CDT键为 feature view 名每个值为feature → native映射Binevent_tsMap CDT键为 feature view 名每个值为 int64 毫秒级时间戳Bincreated_ts顶层 int64 毫秒级时间戳最近一次feast materialize5.1 示例记录对于单个实体、携带两个 feature viewdriver_stats和pricing的特征key: (nsfeast, setmy_feature_repo_latest, user_keyserialize_entity_key as bytearray) bins: features: driver_stats: rating: 4.91 trips_last_7d: 132 pricing: surge_multiplier: 1.2 event_ts: driver_stats: 1737374400000 # 2025-01-20T12:00:00Z pricing: 1737447000000 # 2025-01-21T08:30:00Z created_ts: 1737460805000 # 2025-01-21T12:00:05Z5.2 关键设计决策每条记录对应一个实体每个 bin 对应一个概念。features和event_ts是 Aerospike Map CDT bin而不是动态 bindynamic bins这使存储保持在 Aerospike 15 字节 bin 名限制之内无论项目拥有多少个 feature view。通过 Map CDT 操作实现部分 upsert。写入使用batch_writemap_put_items(features, {fv: {...}})与map_put(event_ts, fv, epoch_ms)。同一实体上不同 feature view 的并发写入永远不会互相覆盖——每次写入只修改自己所在 map 的键。源码 _build_batch_writes 展示了这一操作列表的完整构造且 Map 以MAP_KEY_ORDERED策略创建使map_get_by_key/map_remove_by_key在 map 规模上保持 O(log N)见 _ORDERED_MAP_POLICY。实体键字节作为 Aerospike 用户键。Feast 的serialize_entity_key输出作为bytearray用户键而非bytes传入——Python 客户端对bytes键只哈希第一个字节会把不同实体折叠在一起源码 _aerospike_key 注释明确记录了这一点。时间戳以 int64 毫秒级存储。Aerospike 没有原生 datetime 类型tz-naive 时间戳按OnlineStore契约视为 UTC。转换函数 _datetime_to_epoch_ms 与 _epoch_ms_to_datetime 在单元测试 test_aerospike_online_retrieval.py 中有往返round-trip与 naive 视为 UTC 的验证。5.3 TTL 与过期ttl_seconds在每次online_write_batch调用时作为记录级元数据写入源码 _resolve_ttlttl_secondsAerospike TTL效果未设置 /nullTTL_NAMESPACE_DEFAULT记录继承 namespace 配置的default-ttl。0TTL_NEVER_EXPIRE记录一直保留直到被显式删除。0对应秒数记录由服务器nsup线程驱逐。当前版本没有按 feature view 的 TTL 覆盖——该设置对 online store 发起的每次写入统一生效。5.4 索引不创建任何二级索引。所有访问都通过主键即序列化后的实体键进行。6. 异步支持异步读写通过将 Aerospike Python 客户端的阻塞调用放到默认线程池执行器loop.run_in_executor实现。底层 C 客户端在网络 I/O 期间会释放 GIL因此await store.online_read_async(...)能保持事件循环响应。当前未使用原生 asyncio Aerospike 客户端。源码 _online_write_batch_async / _online_read_async 使用functools.partial包装同步方法后提交给执行器。同步与异步方法均得到完整支持online_read/online_read_asynconline_write_batch/online_write_batch_asyncinitialize/close——initialize(config)会急切地打开连接让 feature server 在启动时支付 TCP/握手成本close()释放缓存的客户端。源码 initialize / close 同样通过run_in_executor执行且close()以_client_lock保护避免并发关闭竞态。客户端创建本身是惰性 加锁的_get_client 在首次使用时创建并缓存单例客户端双重检查加锁double-checked locking防止线程化 feature server 中并发首调泄漏多余连接。7. 功能矩阵在线存储支持的功能集合在 overview.md#functionality 中有详细描述。以下矩阵展示 Aerospike online store 支持的功能功能Aerospike向在线存储写入特征值yes从在线存储读取特征值yes更新在线存储基础设施如表yes拆除在线存储基础设施如表yes生成基础设施变更计划no支持按需转换on-demand transformsyes可被 Python SDK 读取yes可被 Java 读取no可被 Go 读取no支持无实体 feature viewyes支持并发写入同一键yes支持检索时 TTLyes支持删除过期数据yes按 feature view 共置collocatedno按 feature service 共置no按实体键共置yes要与其他在线存储对比该功能集合请参见完整的功能矩阵。8. 测试与验证仓库为 Aerospike online store 提供了两层测试保障可作为理解实现行为的辅助材料单元测试test_aerospike_online_retrieval.py覆盖时间戳/TTL 辅助函数、面向列的 proto 重塑reshape以及使用 mock 客户端调度的写/读/管理路径另有一个标注_requires_docker的端到端测试在 Docker 不可用时跳过。通用在线存储测试aerospike.py测试仓库配置使用aerospike/aerospike-server:8.0.0.10_1社区版镜像通过 testcontainers 拉起单节点集群等待日志标记migrations: complete后复用其内置的testnamespace集成配置定义在 aerospike_repo_configuration.py 中。9. 小结Aerospike online store 以单 set 每项目 实体键共置 Map CDT 部分 upsert为核心设计在保留 Aerospike 高性能读写的同時提供了 TTL 管理、按 feature view 的 namespace/set 隔离、prewriting hook 扩展与完整同步/异步 API。接入时需重点留意namespace 必须预先在集群上创建、超时参数需同时配置socket_timeout_ms才能让重试生效、批量请求需低于服务器batch-max-requests上限。本文所有配置与行为均可在 AerospikeOnlineStoreConfig 源码 与配套测试中找到直接依据。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考