Airflow能否在DAG文件中嵌入含自定义Trigger的自定义Operator?报错求助
问题分析与解决方案
问题原因
Airflow的Triggerer进程独立运行,不会将Dags目录下的文件添加到Python模块搜索路径中。当你在DAG文件里直接定义MyTrigger时,Scheduler和Webserver能执行DAG代码,但Triggerer进程尝试通过my_dag.MyTrigger导入类时,找不到对应模块,因此抛出ModuleNotFoundError。
解决方案
把自定义Trigger和Operator迁移到Composer的Plugins目录,该目录会被自动添加到所有Airflow组件(包括Triggerer)的Python路径中,确保Triggerer能正确导入类。
步骤1:创建插件模块文件
在Cloud Composer的GCS存储桶中找到plugins目录(路径为gs://<你的Composer存储桶名>/plugins/),创建my_custom_operators.py文件,内容如下:
from airflow.triggers.base import BaseTrigger, TriggerEvent from airflow.utils.context import Context from airflow.models.baseoperator import BaseOperator import os class MyTrigger(BaseTrigger): def __init__(self, path): super().__init__() self.path = path def serialize(self): # 修改类路径为插件模块下的完整路径 return ("my_custom_operators.MyTrigger", {"path": self.path}) async def run(self): print("MyTrigger.run") exists = os.path.exists(self.path) yield TriggerEvent(exists) class MyOperator(BaseOperator): def __init__(self, path, deferrable, **kwargs) -> None: super().__init__(**kwargs) self.path = path self.deferrable = deferrable def execute(self, context: Context) -> None: print("MyOperator.execute: before deferring") if self.deferrable: self.defer( trigger=MyTrigger(path=self.path), method_name="execute_complete", ) else: print("deferrable=False") print("MyOperator.execute: after deferring") def execute_complete(self, context: Context, event: dict) -> None: # 修正原代码中{dict}的错误,改为{event} print(f"MyOperator.execute_compute: {event}") return
步骤2:修改DAG文件
更新my_dag.py,从插件模块导入自定义Operator:
from airflow import DAG from my_custom_operators import MyOperator from datetime import datetime dag = DAG('my_dag', description='Example', schedule_interval='@once', start_date=datetime(2023, 8, 17)) my_operator = MyOperator(task_id='my_task', path="/home/airflow/composer_kube_config", deferrable=True, dag=dag) my_operator
步骤3:验证部署
将修改后的文件上传到对应GCS目录,Composer会自动同步插件和DAG。触发DAG运行后,Triggerer就能正常加载MyTrigger类,执行触发逻辑。
额外说明
- 原代码中
execute_complete方法里的{dict}是语法错误,会打印dict类的字符串表示而非事件内容,已在示例中修正为{event}。 - Plugins目录是Airflow官方推荐的自定义组件存放位置,所有Airflow进程都能访问这里的模块,避免了跨进程的导入问题。
内容的提问来源于stack exchange,提问作者Dagang Wei
相关产品推荐
相关产品推荐

