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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 11:51:02