使用Apache Airflow导出BigQuery数据至AWS S3失败的问题
解决BigQuery通过Airflow导出到AWS S3的区域错误问题
问题场景
在GCP BigQuery控制台UI中运行导出到AWS S3的查询可正常执行,但通过Apache Airflow执行相同查询时触发报错:
google.api_core.exceptions.BadRequest: 400 EXPORT to AWS S3 is only supported for tables present in BigQuery Omni AWS regions.; reason: invalidQuery, message: EXPORT to AWS S3 is only supported for tables present in BigQuery Omni AWS regions.
对应的Airflow代码如下:
from airflow import DAG from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator with DAG( dag_id="sql_dag", schedule=None ): query = """ EXPORT DATA WITH CONNECTION `aws-eu-west-1.export_to_s3` OPTIONS(uri=\"s3://XXXXXXX/export_1/*.json\", format=\"json\") AS select 10,10 """ BigQueryInsertJobOperator( task_id=f"bq_extract", location="EU", gcp_conn_id="bq", configuration={ "query": { "query": query, "useLegacySql": False, } } )
问题原因
控制台执行时,BigQuery会自动匹配连接对应的Omni AWS区域(从连接aws-eu-west-1.export_to_s3的前缀可看出对应区域是aws-eu-west-1),但Airflow的BigQueryInsertJobOperator中显式指定了location="EU"——这是GCP的通用区域,并非BigQuery Omni的AWS区域,导致导出操作不被支持。
解决方案
将Airflow任务中的location参数修改为对应的BigQuery Omni AWS区域,即aws-eu-west-1,与你的连接区域保持一致。
修改后的代码如下:
from airflow import DAG from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator with DAG( dag_id="sql_dag", schedule=None ): query = """ EXPORT DATA WITH CONNECTION `aws-eu-west-1.export_to_s3` OPTIONS(uri=\"s3://XXXXXXX/export_1/*.json\", format=\"json\") AS select 10,10 """ BigQueryInsertJobOperator( task_id=f"bq_extract", location="aws-eu-west-1", # 此处修改为Omni AWS区域 gcp_conn_id="bq", configuration={ "query": { "query": query, "useLegacySql": False, } } )
内容的提问来源于stack exchange,提问作者raphaelauv
相关产品推荐
相关产品推荐

