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

如何无需按100条批量导入数据至AWS Timestream?

将1亿条代币历史价格批量导入AWS Timestream的高效方案

问题描述

我需用Python的boto3将1亿条代币历史价格数据存入AWS Timestream,每条记录包含token_address、timestamp、price三个字段,时间范围覆盖2019年至今,所有数据存储在CSV文件中。按照AWS配额限制,单批次写入最多仅支持100条记录,官方文档给出的优化方案是「写入具有共同属性的批量记录」,但我的记录仅结构一致、值均不同,无法定义common_attributes。想咨询是否存在无需按100条逐批导入的替代方法?

可行解决方案

虽然Timestream单批次写入的100条硬限制无法突破,但可以通过以下方式大幅提升导入效率,避免手动低效处理小批量数据:

1. 并行批量写入

借助Python多线程/多进程或异步框架(如asyncio+aioboto3)同时发起多个写入请求,通过并发提升整体吞吐量。例如开启10-20个并发线程,每个线程独立处理100条的批次:

import boto3
from concurrent.futures import ThreadPoolExecutor

# 初始化Timestream客户端
client = boto3.client('timestream-write')
DB_NAME = 'your_token_db'
TABLE_NAME = 'token_price_table'

def write_batch(batch):
    try:
        client.write_records(
            DatabaseName=DB_NAME,
            TableName=TABLE_NAME,
            Records=batch
        )
    except client.exceptions.RejectedRecordsException as e:
        # 记录并处理被拒绝的记录,比如重试或存入错误日志
        print(f"被拒绝记录详情: {e.response['RejectedRecords']}")

# 从CSV读取数据并拆分100条批次
batches = []
current_batch = []
with open('token_historical_prices.csv', 'r') as f:
    next(f)  # 跳过表头
    for line in f:
        token_addr, ts, price = line.strip().split(',')
        # 构造Timestream要求的记录格式
        record = {
            'Dimensions': [{'Name': 'token_address', 'Value': token_addr}],
            'MeasureName': 'price',
            'MeasureValue': price,
            'MeasureValueType': 'DOUBLE',
            'Time': str(int(float(ts) * 1000))  # 转换为毫秒级时间戳
        }
        current_batch.append(record)
        if len(current_batch) == 100:
            batches.append(current_batch)
            current_batch = []
if current_batch:  # 处理最后一批不足100条的记录
    batches.append(current_batch)

# 用线程池并行执行写入
with ThreadPoolExecutor(max_workers=15) as executor:
    executor.map(write_batch, batches)

2. 使用Timestream批量导入工具

如果CSV文件已上传至S3,可直接使用Timestream批量导入功能,无需自行编写分批逻辑,该工具会自动处理批次和并发,适配亿级规模数据导入:

  • 将CSV文件上传至S3存储桶
  • 登录Timestream控制台,创建导入任务,指定S3路径、目标数据库和表
  • 配置数据映射规则:将token_address设为维度,timestamp设为时间列,price设为度量值
  • 启动任务后,Timestream会自动完成全量数据导入

3. 优化写入重试策略

针对Timestream的限流异常(ThrottlingException),可自定义指数退避重试逻辑,避免限流拖慢导入进度。boto3支持自定义重试配置:

from botocore.config import Config

# 配置自适应重试策略,最大重试10次
config = Config(
    retries={
        'max_attempts': 10,
        'mode': 'adaptive'
    }
)
client = boto3.client('timestream-write', config=config)

注意事项

  • 使用批量导入工具时,需确保CSV的时间戳格式符合Timestream要求(支持UNIX秒/毫秒/微秒级时间戳,需在导入配置中明确指定)
  • 并行写入时,避免过度开启并发数,建议从低并发测试后逐步调整,防止触发账户级限流
  • 务必记录并处理被拒绝的记录,避免数据丢失

内容的提问来源于stack exchange,提问作者Vince M

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 04:01:14