使用Airflow BigQueryInsertJobOperator导出GA数据遇权限及参数错误
问题:Airflow中使用BigQueryInsertJobOperator提取GA数据遇参数及权限异常
背景
尝试用Airflow的BigQueryInsertJobOperator替代BigQueryExecuteQueryOperator从BigQuery获取GA数据,但官方文档对前者的用法说明不够清晰。
DAG执行流程
- 列出目标数据集中的所有表名;
- 使用BigQueryInsertJobOperator执行GA数据查询,查询语法如下:
`{my-project}.{my-dataset}.events_*` WHERE _TABLE_SUFFIX BETWEEN '{start}' AND '{end}'
对应的Operator代码:
select_query_job = BigQueryInsertJobOperator( task_id="select_query_job", gcp_conn_id='big_query', configuration={ "query": { "query": build_query.output, "useLegacySql": False, "allowLargeResults": True, "useQueryCache": True, } } )
- 从Xcom获取上述查询任务的job id,再通过配置了
extract的BigQueryInsertJobOperator提取查询结果,但执行失败。
尝试的两种Operator写法
# 写法一 retrieve_job_data = BigQueryInsertJobOperator( task_id="get_job_data", gcp_conn_id='big_query', job_id=select_query_job.output, project_id=project_name, configuration={ "extract": { } } ) # 写法二 retrieve_job_data = BigQueryInsertJobOperator( task_id="get_job_data", gcp_conn_id='big_query', configuration={ "extract": { "jobId": select_query_job.output, "projectId": project_name } } )
报错信息
任务执行错误日志
google.api_core.exceptions.BadRequest: 400 POST https://bigquery.googleapis.com/bigquery/v2/projects/{my-project}/jobs?prettyPrint=false: Required parameter is missing [2022-08-16, 09:44:01 UTC] {taskinstance.py:1415} INFO - Marking task as FAILED. dag_id=BIG_QUERY, task_id=get_job_data, execution_date=20220816T054346, start_date=20220816T054358, end_date=20220816T054401 [2022-08-16, 09:44:01 UTC] {standard_task_runner.py:92} ERROR - Failed to execute job 628 for task get_job_data (400 POST https://bigquery.googleapis.com/bigquery/v2/projects/{my-project}/jobs?prettyPrint=false: Required parameter is missing; 100144)
直接访问API链接的错误返回
{ "error": { "code": 401, "message": "Request is missing required authentication credential. Expected OAuth 2 access token, login cookie or other valid authentication credential.", "errors": [ { "message": "Login Required.", "domain": "global", "reason": "required", "location": "Authorization", "locationType": "header" } ], "status": "UNAUTHENTICATED", "details": [ { "@type": "type.googleapis.com/google.rpc.ErrorInfo", "reason": "CREDENTIALS_MISSING", "domain": "googleapis.com", "metadata": { "service": "bigquery.googleapis.com", "method": "google.cloud.bigquery.v2.JobService.ListJobs" } } ] } }
疑问
gcp_conn_id在其他Operator中可正常工作,且已指定project_id,却出现401权限认证提示,不符合预期。
解决方案
1. 核心问题:误解了extract配置的作用
BigQueryInsertJobOperator的extract配置不是用来获取已有查询任务的结果,而是需要明确指定提取的数据源(查询结果存储的表)和输出目标位置。你之前的两种写法都缺少关键参数,导致BigQuery返回400错误;而401错误是直接访问API链接时的未授权状态,和任务执行的权限无关。
2. 修正后的完整流程代码
步骤1:修改查询任务,指定结果存储表
首先在查询任务中配置destinationTable,将查询结果持久化到临时表(避免临时表自动过期):
select_query_job = BigQueryInsertJobOperator( task_id="select_query_job", gcp_conn_id='big_query', configuration={ "query": { "query": build_query.output, "useLegacySql": False, "allowLargeResults": True, "useQueryCache": True, # 指定查询结果保存的目标表 "destinationTable": { "projectId": project_name, "datasetId": "your_temp_dataset", # 提前创建好的临时数据集 "tableId": f"temp_ga_results_{{{{ ds_nodash }}}}" # 用日期后缀避免冲突 }, "writeDisposition": "WRITE_TRUNCATE" # 每次执行覆盖已有表 } } )
步骤2:配置extract任务提取数据
提取任务的extract配置需要指向上述临时表,同时指定输出到GCS的路径:
retrieve_job_data = BigQueryInsertJobOperator( task_id="get_job_data", gcp_conn_id='big_query', configuration={ "extract": { "sourceTable": { "projectId": project_name, "datasetId": "your_temp_dataset", "tableId": f"temp_ga_results_{{{{ ds_nodash }}}}" }, "destinationUris": ["gs://your-gcs-bucket/ga-results/ga_data_{{{{ ds_nodash }}}}.csv"], "destinationFormat": "CSV", # 支持JSON、AVRO等格式 "printHeader": True # 是否输出表头 } } )
3. 权限校验
确保gcp_conn_id对应的服务账号拥有以下权限:
- BigQuery数据读取权限(
bigquery.tables.getData) - BigQuery作业执行权限(
bigquery.jobs.create) - GCS写入权限(
storage.objects.create)
4. 可选:合并查询与提取为单个任务
如果不需要拆分任务,可以在同一个BigQueryInsertJobOperator中同时配置查询和提取逻辑:
query_and_extract_job = BigQueryInsertJobOperator( task_id="query_and_extract_job", gcp_conn_id='big_query', configuration={ "query": { "query": build_query.output, "useLegacySql": False, "allowLargeResults": True, "useQueryCache": True, "destinationTable": { "projectId": project_name, "datasetId": "your_temp_dataset", "tableId": f"temp_ga_results_{{{{ ds_nodash }}}}" }, "writeDisposition": "WRITE_TRUNCATE" }, "extract": { "sourceTable": { "projectId": project_name, "datasetId": "your_temp_dataset", "tableId": f"temp_ga_results_{{{{ ds_nodash }}}}" }, "destinationUris": ["gs://your-gcs-bucket/ga-results/ga_data_{{{{ ds_nodash }}}}.csv"], "destinationFormat": "CSV" } } )
内容的提问来源于stack exchange,提问作者Justin
相关产品推荐
相关产品推荐

