You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 07:56:09