寻求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权限
- 连接ID建议设为
- 配置AWS连接:
同样在Connections页面,创建类型为「Amazon Web Services」的连接:- 连接ID建议设为
aws_default(自定义ID需与代码中对应) - 填入AWS的Access Key ID和Secret Access Key,该账号需拥有目标S3存储桶的
s3:PutObject权限
- 连接ID建议设为
- 网络访问:确保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
相关产品推荐
相关产品推荐

