在PythonOperator中调用GCSToBigQueryOperator报logical_date键错误
单独使用GCSToBigQueryOperator时功能正常,但将其放入函数并通过PythonOperator调用时,出现如下错误:
File "/home/airflow/.local/lib/python3.7/site-packages/airflow/providers/google/cloud/transfers/gcs_to_bigquery.py", line 336, in execute
logical_date=context["logical_date"],
KeyError: 'logical_date'
相关代码如下:
def fn_gcs_to_bqstg(): task1 = GCSToBigQueryOperator( task_id='gcstobqstg', bucket=GCS_BUCKET, source_objects=['data/*.json'], destination_project_dataset_table='local_development.table', schema_object='schemas/data_schema.json', source_format='NEWLINE_DELIMITED_JSON', create_disposition='CREATE_IF_NEEDED', write_disposition='WRITE_TRUNCATE', gcp_conn_id='bq-conn', dag=dag ) task1.execute(dict()) gcs_to_bqstg = PythonOperator( task_id='gcs_to_bqstg', python_callable=fn_gcs_to_bqstg)
问题出在手动调用task1.execute(dict())时,传入了空字典作为context参数。GCSToBigQueryOperator的execute方法依赖Airflow运行时提供的完整上下文(包含logical_date等关键字段),空字典没有这些必要的键,因此触发了KeyError。
同时这种用法本身不符合Airflow的设计逻辑:PythonOperator用于执行自定义Python函数,而GCSToBigQueryOperator本身就是独立的Airflow任务算子,不需要嵌套在Python函数里通过PythonOperator调用。
有两种可行的解决方式:
方式一:直接使用GCSToBigQueryOperator作为独立任务
这是最符合Airflow设计的用法,无需额外的PythonOperator,直接将GCSToBigQueryOperator实例作为DAG的任务节点:
gcs_to_bqstg = GCSToBigQueryOperator( task_id='gcstobqstg', bucket=GCS_BUCKET, source_objects=['data/*.json'], destination_project_dataset_table='local_development.table', schema_object='schemas/data_schema.json', source_format='NEWLINE_DELIMITED_JSON', create_disposition='CREATE_IF_NEEDED', write_disposition='WRITE_TRUNCATE', gcp_conn_id='bq-conn', dag=dag )
方式二:若必须通过PythonOperator调用(不推荐)
如果因特殊需求必须在Python函数内调用算子的execute方法,需要让Airflow传入完整上下文,并传递给execute方法:
- 修改
PythonOperator,添加provide_context=True参数,触发Airflow自动传入上下文 - 函数接收
context参数,并将其传给execute方法
修改后的代码:
def fn_gcs_to_bqstg(**context): task1 = GCSToBigQueryOperator( task_id='gcstobqstg', bucket=GCS_BUCKET, source_objects=['data/*.json'], destination_project_dataset_table='local_development.table', schema_object='schemas/data_schema.json', source_format='NEWLINE_DELIMITED_JSON', create_disposition='CREATE_IF_NEEDED', write_disposition='WRITE_TRUNCATE', gcp_conn_id='bq-conn', dag=dag ) task1.execute(context) gcs_to_bqstg = PythonOperator( task_id='gcs_to_bqstg', python_callable=fn_gcs_to_bqstg, provide_context=True # 关键:让Airflow传入上下文 )
内容的提问来源于Stack Exchange,提问作者Deniz Nur

