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

Airflow UI手动标记任务为成功时如何触发API调用?

解决Airflow手动标记任务成功时触发自定义API的问题

针对你在Airflow 2.10.0中遇到的自定义Operator无法在手动标记成功时触发API的问题,核心原因是:手动标记任务成功(Mark Success)是直接更新元数据库的任务状态,不会触发Operator的execute方法、on_success_callback或常规回调——这些回调仅在任务实际执行完成并返回成功状态时才会触发。

以下是两种可行的解决方案:

方案一:使用任务实例状态变更钩子(推荐)

Airflow 2.x提供了任务实例状态变更监听机制,可捕获所有状态变更事件(包括手动修改),进而执行自定义逻辑。

实现步骤:

  1. 创建自定义钩子类,监听状态变更:
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)
  1. 在Airflow配置文件airflow.cfg中注册这个钩子:
    找到task_instance_state_change_hooks配置项,添加自定义钩子的完整类路径:
task_instance_state_change_hooks = your.package.path.ManualSuccessTriggerHook
  1. 重启Airflow Webserver和Scheduler,使配置生效。

方案二:自定义状态校验任务(备选)

如果不想使用全局钩子,可添加一个周期性校验任务,专门检测目标任务的手动标记成功事件:

  1. 创建状态校验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
  1. 将这个Sensor添加到DAG中,设置合适的检查间隔(比如每分钟一次)。

关键说明

  • on_failure_callback仅在任务执行失败时触发,on_kill仅在任务被主动取消时触发,两者都不适用于手动标记成功的场景。
  • 全局钩子会监听所有任务的状态变更,若需针对特定任务触发API,可在钩子逻辑中添加task_instance.task_id或task_instance.dag_id的判断。

内容的提问来源于stack exchange,提问作者Alex Man

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 17:42:37