资讯动态

Feast Kubernetes 物化引擎集成测试完全指南:从 EKS 集群搭建到端到端验证

发布时间:2026/9/17 12:31:00 来源:尧图企业网站定制
Feast Kubernetes 物化引擎集成测试完全指南从 EKS 集群搭建到端到端验证【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast本指南围绕 Feast开源 Feature Store的 Kubernetes 批处理物化引擎完整讲解其集成测试的搭建、执行与验证方法如何用 eksctl 创建单节点 EKS 集群、运行端到端物化一致性测试并深入源码剖析 Kubernetes 计算引擎的配置项与底层作业执行原理。读完本文你将掌握 Feast 在 AWS 环境下使用 Kubernetes 集群完成从离线存储到在线存储物化的完整测试链路。一、Kubernetes 引擎在 Feast 中的定位Feast 支持多种批处理计算引擎batch engine来执行物化materialization任务包括 Spark、Ray、Snowflake、AWS Lambda 与 Kubernetes。其中 Kubernetes 引擎的独特价值在于它复用现有 Kubernetes 集群作为计算资源池将物化任务以 Kubernetes Job 的形式调度执行每个 Job 按离线数据分片parquet 文件生成对应数量的 Pod 并行处理从而为 AWS 用户提供了一种无需预置专用计算框架如 Spark on EMR的弹性物化方案。从 sdk/python/feast/repo_config.py 的get_batch_engine_config_from_type可以看出batch_engine配置通过类型注册表动态解析为对应的引擎配置类{type: k8s}会被解析为KubernetesComputeEngineConfig最终由 sdk/python/feast/infra/compute_engines/kubernetes/k8s_engine.py 中的KubernetesComputeEngine负责实际执行。本文所讲解的集成测试文档位于 sdk/python/tests/integration/materialization/kubernetes/README.md它给出了验证该引擎的最小闭环创建 EKS 集群 → 运行测试 → 删除集群。二、测试前置条件准备 EKS 集群2.1 依赖工具运行 Kubernetes 引擎集成测试前需要准备安装了 eksctl 的机器eksctl 是 AWS 官方推荐的 EKS 集群创建工具配置好 AWS 凭证默认使用~/.aws/credentials或环境变量且当前账号具备创建 EKS 集群与 Redshift、DynamoDB 等资源的权限测试过程会真实创建 AWS 云资源并产生费用建议使用专用测试账号。2.2 集群配置文件详解测试目录下已附带 EKS 集群配置文件 sdk/python/tests/integration/materialization/kubernetes/eks-config.yaml内容如下apiVersion: eksctl.io/v1alpha5 kind: ClusterConfig metadata: name: feast-cluster version: 1.22 region: us-west-2 managedNodeGroups: - name: ng-1 instanceType: c6a.large desiredCapacity: 1 privateNetworking: true各字段含义与工程考量字段取值说明metadata.namefeast-cluster集群名称后续删除集群时需与此保持一致metadata.version1.22Kubernetes 控制面版本该测试配置锁定版本以保证 API 兼容性metadata.regionus-west-2AWS 区域与测试中 Redshift、DynamoDB 所在区域保持一致managedNodeGroups[].instanceTypec6a.large计算实例规格2 vCPU / 4 GiB足够跑通单节点物化测试desiredCapacity1单节点即可满足测试需求desiredCapacity: 1同时控制了测试成本privateNetworkingtrue节点使用私有网络不分配公网 IP更安全注意配置锁定 Kubernetes 1.22 与 us-west-2 区域若你的账户默认区域不同或需要更新 Kubernetes 版本请相应调整该文件。2.3 创建与删除集群的命令在测试目录下执行README 原文命令eksctl create cluster -f ./eks-config.yaml该命令会依次完成创建 VPC/子网等网络资源、创建 EKS 控制面、创建托管节点组、将kubeconfig写入默认位置供后续 kubectl 使用。集群就绪通常需要 1020 分钟。测试完成后删除集群eksctl delete cluster feast-clusterdelete cluster会自动清理该集群关联的 CloudFormation 栈包括节点组与网络资源避免资源残留与持续计费。三、运行 Kubernetes 引擎集成测试3.1 测试入口与标记测试用例位于 sdk/python/tests/integration/materialization/kubernetes/test_k8s.py其中test_kubernetes_materialization同时带有两个装饰器pytest.mark.integration pytest.mark.skip(reasonRun this test manually after creating an EKS cluster.)pytest.mark.integration标注为集成测试默认不被普通单元测试收集pytest.mark.skip默认跳过强制要求先手动创建 EKS 集群后再运行。也就是说该测试不是 CI 常驻用例而是需要人工确认环境就绪后主动执行。从仓库结构看本测试属于 sdk/python/tests/integration/materialization 下针对不同物化引擎如 Spark、Snowflake、Lambda 等的集成测试集合之一Kubernetes 引擎在此目录中拥有独立的子目录。3.2 测试环境装配测试通过IntegrationTestRepoConfig声明整套运行环境config IntegrationTestRepoConfig( provideraws, online_store{type: dynamodb, region: us-west-2}, offline_store_creatorRedshiftDataSourceCreator, batch_engine{type: k8s}, registry_locationRegistryLocation.S3, )各配置项对应的真实资源provideraws使用 AWS 基础设施对应 sdk/python/feast/infra/passthrough_provider.py 中定义的 AWS Provideronline_store{type: dynamodb}在线特征存储使用 DynamoDBus-west-2offline_store_creatorRedshiftDataSourceCreator离线数据源使用 Redshiftbatch_engine{type: k8s}核心配置指定物化计算引擎为 Kubernetesregistry_locationRegistryLocation.S3Feature Registry 存储在 S3。随后construct_test_environment(config, None)会基于该配置构建完整的 Feast 测试环境包括 Feature Store 客户端、数据源创建器等。3.3 测试数据集与特征定义测试使用create_basic_driver_dataset()生成模拟司机数据来自 sdk/python/tests/data/data_creator.py并注册一个driver_hourly_statsFeatureViewdriver Entity( namedriver_id, join_keydriver_id, value_typeValueType.INT64, ) driver_stats_fv FeatureView( namedriver_hourly_stats, entities[driver_id], ttltimedelta(weeks52), features[Feature(namevalue, dtypeValueType.FLOAT)], batch_sourceds, ) fs.apply([driver, driver_stats_fv])注意field_mapping{ts_1: ts}将数据集中的ts_1时间戳列映射为 Feast 约定的ts事件时间列这是 FeatureView 正确执行 point-in-time 语义的前提。3.4 一致性验证流程测试的核心是调用 sdk/python/tests/utils/e2e_test_validation.py 中的validate_offline_online_store_consistency其验证逻辑为以测试数据生成的时间戳为分界点split_dtdf[ts_1][4].to_pydatetime() - timedelta(seconds1)将其作为首次物化的截止时间调用fs.materialize(feature_views[fv.name], start_datestart_date, end_dateend_date)触发第一轮全量物化此时计算引擎为 Kubernetes校验在线存储DynamoDB与离线存储Redshift对 driver_id1、2、3 的特征值是否一致见_check_offline_and_online_features中对get_online_features与get_historical_features结果的断言调用fs.materialize_incremental(feature_views[fv.name], end_datenow)触发第二轮增量物化断言 Registry 中materialization_intervals被正确更新为两个时间区间[start_date, end_date]与[end_date, now]从而验证物化区间元数据持久化正确最后再次校验 driver_id3 的增量特征值。该测试同时覆盖了materialize与materialize_incremental两种入口是验证 Kubernetes 引擎端到端正确性的核心用例。运行结束后fs.teardown()清理测试资源。四、底层原理Kubernetes 物化作业如何工作4.1 引擎初始化与离线数据拉取KubernetesComputeEngine.__init__k8s_engine.py中调用k8s_config.load_config()加载本地 kubeconfig并实例化CoreV1Api管理 ConfigMap/Pod与BatchV1Api管理 Job。_materialize_onek8s_engine.py的流程为从 Registry 解析 FeatureView 的实体Entity计算 join key、特征列与时间戳列调用offline_store.pull_latest_from_table_or_query(...)从离线存储本测试为 Redshift拉取最新特征数据将结果通过offline_job.to_remote_storage()写出为多个 parquet 分片文件每个分片对应一个物化 Pod 的工作负载依据synchronous配置决定执行方式。4.2 同步与异步两种执行模式从源码可以推断引擎支持两种物化调度模式对应配置synchronous异步模式默认synchronous: false直接创建单个 Kubernetes JobJob 内completions pods每个分片一个 Podparallelism min(pods, max_parallelism)采用completionMode: Indexed让每个 Pod 通过JOB_COMPLETION_INDEX环境变量识别自己负责的分片同步模式synchronous: true按job_batch_size默认 100分批创建 Job 并轮询等待完成每 30 秒查询一次状态见_await_path_materialization支持失败时打印 Pod 日志print_pod_logs_on_failure。4.3 作业与工作负载的创建_create_kubernetes_jobk8s_engine.py做了两件事创建 ConfigMap_create_configuration_map将序列化后的feature_store.yaml与materialization_config.yaml含分片路径列表和 FeatureView 名称写入名为feast-{job_id}的 ConfigMap创建 Job_create_job_definition构造batch/v1Job 清单容器执行python main.py镜像默认为feast/feast-k8s-materialization:latestConfigMap 挂载到/var/feast/。此外 Job 定义中还包含安全加固细节securityContext设置allowPrivilegeEscalation: false、drop: [ALL]并仅添加NET_BIND_SERVICE能力ttlSecondsAfterFinished: 3600保证作业结束后 1 小时自动清理。4.4 工作 Pod 内的物化执行每个 Pod 运行 sdk/python/feast/infra/compute_engines/kubernetes/main.pywith open(/var/feast/feature_store.yaml) as f: feast_config yaml.safe_load(f) with open(/var/feast/materialization_config.yaml) as b: materialization_cfg yaml.safe_load(b) ... KubernetesMaterializer( configconfig, feature_viewstore.get_feature_view(materialization_cfg[feature_view]), pathsmaterialization_cfg[paths], worker_indexint(os.environ[JOB_COMPLETION_INDEX]), ).run()KubernetesMaterializer.run()的处理链路为按JOB_COMPLETION_INDEX取出自己负责的 parquet 分片路径用pyarrow读取分片并按MINI_BATCH_SIZE默认 1000 行由引擎配置注入切分为小批量执行字段映射field_mapping将 Arrow 表转换为 Feast 的 protobuf 特征行通过feature_store._get_provider().online_write_batch(...)批量写入在线存储DynamoDB。作业状态由 k8s_materialization_job.py 中的KubernetesMaterializationJob.status()读取 Job 状态并映射为WAITING / RUNNING / SUCCEEDED / ERROR四种状态供同步模式轮询及上层日志输出使用。4.5 镜像构建物化 Pod 使用的镜像通过 sdk/python/feast/infra/compute_engines/kubernetes/Dockerfile 构建基于 Debian 11 slim使用 uv 安装包含aws, gcp, k8s, snowflake, postgresextras 的完整 Feast SDK并将仓库根目录的sdk/python、protos、go、pyproject.toml一并拷贝进镜像。这也解释了为何测试要求集群节点具备拉取该镜像的能力——若使用私有镜像仓库需通过image_pull_secrets配置凭据。五、Kubernetes 引擎完整配置项参考在 feature_store.yaml 中可通过batch_engine字段启用该引擎。以下为 KubernetesComputeEngineConfig 支持的全部配置项均来自源码定义batch_engine: type: k8s # 必选引擎类型选择器 namespace: default # (可选) 创建 Service/ConfigMap/Job 的命名空间 image: feast/feast-k8s-materialization:latest # (可选) 物化作业容器镜像 env: [] # (可选) 注入 Pod 的环境变量列表可用于引用 K8s Secret image_pull_secrets: [] # (可选) 拉取镜像所需 Secret resources: {} # (可选) 物化容器的资源 requests/limits service_account_name: # (可选) 运行 Job 的 ServiceAccount常用于 IRSA 绑定 IAM 角色 annotations: {} # (可选) 作用于 Job 容器的注解如 IAM 角色关联、运维元数据 include_security_context_capabilities: true # (可选) 是否在 init/Job 容器中加入安全上下文 capabilities labels: {} # (可选) 追加到 Kubernetes 对象的额外标签 max_parallelism: 10 # (可选) 单个 Job 内并行运行的 Pod 上限 synchronous: false # (可选) 为 true 时逐个特征等待物化完成再继续 retry_limit: 2 # (可选) 物化工作 Pod 的最大重试次数 mini_batch_size: 1000 # (可选) 每次写入操作处理的批处理行数 active_deadline_seconds: 86400 # (可选) 物化作业允许运行的最长时间 job_batch_size: 100 # (可选) 同步模式下每个 Job 处理的 Pod 数量 print_pod_logs_on_failure: true # (可选) 同步模式下作业失败时打印 Pod 日志关键参数的工程解读service_account_name与annotations在 AWS 上通常配合 IRSAIAM Roles for Service Accounts使用让物化 Pod 具备访问 S3、Redshift、DynamoDB 的最小权限max_parallelism控制单 Job 并发度Job 的parallelism取值为min(pods, max_parallelism)可防止一次性创建过多 Pod 压垮集群retry_limit对应 Job 的backoffLimitactive_deadline_seconds对应activeDeadlineSeconds防止作业卡死占用资源mini_batch_size通过MINI_BATCH_SIZE环境变量注入 Pod控制每次online_write_batch的批量行数直接影响 DynamoDB 写入吞吐与成本synchronous: true时若job_batch_size max_parallelism源码会告警并自动将job_batch_size提升到max_parallelism见_materialize_one中的校验逻辑。六、从测试到生产的迁移要点集成测试验证的正是生产环境使用 Kubernetes 引擎的最小可行路径将其迁移到真实业务场景时需关注网络与权限测试中 Redshift 与 EKS 节点同处 us-west-2生产环境应保证物化 Pod 能访问离线存储Redshift/S3与在线存储DynamoDB通过安全组、VPC 对等或私有链接打通并配合service_account_name IRSA 做最小权限授权镜像管理将feast/feast-k8s-materialization镜像推送到私有仓库并固定版本 tag通过image_pull_secrets配置拉取凭证资源与并发根据离线数据量设置mini_batch_size、max_parallelism、resources避免 Pod 因 OOM 反复重启retry_limit默认仅 2 次观测与清理借助labels、annotations标记物化作业便于检索Job 默认ttlSecondsAfterFinished: 3600自动回收ConfigMap 在同步模式结束后也会被删除版本兼容本测试配置锁定 Kubernetes 1.22 与 Feast 仓库当前版本升级集群或 Feast 版本后应先重跑该集成测试。七、结语Feast 的 Kubernetes 物化引擎将特征物化从需要独立计算框架转变为复用已有 Kubernetes 集群而 sdk/python/tests/integration/materialization/kubernetes 目录下的集成测试README、eks-config.yaml、test_k8s.py提供了一个可复现的最小验证闭环eksctl create cluster一键建集群、运行带pytest.mark.integration的端到端一致性测试、eksctl delete cluster一键清理。结合 k8s_engine.py 与 main.py 的源码研读你可以深入理解 Indexed Job ConfigMap 分发 分片并行的完整机制并将其平稳迁移到生产环境。/DSMLparameter /DSMLinvoke /DSMLtool_calls【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价