Feast 的 AWS Lambda 批式物化引擎(alpha):配置、镜像构建与工作原理
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
本指南围绕 Feast 官方文档 lambda.md 展开,系统讲解基于 AWS Lambda 的批式物化(batch materialization)引擎:它如何依赖离线存储把特征值落盘到 S3、再由 Lambda 函数灌入在线存储,以及如何在feature_store.yaml中完成最小可用配置。读完本文,你将掌握该引擎的完整配置参数、物化镜像的构建方式,以及从源码层面理解其创建函数、并发执行、超时重试和资源销毁的完整生命周期。
引擎定位与核心工作流
AWS Lambda 批式物化引擎在 Feast 中处于alpha 阶段(源码类注释中亦明确标注WARNING: This engine should be considered "Alpha" functionality.,见 lambda_engine.py),其设计思路非常简洁:物化过程的"计算"发生在离线存储侧与 Lambda 侧,Feast 仅负责编排。
整个流程可以拆解为两条链路:
- 离线侧产出:引擎调用离线存储(offline store)执行特征抽取,并通过其
to_remote_storage()方法把结果特征值以 Parquet 文件形式输出到 S3; - 在线侧加载:引擎为每个 S3 上的 Parquet 文件并发调用一个 Lambda 函数,由该函数读取文件并写入在线存储(online store),例如 Redis 或 DynamoDB。
这一点正是该引擎与本地物化、Spark 物化的核心区别:你无需常驻的计算集群,物化负载被拆散成一个个短暂的、按文件粒度执行的 Lambda 调用,天然具备按需伸缩与"用完即走"的特性。
需要特别指出的是,该引擎只承担materialize()(物化),不支持历史特征检索。其get_historical_features()直接抛出NotImplementedError(见 lambda_engine.py),因此在配置该引擎后,训练数据集的生成仍需依赖其他路径。
在 feature_store.yaml 中启用 Lambda 引擎
官方文档给出的最小配置如下(原示例以 Snowflake 为离线存储):
# feature_store.yaml ... offline_store: type: snowflake.offline ... batch_engine: type: lambda lambda_role: [your iam role] materialization_image: [image uri of above Docker image]配置要点说明:
batch_engine.type: lambda:在 Feast 仓库中,批式物化引擎通过batch_engine段落选择。当前源码中该配置类名为LambdaComputeEngineConfig(见 lambda_engine.py),type字段被限制为字面量"lambda";lambda_role:物化 Lambda 函数运行时使用的 IAM 角色 ARN。该角色需要具备读取 S3 上 Parquet 文件、访问在线存储(如 DynamoDB/Redis)以及必要的日志写入权限,Feast 源码在创建函数时直接将其作为Role参数传给 AWS Lambda API;materialization_image:承载 Lambda 处理逻辑的容器镜像 URI(存放于 Amazon ECR)。引擎在创建函数时以PackageType="Image"的方式部署该镜像。
说明:文档中提到的配置类名称为
LambdaMaterializationEngineConfig;在当前仓库源码中对应的类已演进为LambdaComputeEngineConfig,位于 lambda_engine.py。字段语义保持一致,仅命名随 ComputeEngine 抽象重构而更新。
此外,引擎的仓库级配置还支持按 FeatureView 进行运行时覆盖:ComputeEngine基类提供了_get_feature_view_engine_config()方法,将仓库默认的batch_engine_config与BatchFeatureView.batch_engine/StreamFeatureView.stream_engine中定义的覆盖项合并,且FeatureView 级配置优先级更高(见 base.py)。
构建物化镜像:Dockerfile 拆解
Feast 仓库内为 Lambda 引擎提供了可直接使用的 Dockerfile:sdk/python/feast/infra/compute_engines/aws_lambda/Dockerfile,构建后即可作为materialization_image使用。其要点如下:
FROM public.ecr.aws/lambda/python:3.9 RUN yum install -y git COPY --from=ghcr.io/astral-sh/uv:latest /uv /usr/local/bin/uv # Copy app handler code COPY sdk/python/feast/infra/materialization/lambda/app.py ${LAMBDA_TASK_ROOT} # Copy necessary parts of the Feast codebase COPY sdk/python sdk/python COPY protos protos COPY go go COPY pyproject.toml pyproject.toml COPY README.md README.md # Install Feast for AWS with Lambda dependencies RUN --mount=source=.git,target=.git,type=bind uv pip install --system --no-cache-dir -e '.[aws,redis]' # Set the CMD to your handler CMD [ "app.handler" ]几个值得注意的实现细节:
- 基础镜像:使用 AWS 官方 Lambda Python 3.9 运行时镜像(
public.ecr.aws/lambda/python:3.9),保证与 Lambda 运行时环境的兼容性; app.handler入口:镜像的CMD指向app.handler,即 Lambda 的 handler 为app.py中的handler(event, context)函数(见 app.py);- 依赖安装:通过
uv pip install -e '.[aws,redis]'以可编辑模式安装 Feast 本体,并带上 AWS 与 Redis 在线存储依赖。安装过程通过 bind mount 挂载.git目录,这是因为setuptools_scm需要访问 git 元数据来推断 Feast 版本号; - 构建前提:注释明确提示该 Dockerfile 假定从仓库根目录执行构建(因为
COPY的源路径均相对仓库根目录)。另外,当前源码中app.py的实际路径是sdk/python/feast/infra/compute_engines/aws_lambda/app.py,Dockerfile 内注释沿用了旧路径,使用时请以实际仓库结构为准。
源码视角:物化一次到底发生了什么
函数创建与销毁(基础设施生命周期)
LambdaComputeEngine.update()负责在feast apply阶段创建 Lambda 函数(见 lambda_engine.py):
- 函数名固定为
feast-materialize-<project>(超过 64 字符时截断,见构造函数 lambda_engine.py),一个 Feast 项目对应一个物化函数; - 使用
PackageType="Image"与materialization_image进行镜像部署; Timeout固定为 600 秒(DEFAULT_TIMEOUT);- 打上三类标签(Tags):
feast-owned: "True"、project以及feast-sdk-version,方便识别与成本归属; - 创建后通过
get_waiter("function_active")等待函数进入 Active 状态再返回。
对应的teardown_infra()则调用delete_function销毁该函数(见 lambda_engine.py),即feast teardown时完成资源回收。
物化编排:_materialize_one
真正执行物化的是_materialize_one()(见 lambda_engine.py),其步骤与上文工作流一一对应:
- 从 Registry 解析 FeatureView 的实体(Entity),并借助
_get_column_names计算 join key、特征列、时间戳列; - 调用
offline_store.pull_latest_from_table_or_query()执行 point-in-time 正确的最新特征抽取; - 调用
offline_job.to_remote_storage()将结果写出为 S3 上的 Parquet 文件列表;若文件数为 0,直接返回SUCCEEDED状态; - 使用
ThreadPoolExecutor并发调用 Lambda,最大并发数为文件数与 20 的较小值(max_workers = num_files if num_files <= 20 else 20); - 每个调用携带的 payload 包含三项关键信息:
FEATURE_STORE_YAML_BASE64:base64 编码的feature_store.yaml内容(该环境变量名定义在 constants.py);view_name:待物化的 FeatureView 名称;path:当前 Parquet 文件在 S3 上的路径;view_type:固定为"batch";
- 汇总所有调用结果,全部成功则返回
SUCCEEDED,任一失败则返回携带错误信息的ERROR状态。
超时重试:应对 DynamoDB 限流
invoke_with_retries()(见 lambda_engine.py)实现了最多LAMBDA_TIMEOUT_RETRIES = 5次的重试逻辑。其触发条件非常具体:当 Lambda 返回的错误消息包含"Task timed out after"时进行重试。源码注释解释了原因——大批量写入时 DynamoDB 可能限流导致函数超时,一旦表完成扩容,重试即可成功。该机制使得物化流程对在线存储的瞬时限流具备一定的韧性。
Lambda 内部:app.handler做了什么
每个被调用的 Lambda 实例执行 app.py 中的handler:
- 从事件中取出
FEATURE_STORE_YAML_BASE64,解码后在临时目录中还原出feature_store.yaml; - 基于该配置初始化
FeatureStore,并依据view_name获取 FeatureView(batch 或 stream 两种类型均支持); - 解析
path(格式s3://bucket/key),用 PyArrow 读取 Parquet 表,若 FeatureView 定义了field_mapping则先执行字段映射; - 将实体 dtype 转换为对应的 value type;
- 按
DEFAULT_BATCH_SIZE = 10_000行分批,把 Arrow 记录批次转换为 Feast 内部 proto 结构,再通过store._provider.online_write_batch()批量写入在线存储; - 记录成功写入的行数(含
num_updated_rows结构化日志字段),异常时记录堆栈并抛出。
可以看到,单次 Lambda 调用只负责"一个文件 → 在线存储"的搬运与写入,天然适合横向拆分与并行。
限制与适用前提
综合官方文档与源码,使用该引擎前需要明确以下边界:
- alpha 状态:官方文档与源码类注释均明确标注 alpha,接口与行为可能随版本调整;
- 仅支持物化:
get_historical_features未实现,历史特征检索不可用; - 依赖离线存储的远端落盘能力:物化依赖离线存储的
to_remote_storage()将结果写到 S3,因此离线存储必须支持该能力(例如 Snowflake 离线存储),本地文件类离线存储并不适合; - 需要预先构建并推送物化镜像:
materialization_image是必填项,需按照上文 Dockerfile 在仓库根目录构建并推送到 ECR; - IAM 权限面较大:
lambda_role需同时覆盖 S3 读取、在线存储读写与 CloudWatch 日志;同时 Feast 会为每个项目创建/销毁一个 Lambda 函数,需要注意账户级配额与成本。
小结
AWS Lambda 批式物化引擎为 Feast 提供了一种"无常驻集群"的物化选项:离线存储把特征结果落到 S3,Feast 编排并发调用镜像化的 Lambda 函数逐文件灌入在线存储。通过本文的配置示例与源码拆解可以看到,其核心配置仅有type、lambda_role、materialization_image三项,而镜像构建、函数生命周期管理、并发调度与超时重试均已由 Feast 封装完成。适合希望在 AWS 上以低成本、按需伸缩方式运行物化任务,并愿意接受 alpha 状态接口变化的用户进一步评估。
【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考