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

如何在Airflow中通过DAG解压S3上的ZIP文件且无需下载?

如何在Airflow中通过DAG解压S3上的ZIP文件且无需下载?

别发愁啦,完全可以做到不用把ZIP文件下载到本地,直接在云端完成解压操作!我给你两种实用的方案,你可以根据自己的场景来选:

方案一:用PythonOperator自定义逻辑(适合Airflow内直接控制流程)

如果想在Airflow DAG里直接实现解压逻辑,用PythonOperator配合boto3和zipfile就能搞定——核心是用流处理,不把整个文件写到本地磁盘,全程在内存里操作。

下面是完整的DAG示例代码:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
import boto3
from zipfile import ZipFile
from io import BytesIO

def unzip_s3_file(**context):
    # 初始化S3客户端
    s3 = boto3.client('s3')
    
    # 这里替换成你的实际配置
    source_bucket = 'your-source-bucket-name'
    source_zip_key = 'path/to/your/target/file.zip'  # S3上zip文件的完整路径
    target_bucket = 'your-target-bucket-name'  # 可以和源bucket相同
    target_prefix = 'unzipped_results/'  # 解压后文件存放的前缀路径

    # 以流的方式读取S3上的ZIP文件,不下载到本地
    try:
        zip_obj = s3.get_object(Bucket=source_bucket, Key=source_zip_key)
        zip_stream = BytesIO(zip_obj['Body'].read())
        
        # 读取ZIP内容并逐个解压上传
        with ZipFile(zip_stream, 'r') as zip_ref:
            for file_name in zip_ref.namelist():
                # 跳过ZIP里的目录(如果有的话)
                if file_name.endswith('/'):
                    continue
                # 读取单个文件的内容
                file_content = zip_ref.read(file_name)
                # 构造目标文件的S3路径
                target_file_key = f"{target_prefix}{file_name}"
                # 上传到目标S3位置
                s3.put_object(
                    Bucket=target_bucket,
                    Key=target_file_key,
                    Body=file_content
                )
        print(f"解压成功!所有文件已上传到 s3://{target_bucket}/{target_prefix}")
    except Exception as e:
        print(f"解压过程出错:{str(e)}")
        raise  # 抛出异常让Airflow标记任务失败

# 定义DAG
with DAG(
    dag_id='s3_unzip_dag',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,  # 按需手动触发,也可以设置定时比如'@daily'
    catchup=False,
    tags=['s3', 'unzip']
) as dag:
    unzip_task = PythonOperator(
        task_id='unzip_s3_zip_file',
        python_callable=unzip_s3_file,
        provide_context=True
    )

unzip_task

方案二:用LambdaOperator触发AWS Lambda处理(适合大文件场景)

如果你的ZIP文件很大(比如几十GB),Airflow Worker的内存可能扛不住,这时候用AWS Lambda来处理更合适——Lambda可以调整内存(最高10GB)和超时时间(最高15分钟),而且天然支持S3对象的流操作。

步骤1:创建Lambda函数

先在AWS Lambda里写一个处理解压的函数,逻辑和上面类似:

import boto3
from zipfile import ZipFile
from io import BytesIO

s3 = boto3.client('s3')

def lambda_handler(event, context):
    # 从Airflow传递的参数里获取配置
    source_bucket = event['source_bucket']
    source_zip_key = event['source_zip_key']
    target_bucket = event['target_bucket']
    target_prefix = event['target_prefix']

    try:
        zip_obj = s3.get_object(Bucket=source_bucket, Key=source_zip_key)
        zip_stream = BytesIO(zip_obj['Body'].read())
        
        with ZipFile(zip_stream, 'r') as zip_ref:
            for file_name in zip_ref.namelist():
                if file_name.endswith('/'):
                    continue
                file_content = zip_ref.read(file_name)
                target_key = f"{target_prefix}{file_name}"
                s3.put_object(Bucket=target_bucket, Key=target_key, Body=file_content)
        
        return {
            'statusCode': 200,
            'body': f"成功解压:s3://{source_bucket}/{source_zip_key} → s3://{target_bucket}/{target_prefix}"
        }
    except Exception as e:
        return {
            'statusCode': 500,
            'body': f"解压失败:{str(e)}"
        }

记得给Lambda的执行角色配置S3的GetObject和PutObject权限哦!

步骤2:在Airflow DAG里触发Lambda

用LambdaInvokeOperator来调用上面的Lambda函数:

from airflow import DAG
from airflow.providers.amazon.aws.operators.lambda_function import LambdaInvokeOperator
from datetime import datetime

with DAG(
    dag_id='s3_unzip_via_lambda_dag',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False,
    tags=['s3', 'lambda', 'unzip']
) as dag:
    trigger_lambda_unzip = LambdaInvokeOperator(
        task_id='trigger_s3_unzip_lambda',
        function_name='your-lambda-function-name',  # 替换成你的Lambda函数名
        payload={
            'source_bucket': 'your-source-bucket-name',
            'source_zip_key': 'path/to/your/file.zip',
            'target_bucket': 'your-target-bucket-name',
            'target_prefix': 'unzipped_results/'
        },
        aws_conn_id='aws_default'  # 你的Airflow AWS连接ID
    )

trigger_lambda_unzip

一些重要的注意事项

  • 权限配置:不管用哪种方案,执行角色(Airflow Worker角色或Lambda角色)都必须拥有S3源桶的GetObject权限、目标桶的PutObject权限,必要时还要加上ListBucket权限来验证文件存在性。
  • 大文件优化:如果处理超大ZIP文件,用Lambda时记得调大内存(内存越大,CPU和网络带宽也会越高),同时设置足够的超时时间;用Airflow PythonOperator的话,要确保Worker节点有足够的内存,避免内存溢出。
  • 错误处理:可以在代码里加上针对性的异常捕获(比如zipfile.BadZipFile处理损坏的ZIP、s3.exceptions.NoSuchKey处理文件不存在),还能在Airflow里配置on_failure_callback来发送告警邮件或消息。

备注:内容来源于stack exchange,提问作者ennezetaqu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 17:58:03