如何可扩展地批量修改SageMaker Feature Store记录,保留OfflineStore
批量处理SageMaker Feature Store记录:保留OfflineStore、隐藏OnlineStore的可扩展方案
方案1:基于delete_record的异步批量调用(大规模场景首选)
- 第一步:定位目标记录
从OfflineStore的Athena关联表中查询需要标记隐藏的记录主键(Feature Group的record_identifier_feature_name字段),示例查询:SELECT record_id FROM "your_offline_db"."your_feature_group_table" WHERE [筛选条件] -- 比如特定时间范围、特征值匹配规则 - 第二步:批次拆分与异步执行
将查询到的主键列表按批次拆分(建议每批次100-500条,根据API配额调整),用多线程/异步框架批量调用delete_record接口,避免单线程串行效率低下:import boto3 from concurrent.futures import ThreadPoolExecutor from tenacity import retry, stop_after_attempt, wait_exponential sagemaker_fs_client = boto3.client('sagemaker-featurestore-runtime') FEATURE_GROUP_NAME = "your-feature-group-name" @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) def delete_record(record_id): try: sagemaker_fs_client.delete_record( FeatureGroupName=FEATURE_GROUP_NAME, RecordIdentifierValueAsString=str(record_id), EventTime="2024-05-20T00:00:00Z" # 用当前时间或原记录EventTime,确保覆盖最新版本 ) except Exception as e: print(f"删除失败 {record_id}: {str(e)}") raise # 从Athena查询结果获取主键列表 target_record_ids = ["id_001", "id_002", "id_003", ...] BATCH_SIZE = 200 MAX_WORKERS = 10 # 多线程批量处理 with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: for i in range(0, len(target_record_ids), BATCH_SIZE): batch = target_record_ids[i:i+BATCH_SIZE] executor.map(delete_record, batch) - 第三步:结果验证
随机调用get_record接口检查目标主键,确认OnlineStore返回空结果,同时OfflineStore的原记录仍可通过Athena查询到。
方案2:批量写入is_deleted标记(需同步更新特征的场景)
如果需要同时更新其他特征值并标记隐藏,可通过PutRecord批量写入is_deleted=True的记录(Feature Store会自动将该版本同步到OnlineStore,原版本保留在OfflineStore):
import boto3 from concurrent.futures import ThreadPoolExecutor sagemaker_fs_client = boto3.client('sagemaker-featurestore-runtime') FEATURE_GROUP_NAME = "your-feature-group-name" def update_to_deleted(record): try: sagemaker_fs_client.put_record( FeatureGroupName=FEATURE_GROUP_NAME, Record=[ {"FeatureName": "record_id", "ValueAsString": record["id"]}, {"FeatureName": "is_deleted", "ValueAsString": "True"}, {"FeatureName": "event_time", "ValueAsString": record["event_time"]} # 其他特征保留原值,确保覆盖最新版本 ] ) except Exception as e: print(f"更新失败 {record['id']}: {str(e)}") # 从OfflineStore获取需更新的完整记录(含主键、event_time) records_to_update = [{"id": "id_001", "event_time": "2024-05-20T00:00:00Z"}, ...] BATCH_SIZE = 100 with ThreadPoolExecutor(max_workers=8) as executor: for i in range(0, len(records_to_update), BATCH_SIZE): batch = records_to_update[i:i+BATCH_SIZE] executor.map(update_to_deleted, batch)
可扩展优化要点
- 限流控制:根据AWS SageMaker Feature Store的API配额(默认1000 TPS)调整并发数和批次大小,避免触发限流
- 错误重试:给API调用添加指数退避重试逻辑,处理临时网络波动或限流
- 增量处理:定期批量处理时,通过Athena查询增量数据(基于
event_time筛选新增待处理记录),避免全量扫描 - 监控告警:用CloudWatch监控API调用成功率,设置告警阈值,失败时及时触发通知
内容的提问来源于stack exchange,提问作者Lucas de Brito Silva
相关产品推荐
相关产品推荐

