新增 BatchWriteRecord 一次调用批量写入多特征组最多 25 条记录,以及 ListRecords 枚举特征组内记录 ID,配有完整代码示例。
Amazon SageMaker Feature Store 是一个全托管的专用仓库,用于存储、共享和管理机器学习(ML)模型的特征。它提供低延迟的在线服务用于实时推理,支持历史保留和训练特征数据的离线存储,并同时支持流式摄入和批量摄入模式。
随着 ML 平台逐渐成熟,有两个运营层面的问题反复出现。首先,运行高吞吐量特征管道的团队必须在循环中调用 PutRecord(用于向在线存储写入单条特征记录)。这意味着每个特征组每条记录需要一次 API 调用,从而产生连接开销和糟糕的吞吐量。一个欺诈检测管道每秒钟摄入 10,000 条记录,跨越五个特征组,就必须维持每秒 50,000 次独立 API 调用,才能保持特征数据的实时更新。第二个挑战是,使用 In-Memory 存储层的团队无法浏览或枚举存储在在线存储中的记录。如果记录标识符因 bug 或管道故障而丢失,这些记录将永久无法恢复。In-Memory 层没有离线存储可以fallback,没有 Amazon Athena 查询可以运行,也没有 API 可以发现其中存储了什么内容。
今天,我们宣布为 Amazon SageMaker Feature Store 推出两个新的 API:
BatchWriteRecord — 在单次 API 调用中写入跨多个特征组最多 25 条记录,支持部分成功语义、每条记录的生命周期(TTL)控制,以及与 PutRecord 相同的基于 EventTime 的顺序保证。
ListRecords — 使用分页枚举特征组内的记录标识符。支持 Standard(Amazon DynamoDB 支持)和 In-Memory(Redis 支持)两种存储层。
在本文中,我们将通过代码示例逐步介绍每个 API,帮助你快速上手。
要跟随本文中的示例,你需要:
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"sagemaker:BatchWriteRecord",
"sagemaker:PutRecord",
"sagemaker:ListRecords"
],
"Resource": "arn:aws:sagemaker:*:*:feature-group/*"
}
]
}
BatchWriteRecord API 解决了单记录摄入的吞吐量限制。以下各节将解释它解决的问题及其工作原理。
单记录摄入的挑战
Feature Store 中现有的 PutRecord API 每次调用向一个特征组写入一条记录。每次调用执行一次条件写入:只有当请求中包含的 EventTime 比现有记录更新时,记录才会作为"最新"版本被持久化。如果条件不满足,记录仍会作为历史版本写入离线存储。
这种设计提供了强顺序保证,但在大规模情况下会强制产生 N×M 的调用模式(N 条记录 × M 个特征组),造成连接开销和尾延迟,限制了吞吐量。
BatchWriteRecord 的工作原理
BatchWriteRecord 在单次请求中接受最多 25 条记录,同时针对一个或多个特征组。每条记录独立成功或失败。这是一个部分成功 API,意味着单条记录失败不会导致整个请求失败。
该 API 保留了与 PutRecord 相同的基于 EventTime 的顺序保证:
未被处理的请求将作为 UnprocessedEntries 在响应中返回,可以重试。
{
"Entries": [
{
"FeatureGroupName": "click-features",
"Record": [
{"FeatureName": "user_id", "ValueAsString": "user-123"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T12:00:00Z"},
{"FeatureName": "click_count", "ValueAsString": "42"}
],
"TargetStores": ["OnlineStore", "OfflineStore"],
"TtlDuration": {"Unit": "Days", "Value": 7}
},
{
"FeatureGroupName": "login-features",
"Record": [
{"FeatureName": "user_id", "ValueAsString": "user-456"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T12:00:01Z"},
{"FeatureName": "login_count", "ValueAsString": "18"}
],
"TargetStores": ["OnlineStore", "OfflineStore"]
}
]
}
响应仅返回失败的记录:
{
"Errors": [
{
"Entry": {
"FeatureGroupName": "string",
"Record": [
{
"FeatureName": "string",
"ValueAsString": "string",
"ValueAsStringList": ["string"]
}
],
"TargetStores": ["string"],
"TtlDuration": {
"Unit": "string",
"Value": number
}
},
"ErrorCode": "string",
"ErrorMessage": "string"
}
],
"UnprocessedEntries": [
{
"FeatureGroupName": "string",
"Record": [
{
"FeatureName": "string",
"ValueAsString": "string",
"ValueAsStringList": ["string"]
}
],
"TargetStores": ["string"],
"TtlDuration": {
"Unit": "string",
"Value": number
}
}
]
}
未在 Errors 或 UnprocessedEntries 中列出的记录均已成功。你的应用应该只使用指数退避重试失败的记录(对于可重试的错误)。
代码示例:使用 Boto3 进行批量摄入
import boto3
featurestore_runtime = boto3.client("sagemaker-featurestore-runtime")
response = featurestore_runtime.batch_write_record(
Entries=[
{
"FeatureGroupName": "click-features",
"Record": [
{"FeatureName": "user_id", "ValueAsString": "user-123"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T12:00:00Z"},
{"FeatureName": "click_count", "ValueAsString": "42"},
],
"TargetStores": ["OnlineStore", "OfflineStore"],
},
{
"FeatureGroupName": "login-features",
"Record": [
{"FeatureName": "user_id", "ValueAsString": "user-456"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T12:00:01Z"},
{"FeatureName": "login_count", "ValueAsString": "18"},
],
"TargetStores": ["OnlineStore", "OfflineStore"],
},
]
)
if response["Errors"]:
for error in response["Errors"]:
print(f"Record {error['Entry']}, ErrorCode: {error['ErrorCode']} Failed: {error['ErrorMessage']}")
if response["UnprocessedEntries"]:
for unprocessed in response["UnprocessedEntries"]:
print(f"Unprocessed: {unprocessed['FeatureGroupName']}")
if not response["Errors"] and not response["UnprocessedEntries"]:
print("All records written successfully.")
代码示例:跨多个特征组写入
你可以在单次请求中针对多个特征组。记录按特征组分组并独立处理:
featurestore_runtime = boto3.client("sagemaker-featurestore-runtime")
```python
response = featurestore_runtime.batch_write_record(
Entries=[
{
"FeatureGroupName": "user-profile-features",
"Record": [
{"FeatureName": "user_id", "ValueAsString": "user-123"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T12:00:00Z"},
{"FeatureName": "age", "ValueAsString": "34"},
{"FeatureName": "region", "ValueAsString": "us-west-2"},
],
"TargetStores": ["OnlineStore"],
},
{
"FeatureGroupName": "click-features",
"Record": [
{"FeatureName": "user_id", "ValueAsString": "user-123"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T12:00:00Z"},
{"FeatureName": "click_count", "ValueAsString": "42"},
],
"TargetStores": ["OnlineStore", "OfflineStore"],
},
]
)
某个特征组的失败不会影响发往其他特征组的记录。
TTL(Time-to-Live,存活时间)支持
BatchWriteRecord 在三个优先级层次上支持 TTL,按以下优先级顺序生效:
记录级 TTL——在单个条目上通过 TtlDuration 设置,优先级最高。
请求级 TTL——在请求顶层设置的默认 TtlDuration,应用于没有记录级 TTL 的条目。
特征组级 TTL——在特征组本身上配置的 TTL,当前两级 TTL 都未设置时生效。
每个请求最多 25 个条目。此限制适用于单个请求中所有特征组的条目总数。
部分成功语义:与事务性 API 不同,BatchWriteRecord 不会在某些记录失败时回滚已成功的写入。设计重试逻辑时,仅重新提交 Errors 中返回的记录。
与 PutRecord 相似的 IAM 模型:调用者必须在每个目标特征组的 ARN 上拥有 sagemaker:BatchWriteRecord 和 sagemaker:PutRecord 权限。会在处理之前检查每个特征组的授权。
EventTime 顺序会被保留:BatchWriteRecord 使用条件写入来维护与 PutRecord 相同的最新记录优先(latest-record-wins)语义。陈旧的记录无法覆盖在线存储中较新的记录。
TargetStores 灵活性:每个条目可以独立地以 OnlineStore、OfflineStore 或两者为目标(默认为特征组启用的存储),让你对每条记录的去向拥有细粒度控制。
ListRecords API 弥补了两种存储层记录发现的空白。以下章节解释它解决的问题及其工作原理。
记录发现面临的挑战
Feature Store 支持 PutRecord、GetRecord 和 DeleteRecord,但所有这些 API 都要求调用者知道确切的记录标识符。没有 API 可以浏览或枚举特征组内的记录。
对于 Standard 层,变通方法是使用 Amazon Athena 查询离线存储。这需要离线存储配置,会增加成本,且不是实时的。
对于 In-Memory 层,情况更为关键。默认情况下没有对应的离线存储。如果记录标识符丢失,这些记录将完全不可恢复。你无法发现它们,也无法删除它们。这导致幽灵数据、存储成本浪费,以及数据主体请求删除时的潜在合规风险。
ListRecords 的工作原理
ListRecords 使用分页枚举特征组内的记录标识符。它仅返回活跃的、未删除的、未过期的、可以与 GetRecord 或 DeleteRecord 配合使用的记录。
该 API 适用于两种存储层:
Standard 层(Amazon DynamoDB):扫描在线存储,返回每条记录最新版本的标识符。软删除和过期记录自动排除。
In-Memory 层(Redis):扫描键并过滤掉软删除记录和内部系统键。从键名中提取记录标识符返回。
请求和响应结构
POST /FeatureGroup/{FeatureGroupName}/ListRecords
{
"MaxResults": 50
}
{
"MaxResults": 50,
"NextToken": "eyJjdXJzb3IiOi4uLn0="
}
{
"RecordIdentifiers": [
"user-001",
"user-002",
"user-003"
],
"NextToken": "eyJuZXh0IjoiLi4ufQ=="
}
当响应中不存在 NextToken 时,分页完成。
代码示例:枚举特征组中的所有记录
import boto3
featurestore_runtime = boto3.client("sagemaker-featurestore-runtime")
all_identifiers = []
next_token = None
while True:
params = {
"FeatureGroupName": "user-profile-features",
"MaxResults": 100,
}
if next_token:
params["NextToken"] = next_token
response = featurestore_runtime.list_records(**params)
all_identifiers.extend(response["RecordIdentifiers"])
next_token = response.get("NextToken")
if not next_token:
break
print(f"Found {len(all_identifiers)} active records.")
代码示例:清理孤立记录
一个常见的用例是识别和删除不再需要的记录。这对于 In-Memory 层的特征组至关重要,因为孤立记录会无限期地保留:
import boto3
featurestore_runtime = boto3.client("sagemaker-featurestore-runtime")
# 步骤 1:枚举所有记录标识符
all_ids = []
next_token = None
while True:
params = {"FeatureGroupName": "session-features", "MaxResults": 100}
if next_token:
params["NextToken"] = next_token
response = featurestore_runtime.list_records(**params)
all_ids.extend(response["RecordIdentifiers"])
next_token = response.get("NextToken")
if not next_token:
break
# 步骤 2:与应用程序的活跃会话列表进行比对
active_sessions = get_active_sessions() # 你的应用程序逻辑
orphaned = [rid for rid in all_ids if rid not in active_sessions]
# 步骤 3:删除孤立记录
for record_id in orphaned:
featurestore_runtime.delete_record(
FeatureGroupName="session-features",
RecordIdentifierValueAsString=record_id,
EventTime="2026-06-05T12:00:00Z",
)
print(f"Deleted {len(orphaned)} orphaned records.")
页面大小:通过 MaxResults 配置(默认 10,最大 100)。
令牌格式:不透明的加密字符串。不要解析或构造令牌,原样传递。
顺序:结果不保证按任何特定顺序返回。
并发写入:如果在分页期间记录被写入或删除,你可能会观察到重复或间隙。这是记录在案的行为。
令牌范围:令牌与特定的特征组和账户绑定,不能跨两者重复使用。
仅返回记录标识符。当前版本仅返回记录标识符,不含特征值。使用 GetRecord 或 BatchGetRecord 检索你所需标识符的完整记录。
自动过滤。该 API 排除软删除记录、过期记录(Standard 层 TTL)和内部系统键(In-Memory 层)。你仅看到活跃的、可检索的记录。
IAM 权限。调用者必须在特征组 ARN 上拥有 sagemaker:ListRecords 权限。
两种存储层都支持。无论特征组使用 Standard 还是 In-Memory 存储,从调用者的角度来看 ListRecords 的行为一致。
这两个 API 自然互补。考虑一个合规工作流,验证在多个特征组中用户的完整数据删除:
import boto3
featurestore_runtime = boto3.client("sagemaker-featurestore-runtime")
feature_groups = ["user-profiles", "click-history", "purchase-signals"]
user_to_delete = "user-789"
# 步骤 1:在所有特征组中查找并删除该用户
for fg_name in feature_groups:
all_ids = []
next_token = None
while True:
params = {"FeatureGroupName": fg_name, "MaxResults": 100}
if next_token:
params["NextToken"] = next_token
response = featurestore_runtime.list_records(**params)
all_ids.extend(response["RecordIdentifiers"])
next_token = response.get("NextToken")
if not next_token:
break
if user_to_delete in all_ids:
featurestore_runtime.delete_record(
FeatureGroupName=fg_name,
RecordIdentifierValueAsString=user_to_delete,
EventTime="2026-06-05T23:59:59Z",
)
print(f"Deleted '{user_to_delete}' from {fg_name}")
featurestore_runtime.batch_write_record(
Entries=[
{
"FeatureGroupName": "deletion-audit-log",
"Record": [
{"FeatureName": "request_id", "ValueAsString": "del-001"},
{"FeatureName": "event_time", "ValueAsString": "2026-06-05T23:59:59Z"},
{"FeatureName": "user_id", "ValueAsString": user_to_delete},
{"FeatureName": "status", "ValueAsString": "completed"},
{"FeatureName": "feature_groups_cleaned", "ValueAsString": "3"},
],
"TargetStores": ["OnlineStore", "OfflineStore"],
}
]
)
为避免持续产生费用,请在完成本演练后删除所创建的特征组。对于 In-Memory 层的特征组,请先使用 ListRecords 枚举记录,然后使用 DeleteRecord 将其删除,再删除特征组。
BatchWriteRecord 和 ListRecords 为 Amazon SageMaker Feature Store 的数据平面提供了关键增强。BatchWriteRecord 通过保持基于 EventTime 的排序保证(确保在线存储正确),将高吞吐量摄取的 API 调用量减少至原来的 1/25。ListRecords 解锁了记录发现和生命周期管理功能。这对于 In-Memory 层的客户至关重要——他们此前无法枚举或清理自己的数据。
这两个 API 共同支持了此前难以实现或无法实现的模式:减少连接数并降低延迟的批量摄取管道、能够验证数据完全删除的合规工作流,以及能够实时浏览特征组内容的运维工具。
更多信息请参阅 Feature Store 文档、Feature Store API 参考、离线存储配置文档以及 What's New 公告。
有关 Feature Store 能力的背景信息,请探索以下相关文章: