如何在AWS MWAA中直接操作S3内Excel文件实现ETL数据加载?
解决方案:在AWS MWAA中直接操作S3上的Excel文件
核心思路是通过**内存流(BytesIO)**结合boto3操作S3,避免本地下载/上传文件的步骤:直接将S3中的Excel文件读入内存,用openpyxl修改后再写回S3。
步骤1:确认依赖
确保MWAA环境的requirements.txt中包含openpyxl(boto3是MWAA默认预装的):
openpyxl>=3.1.2
步骤2:修改ETL加载函数
将原本地文件操作改为内存流+S3操作,完整代码如下:
import boto3 from io import BytesIO import openpyxl from airflow.models import TaskInstance def export_andbeyond_to_excel(ti: TaskInstance): df = ti.xcom_pull(task_ids='transform_andbeyond') s3_file_path = ti.xcom_pull(task_ids='latest_file') # 格式应为s3://bucket-name/path/to/file.xlsx # 解析S3路径,拆分bucket和文件key s3_parts = s3_file_path.replace("s3://", "").split("/", 1) bucket_name = s3_parts[0] file_key = s3_parts[1] # 初始化S3客户端 s3_client = boto3.client('s3') # 从S3读取文件到内存流 response = s3_client.get_object(Bucket=bucket_name, Key=file_key) excel_data = response['Body'].read() master_wb = openpyxl.load_workbook(BytesIO(excel_data)) master_sheet = master_wb.active # 原逻辑保持不变 table = master_sheet['Q71':'Q89'] for i, row in enumerate(table, start=70): site = row[0].value if site in df['Sites'].values: index = df.index[df['Sites'] == site].tolist()[0] master_sheet[f'R{i+1}'] = df.at[index, 'Impressions'] master_sheet[f'S{i+1}'] = df.at[index, 'Revenue'] else: master_sheet[f'R{i+1}'] = 0 master_sheet[f'S{i+1}'] = 0 # 将修改后的工作簿写入内存流,再上传回S3 output_stream = BytesIO() master_wb.save(output_stream) output_stream.seek(0) # 重置流指针到开头 s3_client.put_object(Bucket=bucket_name, Key=file_key, Body=output_stream)
关键注意事项
- 权限配置:确保MWAA的执行角色拥有目标S3桶的
s3:GetObject和s3:PutObject权限,可在IAM角色的政策中添加对应权限。 - 路径格式检查:
latest_file返回的必须是标准S3路径(如s3://my-bucket/excel/master_file.xlsx),若返回的是其他格式(仅key),需调整解析逻辑。 - 内存占用:若Excel文件过大,需注意MWAA任务的内存配置,避免内存溢出。
内容的提问来源于stack exchange,提问作者Zulkifli Arshad
相关产品推荐
相关产品推荐

