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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 17:51:45