Airflow BigQueryInsertJobOperator调用SP报destinationTable KeyError
问题描述
使用Airflow的BigQueryInsertJobOperator调用BigQuery存储过程时执行报错,多次调整SQL语法均无法解决;但将存储过程内的SQL逻辑提取为独立脚本文件执行时,运行正常无报错。
测试用DAG代码
PROJECT_ID = os.environ.get("GCP_PROJECT_ID", "project-id") Dataset = os.environ.get("GCP_Dataset", "dataset") with DAG(dag_id='dag_id',default_args=default_args,schedule_interval="@daily", start_date=days_ago(1), catchup=False ) as dag: Call_SP = BigQueryInsertJobOperator( task_id='Call_SP', configuration={ "query": { "query": "CALL `" + PROJECT_ID + "." + Dataset + "." + "SP`();", #"query": "{% include 'Scripts/Script.sql' %}", "useLegacySql": False, } } ) Call_SP
报错日志
任务运行时生成的CALL语句符合预期,但抛出键不存在的错误,日志如下:
[2022-06-30, 17:32:07 UTC] {bigquery.py:2247} INFO - Executing: {'query': {'query': 'CALL `project-id.dataset.SP`();', 'useLegacySql': False}} [2022-06-30, 17:32:07 UTC] {bigquery.py:1560} INFO - Inserting job airflow_Call_SP_2022_06_30T16_47_13_748417_00_00_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx [2022-06-30, 17:32:13 UTC] {taskinstance.py:1776} ERROR - Task failed with exception Traceback (most recent call last): File "/opt/python3.8/lib/python3.8/site-packages/airflow/providers/google/cloud/operators/bigquery.py", line 2269, in execute table = job.to_api_repr()["configuration"]["query"]["destinationTable"] KeyError: 'destinationTable'
存储过程内部仅包含单条MERGE语句,无其他复杂逻辑。
报错根因
该报错是旧版Airflow Google Provider的逻辑缺陷:BigQueryInsertJobOperator在作业执行完成后,默认会读取返回结果中configuration.query.destinationTable字段,用于血缘采集、结果表关联等后续逻辑。
使用CALL语句调用存储过程属于BigQuery脚本类作业,即使存储过程内部执行MERGE写操作,外层CALL作业的返回配置中也不会携带destinationTable字段,因此触发KeyError。
直接执行MERGE语句的普通查询作业会正常返回目标表字段,因此独立脚本运行无报错。
修复方案
根据当前安装的Airflow Google Provider版本选择对应处理方式:
- 对于 apache-airflow-providers-google < 8.2.0 版本:在算子参数中显式添加
create_insert_job=False,关闭自动结果表解析逻辑,修改后代码示例:
Call_SP = BigQueryInsertJobOperator( task_id='Call_SP', create_insert_job=False, configuration={ "query": { "query": "CALL `" + PROJECT_ID + "." + Dataset + "." + "SP`();", "useLegacySql": False, } } )
- 对于 apache-airflow-providers-google >= 8.2.0 版本:官方已修复该逻辑缺陷,会自动识别脚本类作业跳过
destinationTable字段读取。如果仍遇到同类报错,可显式指定destinationTable=None,或升级provider包到最新稳定版本。 - 临时兼容方案:如果不方便修改算子参数或升级版本,可将存储过程调用逻辑放入独立SQL文件,通过注释中Jinja模板
{% include 'Scripts/Script.sql' %}的方式加载,同时在query配置中显式指定MERGE操作的实际目标表为destinationTable,即可绕过字段校验。
注意:不要为了绕过报错给CALL语句配置无关的目标表,会导致BigQuery将存储过程的空返回结果写入目标表,造成原有数据覆盖。
内容的提问来源于stack exchange,提问作者pyFlummox
相关产品推荐
相关产品推荐

