如何用Python boto3分批次读取并处理S3中的大型CSV文件
解决方案
核心思路
避免一次性将整个CSV文件加载到内存,利用S3对象的流式特性结合Python的csv模块,逐行解析并分批处理数据,全程在内存中完成,无需写入磁盘。
具体实现
1. 流式读取CSV并分批处理
S3返回的data['Body']是一个可迭代的字节流对象,我们可以用codecs将其包装为UTF-8文本流,直接传给csv.DictReader,让它逐行解析,而非一次性加载全部内容。再通过自定义生成器控制每批处理的行数(比如1万行)。
import csv import codecs # 假设你已经初始化了S3客户端s3,以及配置config def process_batch(batch): # 替换为你的实际处理逻辑,比如数据清洗、分析等 for row in batch: # 示例:处理单条数据(根据需求修改) pass def batch_generator(reader, batch_size=10000): batch = [] for row in reader: batch.append(row) if len(batch) >= batch_size: yield batch batch = [] # 处理最后一批不足batch_size的数据 if batch: yield batch # 获取S3对象 data = s3.get_object(Bucket=config.BUCKET_NAME, Key=source_file) # 将字节流包装为UTF-8文本流 text_stream = codecs.getreader('utf-8')(data['Body']) # 创建DictReader,自动读取表头并生成字典格式的行数据 csv_reader = csv.DictReader(text_stream) # 分批次处理数据 for batch in batch_generator(csv_reader, batch_size=10000): process_batch(batch)
2. 方案优势说明
- 流式读取:
data['Body']会逐步从S3拉取数据,而非一次性下载整个文件,内存占用仅与当前批次的行数相关。 - 行边界安全:
csv.DictReader会自动处理CSV的行分隔符,不会出现download_fileobj按字节分块导致的行割裂问题。 - 无磁盘依赖:全程在内存中处理数据,无需临时文件,解决磁盘空间不足的问题。
对比原方案的优化点
原方案中contents = data['Body'].read().decode("utf-8")会将整个文件加载到内存,大文件直接触发内存错误;而新方案通过流式解析,内存占用可严格控制在批次大小范围内,适配任意大小的CSV文件。
内容的提问来源于stack exchange,提问作者eduardosufan
相关产品推荐
相关产品推荐

