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

Airflow 1.10.15捕获终止信号同步关闭触发的Azure DevOps流水线

适配Airflow 1.10.15的原生解决方案

核心逻辑:Airflow 1.x中任务被手动终止、超时终止时,任务主进程会收到SIGTERM信号,我们可以在自定义Operator内注册信号回调,捕获终止信号后主动调用Azure DevOps的取消流水线接口,即可避免流水线成为孤儿进程持续运行。

具体实现步骤

  • 第一步:调用启动Azure DevOps流水线的POST接口时,将返回的运行IDrun_id存储为Operator的实例变量,供后续取消操作使用
  • 第二步:在Operator内注册SIGTERM信号处理函数,捕获到终止信号时触发取消流水线的逻辑
  • 第三步:轮询流水线状态的循环中增加终止信号标记位判断,收到信号后直接退出循环

代码示例

import signal
import time
from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults

class AzureDevopsTriggerOperator(BaseOperator):
    @apply_defaults
    def __init__(self, azure_devops_conn_id, pipeline_id, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self.azure_devops_conn_id = azure_devops_conn_id
        self.pipeline_id = pipeline_id
        self.current_run_id = None # 存储当前启动的流水线运行ID
        self._sigterm_received = False # 终止信号标记位

    def _handle_sigterm(self, signum, frame):
        # SIGTERM信号回调函数
        self.log.info("收到任务终止指令,开始取消关联的Azure DevOps流水线")
        self._sigterm_received = True
        if self.current_run_id:
            self._call_cancel_pipeline_api(self.current_run_id)
            self.log.info(f"已提交流水线取消请求,运行ID:{self.current_run_id}")

    def _call_cancel_pipeline_api(self, run_id):
        """实现Azure DevOps取消运行中流水线的接口调用逻辑"""
        # 此处替换为你自己的Azure API调用代码,接口通常为PATCH格式:
        # https://dev.azure.com/{组织名}/{项目名}/_apis/pipelines/runs/{run_id}?api-version=7.1
        # 请求体设置state为Canceling即可
        pass

    def _start_pipeline(self):
        """实现启动Azure DevOps流水线的接口调用,返回运行ID"""
        # 替换为你自己的启动逻辑
        pass

    def _get_pipeline_status(self, run_id):
        """实现流水线运行状态查询逻辑,返回当前状态"""
        # 替换为你自己的查询逻辑
        pass

    def execute(self, context):
        # 备份原有SIGTERM处理器,后续恢复避免影响其他任务
        original_sigterm = signal.getsignal(signal.SIGTERM)
        signal.signal(signal.SIGTERM, self._handle_sigterm)

        try:
            # 启动流水线
            self.current_run_id = self._start_pipeline()
            self.log.info(f"成功启动Azure DevOps流水线,运行ID:{self.current_run_id}")

            # 轮询等待流水线结束,或收到终止信号
            while not self._sigterm_received:
                current_status = self._get_pipeline_status(self.current_run_id)
                if current_status in ["completed", "failed", "canceled"]:
                    self.log.info(f"流水线运行结束,最终状态:{current_status}")
                    return
                # 轮询间隔可自行调整,单位秒
                time.sleep(30)
        finally:
            # 恢复原有信号处理器
            signal.signal(signal.SIGTERM, original_sigterm)

注意事项

  • 取消流水线的接口调用要做好异常捕获,避免取消请求失败导致任务无法正常退出
  • 如果你的自定义类继承的是BaseSensorOperator,可以将_sigterm_received判断放在poke方法开头,收到信号直接抛出AirflowSkipException或AirflowException终止任务
  • 该方案同时覆盖手动终止、任务执行超时两种终止场景,符合Airflow 1.10.x的原生运行逻辑,不需要额外引入异步检查逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 06:48:03