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

Airflow中MsSqlOperator执行存储过程无法捕获错误状态的问询

解决MsSqlOperator执行SQL Server作业/存储过程时状态不匹配的问题

问题根源

EXEC dbo.sp_start_job这类SQL语句仅负责提交作业启动请求,只要请求提交成功,MsSqlOperator就会判定任务成功,不会追踪作业后续的实际执行状态。同理,部分存储过程执行失败时不会主动抛出SQL异常,仅返回错误码,也会导致Airflow任务误判为成功。

方案1:修改SQL脚本,内置状态校验(适合快速适配)

编写包含「启动作业→等待完成→校验结果」逻辑的SQL脚本,当作业失败时主动抛出SQL错误,让MsSqlOperator捕获并标记任务失败。

示例SQL脚本(替换ARSHD_USERS为目标作业名):

DECLARE @JobName NVARCHAR(128) = N'ARSHD_USERS';
DECLARE @JobStatus INT;
DECLARE @JobId UNIQUEIDENTIFIER = (SELECT job_id FROM msdb.dbo.sysjobs WHERE name = @JobName);

-- 启动作业
EXEC dbo.sp_start_job @job_name = @JobName;

-- 循环等待作业完成(每10秒检查一次)
WAITFOR DELAY '00:00:10';
WHILE 1=1
BEGIN
    SELECT TOP 1 @JobStatus = current_execution_status
    FROM msdb.dbo.sysjobactivity
    WHERE job_id = @JobId
      AND start_execution_date IS NOT NULL
      AND stop_execution_date IS NULL
    ORDER BY start_execution_date DESC;

    IF @JobStatus IS NULL BREAK; -- 作业已完成
    IF @JobStatus NOT IN (3,4,5) BREAK; -- 非执行/等待状态,退出循环
    WAITFOR DELAY '00:00:10';
END

-- 校验作业最终执行结果
IF EXISTS (
    SELECT 1
    FROM msdb.dbo.sysjobhistory
    WHERE job_id = @JobId
      AND step_id = 0 -- 取作业级别的全局结果
      AND run_status != 1 -- 1=成功,其他为失败/取消/重试
)
BEGIN
    RAISERROR('作业 %s 执行失败', 16, 1, @JobName);
END

将上述SQL作为sql参数传入原MsSqlOperator即可,无需额外任务。

方案2:自定义通用Operator(适合批量维护)

继承MsSqlOperator封装作业/存储过程的执行、状态检查逻辑,实现一次编写、全局复用,避免重复代码。

针对SQL Server作业的自定义Operator

from airflow.providers.microsoft.mssql.operators.mssql import MsSqlOperator
from airflow.providers.microsoft.mssql.hooks.mssql import MsSqlHook
from airflow.exceptions import AirflowException
import time

class MsSqlJobOperator(MsSqlOperator):
    def __init__(self, job_name: str, wait_interval: int = 10, timeout: int = 3600, **kwargs):
        self.job_name = job_name
        self.wait_interval = wait_interval  # 状态检查间隔(秒)
        self.timeout = timeout  # 最大等待时长(秒)
        sql = f"EXEC dbo.sp_start_job N'{job_name}'"
        super().__init__(sql=sql, **kwargs)

    def execute(self, context):
        # 执行启动作业的SQL
        super().execute(context)
        hook = MsSqlHook(mssql_conn_id=self.mssql_conn_id)
        
        # 获取作业ID
        job_id_result = hook.get_first(f"SELECT job_id FROM msdb.dbo.sysjobs WHERE name = N'{self.job_name}'")
        if not job_id_result:
            raise AirflowException(f"未找到作业:{self.job_name}")
        job_id = job_id_result[0]
        
        start_time = time.time()
        # 循环等待作业完成
        while True:
            # 检查是否超时
            if time.time() - start_time > self.timeout:
                raise AirflowException(f"作业 {self.job_name} 执行超时(超过{self.timeout}秒)")
            
            # 查询当前作业状态
            activity = hook.get_first(
                f"""
                SELECT current_execution_status, stop_execution_date
                FROM msdb.dbo.sysjobactivity
                WHERE job_id = '{job_id}'
                  AND start_execution_date IS NOT NULL
                ORDER BY start_execution_date DESC
                """
            )
            
            if not activity:
                raise AirflowException(f"未找到作业 {self.job_name} 的执行记录")
            
            current_status, stop_date = activity
            if stop_date is not None:
                # 作业已完成,校验最终结果
                history = hook.get_first(
                    f"""
                    SELECT run_status
                    FROM msdb.dbo.sysjobhistory
                    WHERE job_id = '{job_id}'
                      AND step_id = 0
                    ORDER BY run_date DESC, run_time DESC
                    """
                )
                if not history or history[0] != 1:
                    raise AirflowException(f"作业 {self.job_name} 执行失败,状态码: {history[0] if history else '未知'}")
                break
            
            time.sleep(self.wait_interval)

使用自定义Operator

run_sp1 = MsSqlJobOperator(
    task_id='Run_SP1',
    mssql_conn_id='mssql_test_conn_msdb',
    job_name='ARSHD_USERS',
    autocommit=True,
    wait_interval=15,  # 每15秒检查一次状态
    timeout=7200  # 最长等待2小时
)

针对存储过程的自定义Operator(可选)

如果是执行普通存储过程,可封装返回码校验逻辑:

class MsSqlStoredProcOperator(MsSqlOperator):
    def __init__(self, proc_name: str, proc_params: dict = None, **kwargs):
        self.proc_name = proc_name
        self.proc_params = proc_params or {}
        # 构造带返回码校验的执行SQL
        param_str = ", ".join([f"@{k} = {repr(v)}" for k, v in self.proc_params.items()])
        sql = f"""
        DECLARE @RC INT;
        EXEC @RC = {self.proc_name} {param_str};
        IF @RC != 0
        BEGIN
            RAISERROR('存储过程 {self.proc_name} 执行失败,返回码: %d', 16, 1, @RC);
        END
        """
        super().__init__(sql=sql, **kwargs)

状态码说明

SQL Server作业/存储过程的关键状态码:

  • sysjobhistory.run_status:1=成功,0=失败,2=重试,3=取消,4=运行中
  • 存储过程返回码:通常0表示成功,非0表示失败(需根据实际存储过程定义调整)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 16:01:02