Airflow UI手动标记任务为成功时如何触发API调用?
解决Airflow手动标记任务成功时触发自定义API的问题
针对你在Airflow 2.10.0中遇到的自定义Operator无法在手动标记成功时触发API的问题,核心原因是:手动标记任务成功(Mark Success)是直接更新元数据库的任务状态,不会触发Operator的execute方法、on_success_callback或常规回调——这些回调仅在任务实际执行完成并返回成功状态时才会触发。
以下是两种可行的解决方案:
方案一:使用任务实例状态变更钩子(推荐)
Airflow 2.x提供了任务实例状态变更监听机制,可捕获所有状态变更事件(包括手动修改),进而执行自定义逻辑。
实现步骤:
- 创建自定义钩子类,监听状态变更:
from airflow.hooks.base import BaseHook from airflow.models import TaskInstance from airflow.utils.state import State class ManualSuccessTriggerHook(BaseHook): def on_task_instance_state_change( self, previous_state: str, new_state: str, task_instance: TaskInstance, session, **kwargs ): # 仅处理"手动标记成功"的场景:新状态为SUCCESS,且之前状态不是SUCCESS if new_state == State.SUCCESS and previous_state != State.SUCCESS: # 替换为你的API调用逻辑 self.log.info(f"任务 {task_instance.task_id} 被手动标记为成功,触发API调用") # call_your_api(task_instance.dag_id, task_instance.task_id)
- 在Airflow配置文件
airflow.cfg中注册这个钩子:
找到task_instance_state_change_hooks配置项,添加自定义钩子的完整类路径:
task_instance_state_change_hooks = your.package.path.ManualSuccessTriggerHook
- 重启Airflow Webserver和Scheduler,使配置生效。
方案二:自定义状态校验任务(备选)
如果不想使用全局钩子,可添加一个周期性校验任务,专门检测目标任务的手动标记成功事件:
- 创建状态校验Sensor:
from airflow.sensors.base import BaseSensorOperator from airflow.models import TaskInstance from airflow.utils.state import State class ManualSuccessSensor(BaseSensorOperator): def poke(self, context): ti = TaskInstance( task_id="applitus-update-job", dag_id=context["dag"].dag_id, execution_date=context["execution_date"] ) ti.refresh_from_db() # 任务被手动标记成功时,start_date会为空(未实际执行) if ti.state == State.SUCCESS and ti.start_date is None: # 触发API调用 self.log.info("检测到手动标记成功,触发API") # call_your_api() return True return False
- 将这个Sensor添加到DAG中,设置合适的检查间隔(比如每分钟一次)。
关键说明
on_failure_callback仅在任务执行失败时触发,on_kill仅在任务被主动取消时触发,两者都不适用于手动标记成功的场景。- 全局钩子会监听所有任务的状态变更,若需针对特定任务触发API,可在钩子逻辑中添加
task_instance.task_id或task_instance.dag_id的判断。
内容的提问来源于stack exchange,提问作者Alex Man
相关产品推荐
相关产品推荐

