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

寻求Airflow GCSToS3Operator代码示例及配置要求说明

Airflow实现GCS到S3的文件传输:代码示例与配置指南

一、前置配置要求

  • 安装依赖包:确保Airflow环境安装了GCP和AWS的官方Provider包,执行以下命令:
    pip install apache-airflow-providers-google apache-airflow-providers-amazon
    
  • 配置GCP连接:
    在Airflow UI的「Admin > Connections」页面,创建类型为「Google Cloud」的连接:
    • 连接ID建议设为google_cloud_default(自定义ID需与代码中对应)
    • 上传GCP服务账号密钥JSON文件,该账号需拥有目标GCS存储桶的storage.objects.get权限
  • 配置AWS连接:
    同样在Connections页面,创建类型为「Amazon Web Services」的连接:
    • 连接ID建议设为aws_default(自定义ID需与代码中对应)
    • 填入AWS的Access Key ID和Secret Access Key,该账号需拥有目标S3存储桶的s3:PutObject权限
  • 网络访问:确保Airflow调度器和Worker能访问GCS与S3服务(公网环境默认支持,私有网络需配置VPC端点等)

二、代码实现示例

1. 使用官方GCSToS3Operator(推荐,适合单文件/批量前缀匹配)

这是Airflow官方提供的专用传输Operator,代码简洁易维护:

from airflow import DAG
from airflow.providers.google.cloud.transfers.gcs_to_s3 import GCSToS3Operator
from datetime import datetime, timedelta

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'gcs_to_s3_transfer',
    default_args=default_args,
    description='Transfer files from GCS to S3',
    schedule_interval=timedelta(days=1),
    start_date=datetime(2023, 1, 1),
    catchup=False,
    tags=['gcs', 's3', 'transfer'],
) as dag:

    # 单文件传输任务
    transfer_single_file = GCSToS3Operator(
        task_id='transfer_single_file',
        gcp_conn_id='google_cloud_default',
        aws_conn_id='aws_default',
        gcs_bucket='your-gcs-bucket-name',
        gcs_object='data/source/sample.csv',  # GCS中的源文件路径
        s3_bucket='your-s3-bucket-name',
        s3_key='data/destination/sample.csv',  # S3中的目标文件路径
        replace=True,  # 允许覆盖S3中已存在的同名文件
    )

    # 批量传输任务(匹配GCS前缀下的所有文件)
    transfer_batch_files = GCSToS3Operator(
        task_id='transfer_batch_files',
        gcp_conn_id='google_cloud_default',
        aws_conn_id='aws_default',
        gcs_bucket='your-gcs-bucket-name',
        gcs_object='data/source/',  # GCS中的文件前缀路径
        s3_bucket='your-s3-bucket-name',
        s3_key='data/destination/',  # S3中的目标前缀路径
        replace=True,
        is_gcs_object_prefix=True,  # 启用前缀匹配模式
    )

    transfer_single_file >> transfer_batch_files

2. 自定义PythonOperator(适合复杂逻辑,比如过滤文件、修改文件名)

如果需要灵活处理文件(如按日期过滤、重命名),可以用Airflow Hook自定义传输逻辑:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.google.cloud.hooks.gcs import GCSHook
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from datetime import datetime, timedelta

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

def transfer_gcs_to_s3_with_logic(**context):
    # 初始化Hook,复用Airflow配置的连接
    gcs_hook = GCSHook(gcp_conn_id='google_cloud_default')
    s3_hook = S3Hook(aws_conn_id='aws_default')

    # 配置传输参数
    gcs_bucket = 'your-gcs-bucket-name'
    gcs_prefix = 'data/source/'
    s3_bucket = 'your-s3-bucket-name'
    s3_prefix = 'data/destination/'

    # 获取GCS前缀下的所有文件
    file_list = gcs_hook.list(bucket_name=gcs_bucket, prefix=gcs_prefix)

    for file_path in file_list:
        # 跳过前缀目录本身
        if file_path == gcs_prefix:
            continue
        # 仅传输CSV文件(自定义过滤逻辑)
        if not file_path.endswith('.csv'):
            continue
        # 读取GCS文件内容
        file_content = gcs_hook.download(bucket_name=gcs_bucket, object_name=file_path)
        # 构造S3目标路径(添加日期前缀)
        current_date = context['ds']
        s3_key = f"{s3_prefix}{current_date}/{file_path.split('/')[-1]}"
        # 上传到S3
        s3_hook.load_bytes(
            bytes_data=file_content,
            key=s3_key,
            bucket_name=s3_bucket,
            replace=True
        )
        print(f"Transferred: {file_path} -> {s3_key}")

with DAG(
    'gcs_to_s3_custom_transfer',
    default_args=default_args,
    description='Custom GCS to S3 transfer with filtering logic',
    schedule_interval=timedelta(days=1),
    start_date=datetime(2023, 1, 1),
    catchup=False,
    tags=['gcs', 's3', 'custom'],
) as dag:

    custom_transfer_task = PythonOperator(
        task_id='custom_transfer_task',
        python_callable=transfer_gcs_to_s3_with_logic,
        provide_context=True,
    )

内容的提问来源于stack exchange,提问作者pas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 14:41:04