使用MWAA处理S3中10+GB每日CSV新增列失败求助
解决S3大CSV文件新增列并回存的问题
核心问题分析
直接用object['Body'].read().decode('utf-8')会把10+GB的完整文件加载到内存,而MWAA的worker节点内存资源有限,必然导致任务因内存不足失败。必须改用流式分块读取-处理-写入的方式,避免全量加载。
具体解决方案
1. 流式分块处理+S3分块上传
通过逐行读取CSV、批量写入分块的方式,控制内存占用,同时用S3分块上传高效写入大文件:
import boto3 import csv from io import StringIO, BytesIO s3 = boto3.client('s3') bucket_name = '你的存储桶名称' input_key = '原始文件路径/xxx.csv' output_key = '处理后文件路径/xxx_processed.csv' # 流式读取原始文件 response = s3.get_object(Bucket=bucket_name, Key=input_key) body_stream = response['Body'] # 初始化CSV读取器,先处理表头 header_line = body_stream.readline().decode('utf-8') csv_reader = csv.reader(StringIO(header_line)) header = next(csv_reader) header.append('新增列名') # 替换为你的列名 # 启动S3分块上传 mpu = s3.create_multipart_upload(Bucket=bucket_name, Key=output_key) parts = [] part_number = 1 chunk_buffer = BytesIO() # 先写入处理后的表头 chunk_buffer.write(','.join(header).encode('utf-8') + b'\n') # 逐行处理数据,满阈值就上传分块 for line in body_stream: row = next(csv.reader(StringIO(line.decode('utf-8')))) row.append('新增列值') # 替换为你的列值逻辑 chunk_buffer.write(','.join(row).encode('utf-8') + b'\n') # 当缓冲区达到100MB时上传分块(可根据内存调整阈值) if chunk_buffer.tell() >= 100 * 1024 * 1024: chunk_buffer.seek(0) part = s3.upload_part( Bucket=bucket_name, Key=output_key, UploadId=mpu['UploadId'], PartNumber=part_number, Body=chunk_buffer ) parts.append({'PartNumber': part_number, 'ETag': part['ETag']}) part_number += 1 chunk_buffer = BytesIO() # 上传剩余的最后一块数据 if chunk_buffer.tell() > 0: chunk_buffer.seek(0) part = s3.upload_part( Bucket=bucket_name, Key=output_key, UploadId=mpu['UploadId'], PartNumber=part_number, Body=chunk_buffer ) parts.append({'PartNumber': part_number, 'ETag': part['ETag']}) # 完成分块上传 s3.complete_multipart_upload( Bucket=bucket_name, Key=output_key, UploadId=mpu['UploadId'], MultipartUpload={'Parts': parts} )
2. MWAA环境优化
- 升级worker节点实例规格:选择内存更大的实例(如m5.xlarge及以上),给分块处理足够的内存空间。
- 延长任务超时时间:在Airflow DAG中设置
execution_timeout=timedelta(hours=4)(根据实际处理时长调整),避免任务被提前终止。
3. 简化实现的可选方案
使用smart_open库简化S3流式读写逻辑,需在MWAA的requirements.txt中添加smart_open[s3]依赖:
import csv from smart_open import open bucket_name = '你的存储桶名称' input_path = f's3://{bucket_name}/原始文件路径/xxx.csv' output_path = f's3://{bucket_name}/处理后文件路径/xxx_processed.csv' with open(input_path, 'r') as infile, open(output_path, 'w') as outfile: reader = csv.reader(infile) writer = csv.writer(outfile) # 处理表头 header = next(reader) header.append('新增列名') writer.writerow(header) # 逐行处理数据 for row in reader: row.append('新增列值') writer.writerow(row)
内容的提问来源于stack exchange,提问作者Abhi
相关产品推荐
相关产品推荐

