如何使用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
相关产品推荐
相关产品推荐

