如何使用Airflow运算符完成BigQuery到Amazon S3的数据传输
实现方案
注意:你当前使用的BigQueryGetDataOperator仅适合小数据量场景,该运算符会将全量数据读取到Airflow工作节点内存中,数据量超过1G就容易出现OOM故障,生产环境优先使用导出中转方案
推荐生产级流程(无内存瓶颈)
- 第一步:用
BigQueryToGCSOperator将BigQuery表批量导出到GCS的临时存储桶,支持按大小分片、指定导出格式(CSV/Parquet/JSON都可) - 第二步:用
GCSToS3Operator将GCS临时路径下的导出文件同步到S3目标路径 - 第三步:可选追加清理任务,删除GCS路径下的临时导出文件节省存储成本
可直接复用的代码示例
提前在Airflow连接配置中准备好已有的GCP_CONN_ID,以及AWS连接(命名为AWS_CONN_ID即可),替换代码中自定义变量后即可使用:
# 导入需要的运算符 from airflow.providers.google.cloud.transfers.bigquery_to_gcs import BigQueryToGCSOperator from airflow.providers.amazon.aws.transfers.gcs_to_s3 import GCSToS3Operator from airflow.providers.google.cloud.operators.gcs import GCSDeleteObjectsOperator # 你之前的检查任务保留,调整SQL为计数判断减少资源消耗 bq_check_date = BigQueryCheckOperator( task_id='bq_check_date', sql=''' SELECT COUNT(*) > 0 FROM `myproject.test.test_table` ''', use_legacy_sql=False, bigquery_conn_id=GCP_CONN_ID, dag=dag ) # 1. 导出BQ到GCS临时路径 export_bq_to_gcs = BigQueryToGCSOperator( task_id="export_bq_to_gcs", source_project_dataset_table=f"{GCP_PROJECT}.{BQ_DATASET}.{BQ_TABLE}", destination_cloud_storage_uris=[f"gs://{YOUR_GCS_TEMP_BUCKET}/bq_export/{BQ_TABLE}/*.parquet"], export_format="PARQUET", # 可以换成CSV、JSON等格式 compression="SNAPPY", # 开启压缩减少传输时间 location=LOCATION, bigquery_conn_id=GCP_CONN_ID, dag=dag ) # 2. 同步GCS文件到S3 transfer_gcs_to_s3 = GCSToS3Operator( task_id="transfer_gcs_to_s3", bucket=YOUR_GCS_TEMP_BUCKET, prefix=f"bq_export/{BQ_TABLE}/", dest_s3_key=f"s3://{YOUR_S3_BUCKET}/bq_data/{BQ_TABLE}/", gcp_conn_id=GCP_CONN_ID, aws_conn_id=AWS_CONN_ID, replace=True, # 覆盖S3已有同名文件 dag=dag ) # 3. 可选:清理GCS临时文件 clear_gcs_temp = GCSDeleteObjectsOperator( task_id="clear_gcs_temp", bucket_name=YOUR_GCS_TEMP_BUCKET, prefix=f"bq_export/{BQ_TABLE}/", gcp_conn_id=GCP_CONN_ID, dag=dag ) # 设置任务依赖 bq_check_date >> export_bq_to_gcs >> transfer_gcs_to_s3 >> clear_gcs_temp
小数据量场景简化方案
如果你的表数据量低于100M,不想用GCS中转,可以保留你写的BigQueryGetDataOperator,后续追加S3PutObjectOperator把读取到的内存数据写入S3:
from airflow.providers.amazon.aws.operators.s3 import S3PutObjectOperator import json # 你之前的get_data任务保留 get_data = BigQueryGetDataOperator( task_id="get_data", dataset_id=BQ_DATASET, table_id=BQ_TABLE, location=LOCATION, bigquery_conn_id=GCP_CONN_ID, dag = dag, ) # 将读取到的列表数据转成JSON写入S3 write_to_s3 = S3PutObjectOperator( task_id="write_to_s3", bucket_name=YOUR_S3_BUCKET, key=f"bq_data/{BQ_TABLE}.json", data=json.dumps({{ ti.xcom_pull(task_ids='get_data') }}), # 从XCom拉取上游取到的数据 aws_conn_id=AWS_CONN_ID, dag=dag ) # 依赖设置 bq_check_date >> get_data >> write_to_s3
内容的提问来源于stack exchange,提问作者Abdulsalam Fadhel
相关产品推荐
相关产品推荐

