Amazon SageMaker Feature Store 新增了两个面向数据写入与发现的 API:BatchWriteRecord 可以在一次调用中向多个 Feature Group 写入最多 25 条记录,ListRecords 则可以枚举某个 Feature Group 中的记录标识符。两者结合后,离线导入、数据校验和特征数据盘点都不必再依赖大量单条 API 调用。
BatchWriteRecord:一次写入多个 Feature Group
过去,应用通常需要循环调用单条写入接口。数据量稍大时,这种方式会带来更多网络往返,也让客户端重试、限流和错误收集变得复杂。
BatchWriteRecord 的核心能力是把多个记录放进一次请求中,并且允许这些记录属于不同的 Feature Group。单次请求最多包含 25 条记录。每条记录仍然由特征名、特征值等字段组成,因此调用方需要提前保证字段类型和 Feature Group schema 一致。
一个典型的请求可以包含两个 Feature Group:
import boto3
runtime = boto3.client("sagemaker-featurestore-runtime", region_name="us-east-1")
response = runtime.batch_write_record(
Records=[
{
"FeatureGroupName": "customer-profile",
"Record": [
{"FeatureName": "customer_id", "ValueAsString": "c-1001"},
{"FeatureName": "country", "ValueAsString": "CN"},
{"FeatureName": "signup_days", "ValueAsString": "180"},
],
},
{
"FeatureGroupName": "customer-activity",
"Record": [
{"FeatureName": "customer_id", "ValueAsString": "c-1001"},
{"FeatureName": "orders_30d", "ValueAsString": "12"},
{"FeatureName": "last_active_at", "ValueAsString": "2025-01-15T10:30:00Z"},
],
},
]
)
print(response)
运行前请确认:
- AWS 凭证已经通过 IAM role、环境变量或 AWS profile 配置;
customer-profile和customer-activity已经存在;- 记录中的特征名与 Feature Group schema 匹配;
- 当前安装的
boto3版本已经包含这些 API。
可以先升级 SDK,并检查客户端是否暴露了对应方法:
python -m pip install --upgrade boto3 botocore
python - <<'PY'
import boto3
client = boto3.client("sagemaker-featurestore-runtime", region_name="us-east-1")
print("batch_write_record:", hasattr(client, "batch_write_record"))
print("list_records:", hasattr(client, "list_records"))
PY
真实导入任务通常需要把输入切成 25 条一批,并保留失败记录。下面是一个可改造的批处理骨架:
from itertools import islice
import boto3
def chunks(items, size):
iterator = iter(items)
while batch := list(islice(iterator, size)):
yield batch
def write_records(records, region="us-east-1"):
client = boto3.client("sagemaker-featurestore-runtime", region_name=region)
failures = []
for batch in chunks(records, 25):
result = client.batch_write_record(Records=batch)
# 具体失败字段以当前 SDK 返回结构为准;保留完整响应便于审计和重试。
if result.get("Errors"):
failures.extend(result["Errors"])
return failures
records = [
{
"FeatureGroupName": "customer-profile",
"Record": [
{"FeatureName": "customer_id", "ValueAsString": f"c-{i}"},
{"FeatureName": "country", "ValueAsString": "CN"},
],
}
for i in range(1, 51)
]
failed = write_records(records)
print(f"failed records: {len(failed)}")
这里的重点不是简单地把循环搬进一个函数,而是让批量写入具备可观测性:记录批次边界、保存失败响应,并对失败记录执行有上限的重试。对于会被重复执行的任务,还应设计稳定的记录标识和幂等策略。
ListRecords:先发现记录,再决定读取范围
ListRecords 用于枚举一个 Feature Group 中的记录标识符。它适合做数据盘点、迁移前检查、回填任务准备和测试环境清理等工作。
它与读取完整特征数据是两个不同动作:列表接口帮助你发现有哪些记录,后续是否读取每条记录的全部特征,应根据业务需要决定。对于大 Feature Group,不要一次性把所有结果放进内存,而应使用分页 token。
下面的示例展示了一个分页读取方式:
import boto3
def list_all_records(feature_group_name, region="us-east-1", page_size=100):
client = boto3.client("sagemaker-featurestore-runtime", region_name=region)
next_token = None
while True:
params = {
"FeatureGroupName": feature_group_name,
"MaxResults": page_size,
}
if next_token:
params["NextToken"] = next_token
response = client.list_records(**params)
for record in response.get("Records", []):
yield record
next_token = response.get("NextToken")
if not next_token:
break
for item in list_all_records("customer-profile"):
print(item)
不同版本的 SDK 可能对返回对象的字段暴露方式有所差异。接入前应通过当前版本的 boto3 文档或 client.list_records 的服务模型确认请求参数和返回字段;分页逻辑本身应保留,因为列表接口通常不会保证一次返回全部记录。
两个 API 放进数据流水线时要注意什么
批量大小不等于无限吞吐
单次最多 25 条记录减少了网络往返,但它不代表可以无节制地并发发送请求。生产任务仍需要根据服务限额、记录大小和错误率控制并发,并采用指数退避处理临时性错误。
失败处理要细到记录级别
批次中出现失败记录时,不要只把整个批次标记为失败后盲目重放。应保存请求批次、Feature Group、记录主键和服务返回的错误信息,再只重试可重试的记录。永久性 schema 错误则应该进入隔离队列,而不是持续重试。
发现记录不等于稳定排序
ListRecords 的目标是枚举标识符。不要把返回顺序当作业务排序,也不要假设分页结果在并发写入期间具有完整的一致性视图。如果需要可重复的数据快照,应在上游生成明确的批次边界或使用专门的数据导出流程。
落地清单
可以按下面的顺序接入:
- 升级并锁定包含新 API 的 boto3/botocore 版本。
- 用少量测试数据验证跨 Feature Group 的批量请求。
- 把输入切分为不超过 25 条的请求批次。
- 保存批次响应和记录级错误,设置有限次数的重试。
- 为
ListRecords实现 token 分页,不把全量结果一次性载入内存。 - 为导入任务补充指标,例如成功数、失败数、重试数和处理耗时。
BatchWriteRecord 适合降低多条写入的调用开销,ListRecords 适合让 Feature Group 的内容变得可发现。两者结合时,最重要的工程工作仍然是 schema 校验、分页、限流和失败恢复,而不是仅仅替换一个 API 名称。