Feast 接入 Hazelcast 在线存储(Online Store)实战指南:配置、连接模式与实现原理
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
Feast 的 Hazelcast Online Store 是社区贡献的在线存储实现,让用户可以把 Hazelcast 为核心,结合源码实现(hazelcast_online_store.py)与官方模板(hazelcast 模板),系统讲解如何配置feature_store.yaml接入 Hazelcast、如何理解 TTL 过期机制、三种集群连接模式的差异,以及读写/清理背后的数据模型与调用链。读完本文,你将能够独立完成一个以 Hazelcast 为在线存储的 Feast 特征仓库的初始化、配置、物化与在线读取。
一、Hazelcast Online Store 是什么
Hazelcast Online Store 是 Feast 的一个在线存储贡献实现。与 Redis、DynamoDB 等其他在线存储一样,它在 Feast 的OnlineStore抽象之下工作:一旦在feature_store.yaml中给出了 Hazelcast 客户端配置,其余一切——schema 创建、向 Hazelcast 读写数据、删除(remove)操作——都会像其他在线存储一样由 Feast 自动处理,用户无需关心底层细节。
在仓库中,这个贡献由三个文件组成:
- hazelcast_online_store.py:核心实现,包含配置模型
HazelcastOnlineStoreConfig与存储实现HazelcastOnlineStore; - hazelcast_repo_configuration.py:面向通用集成测试的仓库配置;
- init.py:包标记文件。
使用前需要有一个正在运行的 Hazelcast 集群。你可以通过 Hazelcast Viridian Serverless 快速创建一个云集群,也可以在自己的本地或远程机器上部署一个集群(容器方式可直接使用官方镜像hazelcast/hazelcast,默认客户端端口为 5701)。该在线存储同时支持连接本地/远程集群与 Hazelcast Viridian Serverless 集群两种方式。
二、快速开始:创建 Feature Repository
接入 Hazelcast 在线存储的最快方式,是使用 Feast CLI 初始化一个新的特征仓库。在 Feast 安装完成后,执行:
feast init FEATURE_STORE_NAME -t hazelcast该命令会以交互方式引导你填写访问 Hazelcast 集群所需的配置细节,并生成带online_store配置的feature_store.yaml。交互引导逻辑实现在 templates/hazelcast/bootstrap.py 的collect_hazelcast_online_store_settings()中,主要询问以下内容:
| 交互问题 | 对应配置项 | 默认值 |
|---|---|---|
连接本地集群[L]还是 Viridian 集群[V]? | 决定连接模式 | L |
| Viridian 模式:Cluster ID | cluster_name | 无 |
| Viridian 模式:Discovery Token | discovery_token | 无 |
| Viridian 模式:CA / CERT / Key 文件路径 | ssl_cafile_path/ssl_certfile_path/ssl_keyfile_path | 无 |
| 本地模式:Cluster name | cluster_name | dev |
| 本地模式:Cluster members | cluster_members | localhost:5701 |
| 本地模式:是否启用 TLS/SSL | 启用后填写 CA / CERT / Key 路径 | 否 |
| Key TTL seconds | key_ttl_seconds | 0 |
其中 Viridian 模式下要求提供证书三件套;本地模式仅在确认启用 TLS/SSL 后才询问证书路径。若不需要 SSL,引导脚本还会自动移除模板中的ssl_password: ${SSL_PASSWORD}行(见bootstrap.py中apply_hazelcast_store_settings()对remove_lines_from_file的调用)。
替代方案:也可以先按 Feast 官方快速入门(见 getting-started/quickstart.md 与 examples/quickstart/quickstart.ipynb)执行feast init -t FEATURE_STORE_NAME,然后手动编辑feature_store.yaml中的online_store段落。除第 2 步(配置 Hazelcast 在线存储)外,后续所有步骤——特征定义编写、feast apply部署、训练数据生成、物化(materialization)、在线/离线特征获取——都与 Feast 通用快速入门完全一致。
初始化成功后,特征仓库模板目录位于 templates/hazelcast/feature_repo,包含:
feature_store.yaml:仓库配置(初始占位符会被交互引导替换);feature_definitions.py:示例特征定义(driver 实体、driver_hourly_stats特征视图、on-demand 特征视图、FeatureService、PushSource 等);test_workflow.py:一条完整的端到端演示脚本(apply → 历史特征 → 物化 → 在线特征 → push → teardown);data/driver_stats.parquet:由 bootstrap 生成的演示数据。
三、feature_store.yaml 配置详解
3.1 配置参数总览
online_store段落中所有可用参数均定义在 hazelcast_online_store.py 的HazelcastOnlineStoreConfig中,其字段、类型与默认值如下:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
type | string | "hazelcast" | 在线存储类型选择器,必须为hazelcast |
cluster_name | string | "dev" | 要连接的集群名称,Hazelcast 默认集群名为dev |
cluster_members | string 列表 | ["localhost:5701"] | 与集群相连的成员地址列表(本地/远程模式使用) |
discovery_token | string | "" | Hazelcast Viridian 集群的发现令牌(Viridian 模式使用) |
ssl_cafile_path | string | "" | PEM 格式的 CA 证书绝对路径 |
ssl_certfile_path | string | "" | PEM 格式的客户端证书绝对路径 |
ssl_keyfile_path | string | "" | PEM 格式的客户端私钥文件绝对路径 |
ssl_password | string | "" | 若私钥文件加密,用于解密的密码 |
key_ttl_seconds | int | 0 | Hazelcast 键桶 TTL(秒),用于实体过期 |
3.2 连接本地 / 远程集群(TLS 启用示例)
以下示例连接一个名为dev、运行在5701端口且启用了 TLS/SSL 的本地集群:
[...] online_store: type: hazelcast cluster_name: dev cluster_members: ["localhost:5701"] ssl_cafile_path: /path/to/ca/file ssl_certfile_path: /path/to/cert/file ssl_keyfile_path: /path/to/key/file ssl_password: ${SSL_PASSWORD} # The password will be read form the `SSL_PASSWORD` environment variable. key_ttl_seconds: 86400 # The default is 0 and means infinite.3.3 连接 Hazelcast Viridian 集群
若要连接 Hazelcast Viridian 云集群而非本地/远程集群,将cluster_members替换为cluster_id与discovery_token:
[...] online_store: type: hazelcast cluster_name: YOUR_CLUSTER_ID discovery_token: YOUR_DISCOVERY_TOKEN ssl_cafile_path: /path/to/ca/file ssl_certfile_path: /path/to/cert/file ssl_keyfile_path: /path/to/key/file ssl_password: ${SSL_PASSWORD} # The password will be read form the `SSL_PASSWORD` environment variable. key_ttl_seconds: 86400 # The default is 0 and means infinite.注意两点:
ssl_password支持${SSL_PASSWORD}这样的环境变量引用形式,密码会从同名环境变量读取,避免明文落盘;- 模板 templates/hazelcast/feature_repo/feature_store.yaml 中还额外显式设置了
entity_key_serialization_version: 3,这与实现中固定使用实体键序列化版本 3 保持一致。
3.4 三种连接模式的选择逻辑
HazelcastOnlineStore在_get_client()(hazelcast_online_store.py)中按优先级决定连接方式,理解这个逻辑有助于判断配置是否生效:
- Viridian 模式:若
discovery_token != "",则将云发现地址设为api.viridian.hazelcast.com,并使用cluster_name+cloud_discovery_token+ SSL 证书参数建立客户端; - 本地/远程 TLS 模式:若
ssl_cafile_path != "",则使用cluster_members(客户端侧无需显式传入,但成员地址仍由配置提供)与 SSL 三件套建立客户端; - 本地/远程明文模式:以上均不满足时,仅使用
cluster_members与cluster_name建立普通客户端。
三种模式都会开启statistics_enabled=True。同时,客户端连接是**惰性创建 + 双重检查锁(double-checked locking)**的单例模式:_client类属性在首次调用时通过threading.Lock()保护创建,之后所有读写复用同一连接。
四、TTL 配置:特征的自动过期与淘汰
TTL(Time-To-Live)是指每个特征在 map 中保持空闲的最大秒数。它限制特征相对于其最后一次读或写访问时间的存活时长;当某特征的空闲时长超过该上限时,它会被自动过期并淘汰。一个特征是"空闲"的,指的是没有任何get或put操作作用于它。
[...] online_store: [...] key_ttl_seconds: 86400参数约束与行为:
- 取值范围:合法的
key_ttl_seconds是0到Integer.MAX_VALUE之间的整数; - 默认值:
0,含义是永不过期(infinite); - 作用粒度:该 TTL 在写入时通过
IMap.put(key, value, ttl)的第三个参数传入(见 online_write_batch),即每个键(特征条目)独立计时,由 Hazelcast 服务器端负责过期淘汰。
从源码结构看,TTL 是"写入即生效"的:每次物化写入都会用当前配置的 TTL 刷新对应键的过期时间,因此在线特征被持续更新时,其"最近访问时间"会随之顺延,不会被错误淘汰。
五、源码级原理:数据模型与读写调用链
5.1 存储模型:一个 FeatureView 对应一个 IMap
HazelcastOnlineStore将数据存储在 Hazelcast 的分布式IMap中,每个 map 的命名规则为{project}_{feature_view_name}(见_map_name(),hazelcast_online_store.py)。例如项目my_project下的特征视图driver_hourly_stats对应 mapmy_project_driver_hourly_stats。
Map 的 Key由base64(entity_key 序列化结果) + feature_name拼接而成,其中实体键使用序列化版本 3(serialize_entity_key(entity_key, entity_key_serialization_version=3))。因此同一个实体的不同特征在 map 中是相互独立的条目。
Map 的 Value以HazelcastJsonValue(JSON)存储,包含 5 个字段(常量定义见 hazelcast_online_store.py):
| 字段 | 含义 |
|---|---|
entity_key | base64 编码后的实体键字符串 |
feature_name | 特征名 |
feature_value | 特征值:ValueProto.SerializeToString()后再 base64 编码 |
event_ts | 事件时间(UTC 时间戳,秒级浮点数) |
created_ts | 创建时间(UTC 时间戳;None时为0.0) |
5.2 写入:materialize 的落库路径
online_write_batch()(hazelcast_online_store.py)是物化数据写入的入口,流程如下:
- 校验
config.online_store必须是HazelcastOnlineStoreConfig,否则抛出HazelcastInvalidConfig; - 通过
_get_client()获取客户端并client.get_map(...)取得目标 map; - 对每条记录:序列化实体键 → 将
event_ts/created_ts转为 UTC 时间戳 → 序列化每个特征值; - 以
entity_key_str + feature_name为 key,HazelcastJsonValue为 value,调用fv_map.put(key, value, key_ttl_seconds)写入(第三个参数即 TTL); - 每写一条调用一次
progress(1)回调,供物化进度上报使用。
5.3 读取:online_read 的在线获取路径
online_read()(hazelcast_online_store.py)实现在线特征读取:
- 对每个请求的实体键做同样的 base64 序列化;
- 若指定了
requested_features,则构造entity_key_str + feature的键列表;否则默认取该 FeatureView 的全部特征(table.features); - 通过
fv_map.get_all(hz_keys)一次性批量取回全部条目,减少网络往返; - 对每条命中记录:
loads()解析 JSON、base64.b64decode反序列化ValueProto,并以event_ts构造返回的时间戳; - 未命中的实体返回
(None, None),命中则返回(event_ts, {feature_name: value_proto})。
5.4 Schema 管理:update 与 teardown
Feast 的feast apply会触发update()(hazelcast_online_store.py),它通过 Hazelcast SQL 为保留的特征视图创建映射(mapping),使 IMap 具有可查询的 schema:
CREATE OR REPLACE MAPPING {project}_{table.name} ( __key VARCHAR, entity_key VARCHAR, feature_name VARCHAR, feature_value VARCHAR, event_ts DECIMAL, created_ts DECIMAL ) TYPE IMap OPTIONS ( 'keyFormat' = 'varchar', 'valueFormat' = 'json-flat' )其中__key对应 map 的键(即entity_key_str + feature_name),其余字段与 value JSON 一一对应。对需要删除的表,则执行DELETE FROM清理数据并DROP MAPPING IF EXISTS删除映射。teardown()(hazelcast_online_store.py)的行为与 update 中的删除逻辑一致:先DELETE FROM再DROP MAPPING IF EXISTS,实现feast teardown的完整清理。
六、接入后的完整工作流(以模板 demo 为例)
初始化完成并确认 Hazelcast 集群可达后,即可按标准 Feast 流程操作。仓库模板自带的 test_workflow.py 给出了从应用到清理的完整链路,可直接作为验证 Hazelcast 接入是否成功的冒烟脚本:
- 注册特征仓库:
feast apply(或FeatureStore(repo_path=".")后调用store.apply),触发update()在 Hazelcast 中为各 FeatureView 创建 map 映射; - 离线历史特征:
store.get_historical_features(entity_df=..., features=[...]).to_df()生成训练数据(feature_definitions.py中定义了driver_hourly_stats特征视图与transformed_conv_rateon-demand 特征视图,模板同时演示了请求特征val_to_add的参与); - 物化到在线存储:
store.materialize_incremental(end_date=datetime.now()),将离线数据按上述online_write_batch路径写入 Hazelcast map; - 在线特征读取:
store.get_online_features(features=..., entity_rows=[{"driver_id": 1001}, ...]).to_dict(); - 在线读取进阶用法:可通过
FeatureService(如模板中的driver_activity_v1)整体获取;也可用store.push("driver_stats_push_source", event_df, to=PushMode.ONLINE_AND_OFFLINE)将流式事件直接推入在线存储,验证推源(PushSource)场景下的新鲜特征读取; - 清理:
feast teardown,触发teardown()删除数据与映射。
一个实用的本地验证方式:仓库的通用集成测试 tests/universal/feature_repos/universal/online_store/hazelcast.py 展示了如何用 Docker 启动一个真实的 Hazelcast 集群供本地联调——它使用hazelcast/hazelcast容器,设置环境变量HZ_CLUSTERNAME(随机 5 位小写字母)与HZ_NETWORK_PORT_AUTOINCREMENT=true,暴露 5701 端口,并等待日志中出现Cluster name: {cluster_name}即视为就绪。如果你本地有 Docker,完全可以照此方式起一个干净集群,再用生成的cluster_members(宿主机 IP + 映射端口)填入自己的feature_store.yaml。对应的集成测试配置在 hazelcast_repo_configuration.py 中注册,可参与 Feast 通用在线存储测试套件。
七、使用建议与注意事项
- 连接模式选择:本地联调用
cluster_members明文模式即可;生产环境若使用 Viridian,务必配置discovery_token与证书三件套,并通过环境变量注入ssl_password; - TTL 合理取值:在线特征需要"永不过期"时保持默认
0;若特征存在时效性(如按天失效),可仿照示例设置86400(24 小时)。TTL 由 Hazelcast 服务端按空闲时间淘汰,无需 Feast 侧额外清理; - 实体键序列化:模板与实现统一使用实体键序列化版本 3,自行手写
feature_store.yaml时建议显式声明entity_key_serialization_version: 3,避免与其他存储混用或升级时的不一致; - 验证路径:接入后建议先跑一遍模板自带的
test_workflow.py覆盖 apply → 物化 → 在线读取 → teardown 全链路,再接入自有特征;更深入的 Hazelcast 能力可查阅其官方文档(如 SQL 映射、IMap 过期策略等)与仓库内 reference/online-stores 下其他在线存储的对比说明。
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考