Feast 远程离线存储:用 Arrow Flight(gRPC)把离线存储服务化,实现跨网络的历史特征检索
Feast 远程离线存储用 Arrow FlightgRPC把离线存储服务化实现跨网络的历史特征检索【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast本篇基于 Feast 仓库中的examples/remote-offline-store官方示例讲解如何将离线存储Offline Store部署为独立的 Arrow Flight 服务端并在客户端以type: remote的离线存储配置通过网络完成训练数据集点时连接特征的拉取。读完后你将掌握feast serve_offline服务端的启动与配置、客户端feature_store.yaml的完整参数含义以及RemoteOfflineStore在源码层面如何通过 gRPC 的do_put/do_get协议将取数请求委托给远端服务。一、为什么需要 Remote Offline StoreFeast 的离线存储BigQuery、Snowflake、DuckDB 等通常运行在训练环境所在的进程内FeatureStore.get_historical_features()会直接在本地通过对应离线存储的实现执行点时连接point-in-time join。但在生产拓扑中训练任务往往无法或不应该直连离线数仓例如网络隔离、凭据收敛、多团队共享同一套特征资产等场景。Feast 为此提供了 Remote Offline Store它本质上是一个Apache Arrow Flight 服务端 客户端。服务端Offline feature server在拥有真实离线数据的地方运行把OfflineStore接口的能力暴露为 Arrow FlightgRPC端点客户端只需在feature_store.yaml中把offline_store.type设为remote并填写host/port即可让所有离线取数请求——包括get_historical_features、write_logged_features、offline_write_batch等——被委托到远端执行。官方参考文档见 remote-offline-store.md服务端文档见 offline-feature-server.md。整个示例的目录结构为offline_server一个完整的示例 Feast 仓库包含特征定义与本地数据作为远程服务端的部署对象offline_client一个极简客户端其feature_store.yaml使用remote类型离线存储通过 test.py 验证历史特征拉取。二、服务端准备一个可被“远程化”的 Feast 项目2.1 服务端 feature_store.yaml示例服务端仓库 feature_store.yaml 的配置如下project: offline_server # By default, the registry is a file (but can be turned into a more scalable SQL-backed registry) registry: data/registry.db # The provider primarily specifies default offline / online stores storing the registry in a given cloud provider: local online_store: type: sqlite path: data/online_store.db entity_key_serialization_version: 3要点说明project: offline_server项目名客户端的同一份仓库配置必须与之保持一致本例中客户端配置即声明project: offline_serverregistry: data/registry.db默认使用本地文件型 registry注释中也提示可替换为更可扩展的 SQL 后端 registryprovider: local本地 Provider离线/在线存储默认使用文件与 SQLiteentity_key_serialization_version: 3实体键序列化版本客户端与服务端应保持一致否则反序列化实体键可能出错仓库中专门有 entity-reserialization-of-from-v2-to-v3.md 讲解 v2→v3 迁移注意该示例中没有显式声明offline_store在provider: local下离线数据由文件源parquet直接支撑服务端的离线取数逻辑即围绕这些数据源展开。2.2 特征定义 example_repo.pyexample_repo.py 定义了客户端将要跨网络拉取的特征核心结构如下# 实体driver_id 作为主键 driver Entity(namedriver, join_keys[driver_id]) # 从 parquet 文件读取数据 driver_stats_source FileSource( namedriver_hourly_stats_source, pathf{os.path.dirname(os.path.abspath(__file__))}/data/driver_stats.parquet, timestamp_fieldevent_timestamp, created_timestamp_columncreated, ) # 特征视图3 个特征字段开启在线存储 driver_stats_fv FeatureView( namedriver_hourly_stats, entities[driver], ttltimedelta(days1), schema[ Field(nameconv_rate, dtypeFloat32), Field(nameacc_rate, dtypeFloat32), Field(nameavg_daily_trips, dtypeInt64, descriptionAverage daily trips), ], onlineTrue, sourcedriver_stats_source, tags{team: driver_performance}, )此外还定义了RequestSourcevals_to_add字段val_to_add/val_to_add_2只存在于请求时刻的输入数据on_demand_feature_view装饰的transformed_conv_rate在conv_rate基础上加上请求字段产出conv_rate_plus_val1、conv_rate_plus_val2基于PushSource的driver_hourly_stats_fresh视图及driver_activity_v1/v2/v3三个 FeatureService。客户端后面要跨网络拉取的正是driver_hourly_stats的三个基础特征与transformed_conv_rate的两个按需变换特征。数据文件为 driver_stats.parquet在线数据落在data/online_store.dbSQLite。三、启动远程离线服务端3.1 应用特征仓库并启动服务在服务端目录下先执行 apply 注册特征与数据源再启动离线服务引自 README 的原始步骤cd offline_server feast -c feature_repo applyfeast -c feature_repo serve_offline启动成功的样例输出Serving on grpctcp://127.0.0.1:8815feast serve_offline会拉起一个 Arrow Flight 服务默认监听127.0.0.1:8815。从源码看该命令定义在 serve.py可用参数为参数缩写默认值说明--host-h127.0.0.1服务监听主机--port-p8815服务端口常量定义于 constants.pyDEFAULT_OFFLINE_SERVER_PORT 8815--key-k空TLS 私钥证书路径需与--cert同时提供才能以 TLS 模式启动--cert-c空TLS 公钥证书路径只传其中一个会报BadParameterserve_offline_command最终调用store.serve_offline(host, port, tls_key_path, tls_cert_path)该方法在 feature_store.py 中转发到offline_server.start_server(...)完成 gRPC 服务注册与启动。3.2 以环境变量注入 feature_store.yaml容器化部署示例 README 还描述了服务端的另一种初始化方式通过名为FEATURE_STORE_YAML_BASE64的环境变量提供feature_store.yaml文件Base64 编码。服务端会创建一个临时目录并把该 YAML 解包为其中的feature_store.yml再加载。该环境变量名在源码中的常量为 constants.py 的FEATURE_STORE_YAML_ENV_NAME FEATURE_STORE_YAML_BASE64。这种方式适合把服务端打成镜像后以 K8s Pod 运行无需在镜像内挂载仓库目录。四、客户端配置offline_store type: remote4.1 客户端 feature_store.yaml客户端仓库 feature_store.yaml 的完整内容为project: offline_server # By default, the registry is a file (but can be turned into a more scalable SQL-backed registry) registry: ../offline_server/feature_repo/data/registry.db # The provider primarily specifies default offline / online stores storing the registry in a given cloud provider: local offline_store: type: remote host: localhost port: 8815 entity_key_serialization_version: 3关键设计project与服务端一致指向同一个offline_server项目的元数据registry直接复用服务端 apply 后生成的本地 registry 文件../offline_server/feature_repo/data/registry.db。也就是说registry 是客户端本地可见的而真正的数据查询被委托给远端——这是理解 Remote Offline Store 工作方式的核心客户端持有“元数据”服务端持有“数据”offline_store段声明委托配置type: remote且给出host与port。4.2 配置项全解RemoteOfflineStoreConfigtype: remote对应的 Pydantic 配置模型是 remote.py 中的RemoteOfflineStoreConfig可配置字段比示例中用到的更多class RemoteOfflineStoreConfig(FeastConfigBaseModel): type: Literal[remote] remote scheme: Literal[http, https] http # https 时以 grpctls 连接 host: StrictStr # 必填Arrow Flight 服务端地址 port: Optional[StrictInt] None # 服务端端口 cert: StrictStr 服务端以 TLS 模式启动例如自签名证书时客户端需要指向公钥证书 文件通常以 .crt/.cer/.pem 结尾。 connection_retries: int Field(default3, ge0) 针对瞬时 Arrow Flight 错误的重试次数指数退避默认 3。各字段的作用结合 build_arrow_flight_client 的实现host/port拼出 Arrow Flight 连接串。默认scheme: http对应grpctcp://host:port当scheme: https时切换为grpctls即服务端以 TLS 启动时客户端必须同步声明https并提供certcert以二进制读取后作为tls_root_certs传入 Flight 客户端用于信任自签证书connection_retries默认 3允许 0客户端类FeastFlightClient继承pyarrow.flight.FlightClient并叠加arrow_client_error_handling_decorator对get_flight_info、do_get、do_put等调用做错误处理与重试包装见 remote.py。五、客户端执行历史特征拉取test.py 构造了一个包含driver_id、event_timestamp、label_driver_reported_satisfaction、val_to_add、val_to_add_2的实体 DataFrame然后像使用任何本地 Feast 客户端一样拉取历史特征from datetime import datetime from feast import FeatureStore import pandas as pd entity_df pd.DataFrame.from_dict( { driver_id: [1001, 1002, 1003], event_timestamp: [ datetime(2021, 4, 12, 10, 59, 42), datetime(2021, 4, 12, 8, 12, 10), datetime(2021, 4, 12, 16, 40, 26), ], label_driver_reported_satisfaction: [1, 5, 3], val_to_add: [1, 2, 3], val_to_add_2: [10, 20, 30], } ) features [ driver_hourly_stats:conv_rate, driver_hourly_stats:acc_rate, driver_hourly_stats:avg_daily_trips, transformed_conv_rate:conv_rate_plus_val1, transformed_conv_rate:conv_rate_plus_val2, ] store FeatureStore(repo_path.) training_df store.get_historical_features(entity_df, features).to_df()运行方式README 原始步骤cd offline_client python test.py样例输出节选完整输出见 READMEconfig.offline_store is class feast.infra.offline_stores.remote.RemoteOfflineStoreConfig ----- Feature schema ----- class pandas.core.frame.DataFrame RangeIndex: 3 entries, 0 to 2 Data columns (total 10 columns): # Column Non-Null Count Dtype --- ------ -------------- ----- 0 driver_id 3 non-null int64 1 event_timestamp 3 non-null datetime64[ns, UTC] 2 label_driver_reported_satisfaction 3 non-null int64 3 val_to_add 3 non-null int64 4 val_to_add_2 3 non-null int64 5 conv_rate 3 non-null float32 6 acc_rate 3 non-null float32 7 avg_daily_trips 3 non-null int32 8 conv_rate_plus_val1 3 non-null float64 9 conv_rate_plus_val2 3 non-null float64 dtypes: datetime64ns, UTC, float32(2), float64(2), int32(1), int64(4) memory usage: 332.0 bytes None ----- Features ----- driver_id event_timestamp label_driver_reported_satisfaction ... avg_daily_trips conv_rate_plus_val1 conv_rate_plus_val2 0 1001 2021-04-12 10:59:4200:00 1 ... 590 1.022378 10.022378 1 1002 2021-04-12 08:12:1000:00 5 ... 974 2.762213 20.762213 2 1003 2021-04-12 16:40:2600:00 3 ... 127 3.419828 30.419828 [3 rows x 10 columns]第一行输出config.offline_store is class feast.infra.offline_stores.remote.RemoteOfflineStoreConfig证明了客户端确实在使用 remote 类型的离线存储配置返回的 10 列中既有点时连接得到的driver_hourly_stats基础特征也包含transformed_conv_rate按需特征视图的计算结果与实体表字段共同构成可直接用于模型训练的训练集。六、源码解析RemoteOfflineStore 如何把请求委托到远端客户端类的实现在 remote.py它实现了标准OfflineStore接口的远端版本。从源码结构看其通信协议是一个统一的“两步式”命令模型第一步_call_put上传命令描述与输入数据。_call_put 为每次调用生成一个 UUID 作为command_id把api远端应执行的方法名与参数打包成 JSON 构造FlightDescriptor.for_command(...)随后通过 _put_parameters 用client.do_put上传 Arrow 数据entity_dfpandas DataFrame经pa.Table.from_pandas转换、table已是 Arrow 表的数据或在两者都缺省时上传一个仅含key列的占位表。第二步_call_get取回结果。_call_get 先用client.get_flight_info(command_descriptor)拿到 flight 信息与 ticket再client.do_get(ticket)读取流式结果并read_all()为 Arrow 表。检索类方法经由 _send_retrieve_remote 串联这两步。在此协议上RemoteOfflineStore覆盖了以下能力与 remote-offline-store.md 文档描述的能力清单一致实际覆盖方法更多方法远端执行的语义委托方式get_historical_features点时连接构建训练集返回RemoteRetrievalJobto_df()/to_arrow()时才真正发起 putgetstart_date/end_date以 ISO 字符串序列化传输entity_df若为 SQL 字符串则放入entity_df_sql参数pull_all_from_table_or_query从数据源全量拉取putgetpull_latest_from_table_or_query拉取每个实体最新一条putgetwrite_logged_features写回推理日志特征仅do_putfeature_service_name作为参数offline_write_batch批量写入特征视图数据仅do_putvalidate_data_source / get_table_column_names_and_types_from_data_sourceapply 阶段的数据源校验与列 schema 探测前者 put、后者 putget其中get_historical_features是懒执行lazy的典型它只是把feature_view_names、feature_refs、project、full_feature_names、name_aliases等参数封进RemoteRetrievalJob真正的网络交互发生在调用to_df()/to_arrow()时——_to_arrow_internal触发_send_retrieve_remote结果 Arrow 表再to_pandas()回到 pandas。RemoteRetrievalJob.persist甚至复用了同一协议把“把查询结果落地为 SavedDataset”这一动作也放到远端执行remote.py。在安全方面build_arrow_flight_client检查仓库配置中的auth_config当认证类型不是AuthType.NONE时会通过FlightAuthInterceptorFactory给客户端挂上认证拦截器remote.py使每次 Flight 调用携带鉴权信息认证与授权的完整配置可参考 permission.md。七、生产化Kubernetes 部署与权限模型在集群中服务端可通过 Feast Operator 的 FeatureStore CR 直接声明 offlineStore 服务见 offline-feature-server.mdapiVersion: feast.dev/v1 kind: FeatureStore metadata: name: sample-offline-server spec: feastProject: my_project services: offlineStore: server: {}更多 FeatureStore CR 写法可参考 infra/feast-operator/config/samplesK8s 部署的完整背景见 running-feast-in-production.md。服务端暴露的每个端点都有对应的 RBAC 权限要求来自 offline-feature-server.md 的权限矩阵端点资源类型权限说明offline_write_batchFeatureViewWrite Offline向离线存储写批次数据write_logged_featuresFeatureServiceWrite Offline写推理日志特征persistDataSourceWrite Offline把读取结果持久化到离线存储get_historical_featuresFeatureViewRead Offline检索历史特征pull_all_from_table_or_queryDataSourceRead Offline全量拉取数据pull_latest_from_table_or_queryDataSourceRead Offline拉取最新数据八、小结与适用边界架构分工客户端持有 registry 元数据与remote委托配置服务端持有真实离线数据并以 Arrow FlightgRPC暴露OfflineStore接口。取数、写日志、批量写入、数据源校验全部走统一的 put命令数据 get结果协议列式 Arrow 数据避免了二次序列化开销。最小可运行链路服务端feast -c feature_repo apply feast -c feature_repo serve_offline默认127.0.0.1:8815可用--host/--port/--key/--cert覆盖客户端配置offline_store: {type: remote, host, port}后照常调用get_historical_features。生产加固scheme: httpscert建立 TLS 通道connection_retries控制瞬时错误重试默认 3 次、指数退避Operator 的 FeatureStore CR 声明式部署RBAC 按端点粒度授权。适用前提客户端与服务端必须指向同一project的 registry 元数据且entity_key_serialization_version等序列化约定保持一致remote 存储支持的功能集合与服务端 SDK 直连离线存储的能力一致详见 offline-stores overview 中的功能矩阵。参考路径索引示例 README 与运行步骤examples/remote-offline-store/README.md服务端仓库offline_server/feature_repo/feature_store.yaml、offline_server/feature_repo/example_repo.py客户端仓库offline_client/feature_store.yaml、offline_client/test.py客户端实现sdk/python/feast/infra/offline_stores/remote.pyCLI 与常量sdk/python/feast/cli/serve.py、sdk/python/feast/constants.py、sdk/python/feast/feature_store.py参考文档docs/reference/offline-stores/remote-offline-store.md、docs/reference/feature-servers/offline-feature-server.md、docs/getting-started/concepts/permission.md【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考