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 CE
feature_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 semantics)
Aerospike 客户端区分单次尝试(per-attempt,socket_timeout)与总截止时间(total,total_timeout)。*_timeout_ms系列参数映射到total_timeout——即包含重试在内的整体调用预算。必须同时设置socket_timeout_ms,让每次尝试拥有自己更短的截止时间;否则max_retries实际上永远不会触发,因为第一次尝试就被允许消耗完整个总截止时间。
从源码可以印证这一设计:在 _get_client 中,read、write、batch三套 policy 都同时写入total_timeout与max_retries,且仅在配置了socket_timeout_ms时才将其注入 policy:
read_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 chunking)
online_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 | 含义 |
|---|---|
internal | CE/EE 的用户名/密码认证(默认值) |
external | LDAP/Kerberos 等外部认证 |
pki | 基于证书的认证 |
从源码看,当设置了user但未设置password时,客户端构造会直接抛出ValueError("user is set but password is not");password在配置模型中使用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-only)namespace;宽表、冷数据 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 权衡(Tradeoffs)
namespace_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 scan);teardown()会 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.py:
"""Prewriting 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中引用 hook:
project: 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_hook(import 路径错误、函数缺失、目标不可调用)会在第一次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 映射 |
|---|---|
| Namespace | online_store.namespace(必须已在集群上预配置);通过online_store.namespace_overrides按 feature view 覆盖 |
| Set | online_store.set_name_template→ 默认"{project}_{collection_suffix}";通过online_store.set_overrides按 feature view 覆盖 |
| Key | serialize_entity_key(entity_key)作为bytearray用户键 |
Binfeatures | Map CDT,键为 feature view 名,每个值为feature → native映射 |
Binevent_ts | Map CDT,键为 feature view 名,每个值为 int64 毫秒级时间戳 |
Bincreated_ts | 顶层 int64 毫秒级时间戳(最近一次feast materialize) |
5.1 示例记录
对于单个实体、携带两个 feature view(driver_stats和pricing)的特征:
key: (ns="feast", set="my_feature_repo_latest", user_key=<serialize_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,而不是动态 bin(dynamic bins),这使存储保持在 Aerospike 15 字节 bin 名限制之内,无论项目拥有多少个 feature view。 - 通过 Map CDT 操作实现部分 upsert。写入使用
batch_write+map_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_ttl):
ttl_seconds | Aerospike TTL | 效果 |
|---|---|---|
未设置 /null | TTL_NAMESPACE_DEFAULT | 记录继承 namespace 配置的default-ttl。 |
0 | TTL_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 transforms) | yes |
| 可被 Python SDK 读取 | yes |
| 可被 Java 读取 | no |
| 可被 Go 读取 | no |
| 支持无实体 feature view | yes |
| 支持并发写入同一键 | yes |
| 支持检索时 TTL | yes |
| 支持删除过期数据 | yes |
| 按 feature view 共置(collocated) | no |
| 按 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),仅供参考