如何无需按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
相关产品推荐
相关产品推荐

