Airflow导入报错:无法引入BigQueryTableExistenceAsyncSensor
问题背景
尝试从airflow.providers.google.cloud.sensors.bigquery导入BigQueryTableExistenceAsyncSensor,用于检查BigQuery数据集内的表是否存在,但Airflow返回DAG导入错误。
错误信息
Broken DAG: [/home/airflow/gcs/dags/manatal/test_dag.py] Traceback (most recent call last):
File "<frozen importlib._bootstrap>", line 219, in _call_with_frames_removed
File "/home/airflow/gcs/dags/test/test_dag.py", line 4, in <module>
from airflow.providers.google.cloud.sensors.bigquery import BigQueryTableExistenceAsyncSensor
ImportError: cannot import name 'BigQueryTableExistenceAsyncSensor' from 'airflow.providers.google.cloud.sensors.bigquery' (/opt/python3.8/lib/python3.8/site-packages/airflow/providers/google/cloud/sensors/bigquery.py)
原代码示例
from airflow import DAG from util.dags_hourly import create_dag_write_append # 自定义类,其他DAG使用无问题 from airflow.providers.google.cloud.sensors.bigquery import BigQueryTableExistenceAsyncSensor def __init__(self, dataset=None, table_name=None): self.dataset = dataset self.table_name = table_name def check_table_exists(self): return BigQueryTableExistenceAsyncSensor( task_id="check_table_exists_async", project_id='x-staging', dataset_id=self.dataset, table_id=self.table ) with create_dag_write_append('test') as dag: a = BigQueryTableExistenceAsyncSensor( dataset_id='data_lake_staging', table_id='test_table' ) task1 = a.check_table_exists() task1
解决方案
1. 检查并升级Google Provider版本
BigQueryTableExistenceAsyncSensor是在apache-airflow-providers-google 8.6.0版本新增的传感器类。如果环境中该包版本低于此,就会出现导入错误。执行以下命令升级:
pip install --upgrade apache-airflow-providers-google>=8.6.0
2. 修正代码逻辑
原代码存在逻辑错误:错误地将BigQueryTableExistenceAsyncSensor当作自定义类来实例化并调用check_table_exists方法,而实际上它本身就是Airflow的Sensor任务类,可直接在DAG中定义使用。修正后的代码如下:
from airflow import DAG from util.dags_hourly import create_dag_write_append from airflow.providers.google.cloud.sensors.bigquery import BigQueryTableExistenceAsyncSensor with create_dag_write_append('test') as dag: check_table_task = BigQueryTableExistenceAsyncSensor( task_id="check_table_exists_async", project_id='x-staging', dataset_id='data_lake_staging', table_id='test_table' ) check_table_task
内容的提问来源于stack exchange,提问作者Mohammad Iqbal

