Feast 0.12 接入 AWS Redshift 与 DynamoDB:构建云上特征存储的离线/在线双后端
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
2021 年 8 月发布的 Feast 0.12 为开源特征存储引入了两大 AWS 原生后端:以 Redshift 作为离线存储(offline store)与数据源、以 DynamoDB 作为在线存储(online store),并借助 Feature Service 实现多特征视图的按需分组与统一读取。本文基于当前仓库源码与文档,完整讲解这两类后端的配置方式、声明式 API 用法、物化流程、权限模型以及底层实现原理,帮助你在 AWS 上从零搭建一套可训练、可实时推理的 Feast 特征平台。
版本背景:Feast 0.12 的三大关键能力
Feast(Feature Store,面向 AI/ML 的开源特征存储)在 0.12 版本之前,离线存储主要依赖 BigQuery、Snowflake 等方案。0.12 的发布把生态扩展到 AWS 技术栈,核心新增能力有三项:
- AWS Redshift 作为离线存储:支持从 Redshift 拉取历史特征值构建训练数据集,并支持高吞吐的批式推理,同时 Redshift 也可以作为特征数据源(data source)接入。
- AWS DynamoDB 作为在线存储:支持将特征物化到 DynamoDB,以支撑生产环境大规模、高并发的在线推理请求。
- Feature Service 逻辑分组:将多个 Feature View 的特征按需逻辑分组,通过一次请求统一返回分组内全部特征。
这三项能力都可以通过简单的声明式 API 与feature_store.yaml配置变更来启用,无需修改业务代码。
Redshift 作为数据源与离线存储
声明式数据源:RedshiftSource
在 Feast 的声明式 API 中,数据源定义在 feature repo 目录下的 Python 文件中。RedshiftSource即用于描述"要从哪张 Redshift 表或哪个查询中取特征"。官方博客中最简用法如下:
from feast import RedshiftSource my_redshift_source = RedshiftSource(table="redshift_driver_table")从当前仓库源码 redshift_source.py 可以看到,RedshiftSource实际支持远比这丰富的参数,完整签名如下:
from feast import RedshiftSource my_redshift_source = RedshiftSource( name=None, # 可选,默认取 table 名 timestamp_field="", # 事件时间戳字段,用于 point-in-time 连接 table=None, # Redshift 表名(table 与 query 二选一) schema=None, # Redshift schema,默认 "public" created_timestamp_column="", # 行创建时间戳列,用于行去重 field_mapping=None, # 数据源列名 -> 特征表列名映射字典 query=None, # 自定义查询(table 与 query 二选一) description="", # 人类可读描述 tags=None, # 任意元数据键值对 owner="", # 维护者邮箱 database="", # Redshift 数据库名 connection_ref=None, # 凭据解析引用 )几个关键语义(与源码一致):
- table 与 query 必须二选一:源码在
__init__中强制校验,二者都为空会抛出ValueError;query方式不提供性能保证,官方推荐优先使用表引用。 - schema 默认
public:源码注释明确 "The default Redshift schema is named 'public'"。 - 名称默认取自表名:若未显式传
name,则自动使用table作为数据源名。 - 基于查询的数据源:例如
from feast import RedshiftSource my_redshift_source = RedshiftSource( query="SELECT timestamp as ts, created, f1, f2 " "FROM redshift_table", )数据源在注册与校验阶段,会通过 Redshift Data API 执行describe_table(表模式)或SELECT * FROM (<query>) LIMIT 1(查询模式)来获取列名与列类型,从而完成特征类型的推断(见 get_table_column_names_and_types)。类型映射由redshift_to_feast_value_type完成,支持八种基本类型,暂不支持数组类型。
离线存储配置:RedshiftOfflineStoreConfig
要真正把 Redshift 用作离线存储,需要在feature_store.yaml中配置offline_store段。官方参考文档 redshift.md 给出了完整示例:
project: my_feature_repo registry: data/registry.db provider: aws offline_store: type: redshift region: us-west-2 cluster_id: feast-cluster database: feast-database user: redshift-user s3_staging_location: s3://feast-bucket/redshift iam_role: arn:aws:iam::123456789012:role/redshift_s3_access_role对照源码 redshift.py 中的RedshiftOfflineStoreConfig,各字段含义如下:
| 配置项 | 类型 | 必填 | 说明 |
|---|---|---|---|
type | string | 是 | 固定为redshift |
cluster_id | string | 条件必填 | 预置集群(provisioned cluster)标识符 |
user | string | 条件必填 | 预置集群的用户名 |
workgroup | string | 条件必填 | Redshift Serverless 工作组标识符 |
region | string | 是 | Redshift 集群所在 AWS 区域 |
database | string | 是 | Redshift 数据库名 |
s3_staging_location | string | 是 | 用于向 Redshift 导入/导出数据的 S3 路径 |
iam_role | string | 是 | Redshift 访问 S3 所用的 IAM 角色 ARN |
源码中的模型校验器(require_cluster_and_user_or_workgroup)明确了两类部署形态的约束:
- 预置集群(Provisioned):必须同时提供
cluster_id与user; - Redshift Serverless:必须提供
workgroup; - 二者不能同时指定,否则抛出
ValueError。
Redshift Serverless 的配置示例如下:
project: my_feature_repo registry: data/registry.db provider: aws offline_store: type: redshift region: us-west-2 workgroup: feast-workgroup database: feast-database s3_staging_location: s3://feast-bucket/redshift iam_role: arn:aws:iam::123456789012:role/redshift_s3_access_role离线存储功能矩阵
官方文档 redshift.md 列出了 Redshift 离线存储支持的能力,可作为选型与验收依据:
| 能力 | Redshift |
|---|---|
get_historical_features(point-in-time 正确连接) | 支持 |
pull_latest_from_table_or_query(拉取最新特征值) | 支持 |
pull_all_from_table_or_query(读取已保存数据集) | 支持 |
offline_write_batch(将 dataframe 持久化到离线存储) | 支持 |
write_logged_features(持久化日志特征) | 支持 |
RedshiftRetrievalJob导出能力方面:可导出 dataframe、arrow table、arrow batches、SQL、数据仓库,支持本地执行基于 Python 的 on-demand transforms、查询计划预览与分区数据读取;暂不支持导出到数据湖(S3/GCS 等)与 Spark dataframe。
底层实现:point-in-time 连接的 SQL 生成
RedshiftOfflineStore的核心价值在于所有连接都发生在 Redshift 内部。当调用get_historical_features时,pull_latest_from_table_or_query 会生成一条基于ROW_NUMBER() OVER (PARTITION BY <join_key> ORDER BY <timestamp> DESC, <created_ts> DESC)的去重 SQL:先按实体键分区、按时间戳倒序排序,再取出_feast_row = 1的最新行,并用BETWEEN TIMESTAMP限定事件时间窗,从而保证 point-in-time 正确性。实体 dataframe 既可以是 SQL 查询,也可以是 Pandas dataframe——后者会被临时上传到 Redshift 完成连接。
DynamoDB 作为在线存储
配置与启用
要将 DynamoDB 设为在线存储,只需在feature_store.yaml中声明online_store段。博客给出的最小配置:
project: fraud_detection registry: data/registry.db provider: aws online_store: type: dynamodb region: us-west-2对照源码 dynamodb.py 的DynamoDBOnlineStoreConfig,可用配置项远比最小示例丰富:
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
type | string | dynamodb | 在线存储类型选择器 |
region | string | 必填 | AWS 区域名 |
batch_size | int | 100 | 单次 BatchGetItem 请求条目数(DynamoDB 上限 100) |
endpoint_url | string | null | 本地开发(如http://localhost:8000)或 VPC 端点 |
table_name_template | string | {project}.{table_name} | 表名模板 |
consistent_reads | bool | false | 是否强制强一致读 |
warmup_connections | bool | false | 初始化时是否预热连接池 |
tags | dict | null | 附加到每张表的 AWS 资源标签 |
session_based_auth | bool | false | 是否使用基于 session 的客户端认证 |
max_read_workers | int | 10 | 批量读并行线程数 |
max_pool_connections | int | 50 | 异步操作最大连接数 |
keepalive_timeout | float | 30.0 | 异步连接 keep-alive 超时(秒) |
connect_timeout | float | 5 | 建立连接超时(秒) |
read_timeout | float | 10 | 读取超时(秒) |
total_max_retry_attempts | int | 3 | 单请求最大尝试次数,映射 botocoreretries.total_max_attempts |
retry_mode | string | adaptive | 重试模式:legacy/standard/adaptive,映射 botocoreretries.mode |
包含性能调优选项的完整示例(见 dynamodb.md):
project: my_feature_repo registry: data/registry.db provider: aws online_store: type: dynamodb region: us-west-2 batch_size: 100 max_read_workers: 10 consistent_reads: false warmup_connections: true物化(Materialize):把离线特征写入 DynamoDB
将特征物化进 DynamoDB 在线存储只需一条 CLI 命令:
feast materialize该命令对应 feature_store.py 中的materialize方法,内部流程为:从离线存储(此处为 Redshift)拉取指定时间窗内的最新特征值,再调用在线存储的online_write_batch批量写入。DynamoDB 侧的实现位于 online_write_batch,它使用 boto3batch_writer(以entity_id为主键、overwrite_by_pkeys去重)自动重发未处理条目,适合大批量加载;异步版本online_write_batch_async则通过 aiobotocore 客户端与_latest_data_to_write去重后批量写入。
读取路径与性能设计
在线读取由online_read与online_read_async实现(dynamodb.py),关键设计包括:
- 分批与并行:实体键先按
batch_size(默认 100)切成多个批次,单批次直接调用batch_get_item;多批次则用ThreadPoolExecutor(上限max_read_workers)并行执行。 - 投影表达式:当指定
requested_features时,通过ProjectionExpression只拉取所需特征列,减少网络传输(_build_projection_expression)。 - O(1) 结果归位:由于
BatchGetItem不保证返回顺序,读取响应先构造成以entity_id为键的字典再做 O(1) 查找,避免排序开销(_process_batch_get_response)。 - 表结构:每张 Feature View 对应一张 DynamoDB 表,表名由
table_name_template生成(默认{project}.{table_name}),主键为entity_id(字符串 HASH 键),BillingMode为PAY_PER_REQUEST按量计费(见 update)。 - 异步连接池:
initialize会按max_pool_connections、keepalive_timeout等参数构造 aiobotocore 客户端;warmup_connections: true时通过一次轻量describe_limits预热 TCP/TLS 连接池,消除首次检索的冷启动开销。
DynamoDB 在线存储功能矩阵
根据 dynamodb.md:DynamoDB 在线存储支持特征写入、读取、基础设施更新/拆除、on-demand transforms、Python SDK 读取、entityless feature views;暂不支持 Java/Go SDK 读取、并发写同一键、检索期 TTL 与过期数据删除。
Feature Service:逻辑分组与统一读取
当需要从多个 Feature View 中逻辑分组特征时,使用 Feature Service。这样无论调用feature_store.get_historical_features(...)还是feature_store.get_online_features(...),所有已分组特征都会一次性从特征存储返回,避免在客户端手工拼接多个视图的结果。
源码 feature_service.py 明确定义了它的语义:Feature Service 定义一个或多个 Feature View 特征的逻辑分组,该分组在训练或推理时可整体取出。典型用法:
from feast import FeatureService driver_stats_fs = FeatureService( name="driver_stats_service", features=[ driver_hourly_stats_view, driver_daily_stats_view, ], )在 feature repo 中定义后执行feast apply即可注册。请求侧用法:
# 离线训练 training_df = store.get_historical_features( entity_df=entity_df, features=driver_stats_fs, ).to_df() # 在线推理 online_features = store.get_online_features( features=driver_stats_fs, entity_rows=[{"driver_id": 1001}], ).to_dict()从源码看,Feature Service 支持FeatureView、OnDemandFeatureView与LabelView的组合,还支持tags、description、owner等元数据字段,可用于特征治理与检索。
权限模型:让 Feast 安全地操作 AWS 资源
Redshift 离线存储所需权限
根据 redshift.md,Feast 执行不同命令所需的权限如下:
| 命令 | 所需权限 | 资源 |
|---|---|---|
| Apply | redshift-data:DescribeTable、redshift:GetClusterCredentials | Redshift dbuser / dbname / cluster ARN |
| Materialize | redshift-data:ExecuteStatement、redshift-data:DescribeStatement、s3:ListBucket、s3:GetObject、s3:DeleteObject | cluster ARN、S3 bucket |
| Get Historical Features | redshift-data:ExecuteStatement、redshift:GetClusterCredentials、redshift-data:DescribeStatement、S3 读写 | Redshift ARN、S3 bucket |
此外,Redshift 自身需要一个 IAM 角色来执行UNLOAD与COPY命令访问 S3(即offline_store.iam_role),并需建立信任关系,仅允许redshift.amazonaws.com服务代入该角色。
DynamoDB 在线存储所需权限
Feast 操作 DynamoDB 需要的权限(dynamodb.md):
| 命令 | 所需权限 |
|---|---|
| Apply | dynamodb:CreateTable、dynamodb:DescribeTable、dynamodb:DeleteTable、dynamodb:TagResource |
| Materialize | dynamodb:BatchWriteItem |
| Get Online Features | dynamodb:BatchGetItem |
资源统一限定为arn:aws:dynamodb:<region>:<account_id>:table/*。需要说明的是,从源码看(update),Feast 在建表前会先describe_table探测表是否存在——若表由 Terraform 等外部工具预建且 IAM 角色缺少CreateTable权限,可以正常跳过建表;标签更新在权限不足时也只会记录警告而不会中断流程。
安装与快速上手
使用上述 AWS 后端,需要安装带 AWS 依赖的 Feast:
pip install 'feast[aws]'随后可用 AWS 模板初始化一个 feature repo:
feast init REPO_NAME -t aws然后编辑feature_store.yaml(参考上文 Redshift + DynamoDB 双后端配置),在 feature repo 中定义数据源、Feature View 与 Feature Service,依次执行:
feast apply # 注册数据源、特征视图与 Feature Service,并创建 DynamoDB 表 feast materialize # 将 Redshift 中的特征物化进 DynamoDB之后即可通过FeatureStore.get_historical_features(...)构建训练数据集、通过get_online_features(...)支撑在线推理。
总结
Feast 0.12 通过三项能力完成了对 AWS 技术栈的深度接入:Redshift 兼顾数据源与离线存储,让 point-in-time 连接直接在数仓内完成,并提供pull_latest/pull_all/offline_write_batch等完整离线能力;DynamoDB 提供可水平扩展、按量计费的在线存储,通过批量读写、并行读取、投影表达式与连接池预热等机制保障在线推理的吞吐与延迟;Feature Service 则将多视图特征逻辑分组,统一了训练与推理两条读取路径。三者结合feature_store.yaml的声明式配置与feast materialize单命令物化,构成了一个完整的 AWS 原生特征存储闭环。若需对比各离线/在线存储的功能差异,可查阅 offline-stores/overview.md 与 online-stores/overview.md 中的功能矩阵。
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考