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

Airflow进程终止却被标记为成功,如何改为失败状态?

解决Airflow任务收到SIGTERM后仍标记为SUCCESS的问题

从日志来看,问题出在任务进程收到SIGTERM后仍以exitcode=0正常退出,Airflow默认以进程退出码判断状态,因此误标记为SUCCESS。以下是几种可行的解决方式:

方法1:自定义SIGTERM信号处理器,主动抛出异常

在Python任务代码中注册自定义信号处理逻辑,收到SIGTERM时直接抛出异常,让Airflow将任务标记为FAILURE。

示例代码:

import signal
from airflow.operators.python import PythonOperator

def handle_sigterm(signum, frame):
    raise RuntimeError("Task terminated by SIGTERM, marking as failure")

def your_task_logic():
    # 注册SIGTERM信号处理器
    signal.signal(signal.SIGTERM, handle_sigterm)
    # 原有任务逻辑代码
    # ...

task = PythonOperator(
    task_id="tasktest",
    python_callable=your_task_logic,
    dag=dag
)

方法2:重写自定义Operator的on_kill方法

如果使用自定义Operator,可以重写on_kill方法,在任务被终止时主动设置任务状态为失败。

示例代码:

from airflow.models.baseoperator import BaseOperator
from airflow.models.taskinstance import TaskInstance

class CustomTaskOperator(BaseOperator):
    def execute(self, context):
        # 任务执行逻辑
        # ...
    
    def on_kill(self):
        # 获取当前任务实例并强制设置为失败状态
        ti = TaskInstance(task=self, execution_date=self.execution_date)
        ti.set_state("failed", session=self.session)

方法3:在任务逻辑中检查任务状态

通过Airflow上下文获取当前TaskInstance,定期检查任务是否被外部标记为失败,若已终止则主动抛出异常。

示例代码:

from airflow.decorators import task

@task
def your_task(**context):
    ti = context["ti"]
    # 可以在任务逻辑的关键节点添加状态检查
    if ti.state == "failed":
        raise RuntimeError("Task was externally marked as failed")
    # 后续任务逻辑
    # ...

关键原因说明

Airflow默认以任务进程的退出码判定状态:退出码0为SUCCESS,非0则为FAILURE。你的任务收到SIGTERM后,进程未触发异常而是正常退出(exitcode=0),因此被误标记为成功。上述方法的核心都是让任务在收到SIGTERM后,以非0退出码结束或主动设置状态为失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 18:13:15