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

如何使用Airflow将BigQuery查询数据导出到GCS?

适用Operator及实现方案

针对你的需求,有两种常用的Operator组合可选,都可以实现先校验数据、再将指定查询结果导出到GCS的流程:

方案1:BigQueryInsertJobOperator + BigQueryToGCSOperator(Airflow 2.x推荐)

Airflow 1.x版本中导出类对应名称为你提到的BigQueryToCloudStorageOperator,该方案逻辑清晰,可复用中间结果,是最常用的实现方式:

  • 用BigQueryCheckOperator完成你已有的数据存在性校验
  • 用BigQueryInsertJobOperator执行过滤查询,将结果写入BigQuery临时表/正式表
  • 用BigQueryToGCSOperator将生成的表数据导出到GCS指定路径

如果不需要额外过滤、直接导出全表数据,可省去查询生成表的步骤,直接调用导出Operator指定源表即可。

代码示例

# Airflow 2.x 依赖导入
from airflow.providers.google.cloud.operators.bigquery import (
    BigQueryCheckOperator,
    BigQueryInsertJobOperator,
    BigQueryToGCSOperator
)

# Airflow 1.x 对应依赖导入
# from airflow.contrib.operators.bigquery_operator import BigQueryOperator
# from airflow.contrib.operators.bigquery_to_gcs import BigQueryToCloudStorageOperator

# 你已有的校验任务
t1 = BigQueryCheckOperator(
    task_id="bq_check_covid_data_exists",
    sql="""
        SELECT COUNT(*) > 0
        FROM bigquery-public-data.covid19_italy.data_by_region
        WHERE DATE(date) = DATE_ADD(DATE "{{ ds }}", INTERVAL -2 DAY)
    """,
    use_legacy_sql=False,
    dag=dag
)

# 任务2:执行过滤查询,结果写入临时表
t2 = BigQueryInsertJobOperator(
    task_id="bq_query_filtered_data",
    configuration={
        "query": {
            "query": """
                SELECT *
                FROM bigquery-public-data.covid19_italy.data_by_region
                WHERE DATE(date) = DATE_ADD(DATE "{{ ds }}", INTERVAL -2 DAY)
            """,
            "useLegacySql": False,
            "destinationTable": {
                "projectId": "替换为你的GCP项目ID",
                "datasetId": "替换为你的BigQuery数据集ID",
                "tableId": "temp_covid_data_{{ ds_nodash }}"
            },
            "writeDisposition": "WRITE_TRUNCATE"
        }
    },
    dag=dag
)

# 任务3:导出临时表到GCS
t3 = BigQueryToGCSOperator(
    task_id="bq_export_to_gcs",
    source_project_dataset_table="替换为你的GCP项目ID.替换为你的BigQuery数据集ID.temp_covid_data_{{ ds_nodash }}",
    destination_cloud_storage_uris=["gs://替换为你的GCS桶名/导出路径/covid_data_{{ ds }}.csv"],
    export_format="CSV",
    field_delimiter=",",
    print_header=True,
    dag=dag
)

# 定义任务依赖
t1 >> t2 >> t3

方案2:单独使用BigQueryInsertJobOperator

你也可以直接调用BigQuery原生的导出能力,一步完成查询+导出,不需要生成中间表,适合轻量导出场景:

代码示例

t2 = BigQueryInsertJobOperator(
    task_id="bq_export_query_result_to_gcs",
    configuration={
        "extract": {
            "sourceQuery": {
                "query": """
                    SELECT *
                    FROM bigquery-public-data.covid19_italy.data_by_region
                    WHERE DATE(date) = DATE_ADD(DATE "{{ ds }}", INTERVAL -2 DAY)
                """,
                "useLegacySql": False
            },
            "destinationUris": ["gs://替换为你的GCS桶名/导出路径/covid_data_{{ ds }}.csv"],
            "destinationFormat": "CSV",
            "fieldDelimiter": ",",
            "printHeader": True
        }
    },
    dag=dag
)

# 任务依赖直接配置为 t1 >> t2 即可

注意事项

  • 需要提前在Airflow中配置好有权限访问BigQuery和GCS的GCP连接,默认连接ID为google_cloud_default
  • 导出大文件时可在GCS路径中添加通配符*,例如gs://桶名/路径/covid_data_{{ ds }}-*.csv,BigQuery会自动拆分大文件分片导出
  • Airflow 1.x环境中将BigQueryToGCSOperator替换为BigQueryToCloudStorageOperator即可,参数配置逻辑基本一致

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 04:54:01