Airflow可延迟运算符示例BaseTrigger/TriggerEvent导入及类归属疑问
Airflow 可延迟运算符开发相关类说明
缺失类的正确导入路径
官方示例省略了头部导入语句,两个核心类的导入方式如下:
TriggerEvent:触发器完成事件类,归属airflow.triggers.base模块,对应导入语句:from airflow.triggers.base import TriggerEventBaseTrigger:所有自定义触发器的继承基类,和TriggerEvent同属一个模块,对应导入语句:from airflow.triggers.base import BaseTrigger
BaseTrigger与airflow.models.trigger.Trigger的关系
二者是完全等价的同一个类。airflow.models.trigger.Trigger是Airflow早期版本遗留的兼容别名,当前官方推荐统一使用airflow.triggers.base.BaseTrigger作为导入路径,二者功能、接口没有任何差异。
补全导入后的可运行示例代码
import asyncio from airflow.triggers.base import BaseTrigger, TriggerEvent from airflow.utils import timezone class DateTimeTrigger(BaseTrigger): def __init__(self, moment): super().__init__() self.moment = moment def serialize(self): return ("airflow.triggers.temporal.DateTimeTrigger", {"moment": self.moment}) async def run(self): while self.moment > timezone.utcnow(): await asyncio.sleep(1) yield TriggerEvent(self.moment)
开发自定义触发器时,必须实现基类要求的
serialize()和异步run()方法:serialize()负责返回触发器的重建路径和参数,供触发器独立进程序列化加载使用;run()为异步执行逻辑,当满足触发条件时通过yield返回TriggerEvent实例,即可结束触发状态,调度器会恢复对应可延迟运算符的后续执行。
内容的提问来源于stack exchange,提问作者Jiew Meng
相关产品推荐
相关产品推荐

