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

如何在Airflow DAG运行终止前触发自定义逻辑关闭远程任务?

解决方案

Airflow 支持在任务被取消/终止时触发自定义清理逻辑,针对你的场景可以通过以下核心方式实现:

1. 重写Operator的on_kill方法

所有继承自BaseOperator的Operator(包括你的JobOperator)都可以重写on_kill方法,当任务被手动取消、超时终止或DAG运行被取消时,Airflow会自动调用这个方法。你可以在这里加入关闭远程任务的逻辑。

修改你的Operator代码如下:

class JobOperator(base_sensor_operator.BaseSensorOperator):
    def __init__(self, job_config, **kwargs):
        super().__init__(**kwargs)
        self._job_config = job_config
        self._job_started = False
        self._job_info = None  # 确保job_info被类实例持有

    def poke(self, context: Any) -> bool:
        if not self._job_started:
            self._job_info = LaunchJob(self._job_config)
            self._job_started = True
        else:
            status = PollJob(self._job_info)
            if status == "SUCCESS":
                # 成功后的处理逻辑
                return True
            elif status == "FAILED":
                # 失败后的处理逻辑
                return True
            else:
                # 继续轮询
                return False

    def on_kill(self):
        # 当任务被取消时,关闭远程任务
        if self._job_info is not None:
            CloseRemoteJob(self._job_info)  # 替换为你的关闭远程任务的逻辑

2. 处理DAG运行被删除的场景

当你程序化删除DAG运行时:

  • 如果任务正处于运行状态(比如你的Sensor正在轮询),Airflow会先终止该任务实例,触发on_kill方法,此时远程任务会被正常关闭。
  • 如果任务已经完成,删除DAG运行不会触发额外逻辑,此时远程任务可能已经结束,无需处理。

如果需要针对DAG运行删除事件做更全局的处理,可以考虑使用Airflow的事件触发器或编写自定义的API钩子,但对于你的场景,重写on_kill已经足够覆盖绝大多数需求。


内容的提问来源于stack exchange,提问作者Tony Liao

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 21:15:17