Apache Airflow 2:如何不重写execute实现带单次初始操作的自定义Sensor
Airflow 2 自定义Sensor单次前置操作实现方案(不重写execute方法)
核心实现思路
利用Sensor实例属性在同任务的多次poke调用间会保留状态的特性,新增标记位控制仅一次的前置操作执行逻辑,无需改动原有execute方法。
具体实现步骤
- 自定义Sensor类新增实例属性(不要用类属性,避免多实例共用标记出现逻辑混乱),比如命名为
_has_run_init_op,默认值设为False,用于标记前置操作的执行状态 - 重写
poke方法时先判断标记位:- 标记位为
False时,先执行需要仅运行一次的特定操作,执行完成后将标记位改为True - 标记位为
True时,直接执行常规的目标条件校验逻辑即可
- 标记位为
完整代码示例
from airflow.sensors.base import BaseSensorOperator from airflow.utils.decorators import apply_defaults class CustomPreOpSensor(BaseSensorOperator): @apply_defaults def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) # 初始化前置操作执行标记,仅当前实例有效 self._has_run_init_op = False def poke(self, context): # 首次调用poke时执行前置操作 if not self._has_run_init_op: # 替换为你需要仅执行一次的业务逻辑 print("执行仅运行一次的前置操作") # 标记已执行,后续poke不会再触发该部分逻辑 self._has_run_init_op = True # 替换为你原本的目标条件校验逻辑,返回True则任务终止,返回False则进入下一轮探测 condition_res = False return condition_res
注意事项
该方案默认适配Sensor的poke运行模式,如果你使用的是reschedule模式,每次重调度会重新实例化Sensor对象,实例属性标记会失效,这种场景可以将执行标记写入XCom或者Airflow Variable,每次poke先查询标记存储值判断是否需要执行前置操作即可。
内容的提问来源于stack exchange,提问作者sanchit08
相关产品推荐
相关产品推荐

