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

DynamoDB的put_item接口是否存在QPS限制?异步迁移遇吞吐量瓶颈

嘿,我刚好遇到过类似的场景,来给你捋清楚为啥会这样,再给你补个完整的可运行代码示例~

核心原因:单条写入 vs 批量写入的效率差异

DynamoDB的API设计决定了两种写入方式的吞吐量上限:

  • put_item:每次请求只能处理1条记录,哪怕用异步框架,频繁的单个请求会带来大量网络往返开销、连接管理成本,而且DynamoDB单条请求的QPS上限会直接限制整体写入速度,想冲到1000的吞吐量自然很困难。
  • batch_write_item:一次请求最多能处理25条记录(总大小不超过16MB),相当于把多个写入请求合并成一个,大幅减少了请求次数,网络开销骤降,能更高效地利用DynamoDB的写入容量,轻松达到预期的吞吐量目标。
完整的aiobotocore批量写入实现代码

我把你没写完的代码补全,加上批量拆分、未处理项重试等关键逻辑:

import aiobotocore
import asyncio
import json
from itertools import islice

# 从JSON文件加载数据的工具函数
def load_records_from_json(file_path):
    with open(file_path, 'r', encoding='utf-8') as f:
        return json.load(f)

async def batch_write_to_dynamodb(client, table_name, records, batch_size=25):
    records_iter = iter(records)
    while True:
        # 按批次拆分数据,每次取batch_size条
        current_batch = list(islice(records_iter, batch_size))
        if not current_batch:
            break
        
        # 构造batch_write_item的请求格式
        request_payload = {
            table_name: [
                {'PutRequest': {'Item': record}} 
                for record in current_batch
            ]
        }
        
        try:
            response = await client.batch_write_item(RequestItems=request_payload)
            # 处理DynamoDB返回的未写入项(可能因为限流等原因)
            unprocessed = response.get('UnprocessedItems', {}).get(table_name, [])
            if unprocessed:
                # 把未处理的项放回迭代器头部,下次循环继续处理
                records_iter = iter([item['PutRequest']['Item'] for item in unprocessed] + list(records_iter))
                # 加个短暂延迟,避免立即重试触发更严的限流
                await asyncio.sleep(0.1)
        except Exception as e:
            print(f"Batch write failed: {str(e)}")
            # 异常情况下延迟重试,可根据需求调整为指数退避
            await asyncio.sleep(0.5)

async def main():
    # 替换成你的实际配置
    TABLE_NAME = "your-target-table"
    JSON_FILE_PATH = "your-data-file.json"
    
    # 加载待写入的记录
    records = load_records_from_json(JSON_FILE_PATH)
    print(f"已加载 {len(records)} 条待写入记录")
    
    # 创建异步DynamoDB客户端(建议用IAM角色或环境变量存凭证,不要硬编码)
    session = aiobotocore.get_session()
    async with session.create_client(
        "dynamodb",
        region_name="us-east-1",  # 替换成你的区域
        # aws_access_key_id="your-key",
        # aws_secret_access_key="your-secret"
    ) as dynamodb_client:
        await batch_write_to_dynamodb(dynamodb_client, TABLE_NAME, records)
    
    print("所有记录已成功写入DynamoDB!")

if __name__ == "__main__":
    asyncio.run(main())
几个关键注意事项
  • 批次大小限制:如果你的单条记录体积较大(比如接近640KB),要把batch_size调小,确保整个批次的总大小不超过16MB的上限。
  • 未处理项重试:代码里已经做了基础的重试逻辑,如果你需要更稳健的处理,可以改成指数退避(比如第一次等0.1s,第二次0.2s,第三次0.4s...)。
  • 并发优化:如果数据量极大,可以用asyncio.gather同时启动多个批量写入任务,但要注意不要超过DynamoDB表的写入容量单位(WCU),避免被限流。
  • 凭证安全:永远不要硬编码AWS凭证,优先用IAM角色(比如EC2/EKS实例角色)、环境变量或者AWS配置文件来管理凭证。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:28:47