Airflow 1.10.15捕获终止信号同步关闭触发的Azure DevOps流水线
适配Airflow 1.10.15的原生解决方案
核心逻辑:Airflow 1.x中任务被手动终止、超时终止时,任务主进程会收到SIGTERM信号,我们可以在自定义Operator内注册信号回调,捕获终止信号后主动调用Azure DevOps的取消流水线接口,即可避免流水线成为孤儿进程持续运行。
具体实现步骤
- 第一步:调用启动Azure DevOps流水线的POST接口时,将返回的运行ID
run_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
相关产品推荐
相关产品推荐

