Airflow 2.6.2中PubSubPullOperator导入失败的Broken DAG修复求助
问题分析与解决
核心问题
你遇到的ImportError是因为**PubSubPullOperator不属于sensors子模块**,它的正确导入路径是airflow.providers.google.cloud.operators.pubsub,而PubSubPullSensor才归属sensors子模块。
修复步骤
修正导入语句
将原有导入代码:from airflow.providers.google.cloud.sensors.pubsub import PubSubPullSensor, PubSubPullOperator修改为:
from airflow.providers.google.cloud.sensors.pubsub import PubSubPullSensor from airflow.providers.google.cloud.operators.pubsub import PubSubPullOperator验证Provider版本匹配
Airflow 2.6.2对应的apache-airflow-providers-google推荐版本为**>=10.0.0**,可通过以下命令检查当前版本:pip show apache-airflow-providers-google若版本过低,执行升级:
pip install apache-airflow-providers-google --upgradeGKE环境配置确认
- 确保GKE集群内所有Airflow组件(调度器、Worker、Webserver)使用的Python环境一致,避免手动安装的依赖仅在单个节点生效。
- 若采用Docker镜像部署Airflow,需在镜像构建阶段预装
apache-airflow-providers-google,而非事后登录节点手动安装。示例Dockerfile配置:RUN pip install apache-airflow-providers-google>=10.0.0 - 确认Airflow的
core.load_examples配置为False(可选,避免示例DAG干扰),并重启所有Airflow组件使配置生效。
使用示例
以下是同时调用PubSubPullSensor和PubSubPullOperator的简单DAG示例:
from airflow import DAG from airflow.providers.google.cloud.sensors.pubsub import PubSubPullSensor from airflow.providers.google.cloud.operators.pubsub import PubSubPullOperator from datetime import datetime default_args = { 'start_date': datetime(2023, 1, 1), } with DAG('pubsub_example_dag', default_args=default_args, schedule_interval='@daily', catchup=False) as dag: # 监听PubSub主题 listen_pubsub = PubSubPullSensor( task_id='listen_pubsub', project_id='your-gcp-project', subscription='your-subscription', max_messages=5, poke_interval=10, ) # 拉取并处理PubSub消息 pull_pubsub = PubSubPullOperator( task_id='pull_pubsub', project_id='your-gcp-project', subscription='your-subscription', max_messages=5, ack_messages=True, ) listen_pubsub >> pull_pubsub
内容的提问来源于stack exchange,提问作者JYang
相关产品推荐
相关产品推荐

