You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何可扩展地批量修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 22:05:19