如何使用Python优化5GB CSV文件到AWS DynamoDB的导入速度?
5GB CSV导入DynamoDB Python优化方案
基础前提调整
- 先临时提升DynamoDB表的写入容量单元(WCU):按5GB数据、平均单条大小1KB估算,2小时完成导入需要至少700WCU,你可以临时切换到按需模式,或者把预配置WCU调整到1000以上,导入完成后再调回原有配置即可,额外成本极低。
- 更换运行环境:你参考的示例代码是Lambda运行环境,Lambda最长运行时间仅15分钟,完全无法支撑数小时的导入任务,直接换到EC2、ECS任务或者本地性能足够的机器上运行。
代码核心优化点(无需拆分原CSV文件)
- 改用多进程并行写入:原代码是单线程串行写入,是速度慢的核心原因。因为boto3操作受GIL锁限制,多线程无法提升效率,直接用多进程实现:将csv读取生成的行写入多进程安全的队列,启动对应CPU核数2倍的消费进程,从队列拉取数据做批量写入,充分利用带宽和DynamoDB的写入配额。
- 移除冗余操作:原代码每次调用
write_to_dynamo都重复初始化DynamoDB表对象,会产生大量不必要的API请求开销,将表初始化操作提到全局,仅初始化一次即可。 - 优化boto3重试配置:新增自适应重试配置,应对DynamoDB写入时的偶发节流,避免无效等待。初始化dynamodb资源时使用如下配置:
from botocore.config import Config dynamodb = boto3.resource('dynamodb', config=Config( retries = {'max_attempts': 10, 'mode': 'adaptive'} ))
- 优化S3流式读取:原代码直接读取S3对象流的逻辑在处理大文件时容易卡顿,可改用分块拉取S3对象的逻辑,无需拆分原文件,内存占用更低,读取速度更快。
- 调整批量大小:如果你的CSV单行数据较小,可将批量攒批的大小从100调整到200,
batch_writer底层会自动按照DynamoDB单次批量写入最多25条的限制拆分请求,不会触发接口限制。
优化后核心代码示例
import json import boto3 import os import csv import codecs from multiprocessing import Process, Queue from botocore.config import Config # 全局初始化资源,仅执行一次 s3 = boto3.resource('s3') dynamodb = boto3.resource('dynamodb', config=Config( retries = {'max_attempts': 10, 'mode': 'adaptive'} )) bucket = os.environ['bucket'] key = os.environ['key'] tableName = os.environ['table'] table = dynamodb.Table(tableName) BATCH_SIZE = 100 WORKER_NUM = 8 # 根据CPU核数调整,一般为核数的2倍 # 消费者进程:从队列拿数据批量写入DynamoDB def write_worker(queue): batch = [] while True: row = queue.get() if row is None: # 结束标记 if batch: with table.batch_writer() as writer: for item in batch: writer.put_item(Item=item) break batch.append(row) if len(batch) >= BATCH_SIZE: with table.batch_writer() as writer: for item in batch: writer.put_item(Item=item) batch.clear() def main(): # 初始化队列和工作进程 queue = Queue(maxsize=WORKER_NUM * BATCH_SIZE * 2) workers = [] for _ in range(WORKER_NUM): p = Process(target=write_worker, args=(queue,)) p.start() workers.append(p) # 流式读取CSV,不加载全量文件到内存 obj = s3.Object(bucket, key).get()['Body'] for row in csv.DictReader(codecs.getreader('utf-8')(obj)): queue.put(row) # 发送结束标记 for _ in range(WORKER_NUM): queue.put(None) # 等待所有进程完成 for p in workers: p.join() if __name__ == "__main__": main()
相关知识学习建议
- 先掌握DynamoDB的核心基础:包括容量模式、读写配额规则、批量操作的接口限制,这是所有DynamoDB操作优化的前提。
- 学习Python并发编程知识:区分IO密集型和CPU密集型任务的优化逻辑,掌握多进程、队列的使用方法,了解GIL锁对并发操作的影响。
- 熟悉boto3的进阶配置:包括重试策略、客户端复用、会话管理等特性,可以大幅降低AWS服务调用的额外开销。
内容的提问来源于stack exchange,提问作者user37263
相关产品推荐
相关产品推荐

