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

Airflow DAG任务SFTP连接失败但DAG仍显示成功的修复请求

问题分析

你的问题核心在于:Airflow的SSHOperator仅以远程执行命令的退出码作为任务成功/失败的判断依据。从日志看,远程命令docker logs --follow shopping-PROJECT-scheduler-stage本身执行完成后返回了退出码0,因此Airflow标记任务为成功;但该命令运行过程中,内部调用的SFTP客户端连接失败(错误码ERR_GENERIC_CLIENT),这个错误没有被转化为远程命令的非0退出码,导致Airflow无法感知到实际错误。

解决方案

针对这个问题,有三种可行的修改方向,按实现复杂度从低到高排序:

1. 修改远程命令/脚本,让SFTP错误返回非0退出码

如果远程执行的是自定义脚本(比如你的node.js脚本),在捕获到SFTP连接错误时,主动设置非0退出码:

  • 对于node.js脚本,在错误处理逻辑中添加:
    if (err.code === 'ERR_GENERIC_CLIENT') {
      console.error('SFTP连接失败:', err);
      process.exit(1); // 返回非0退出码触发Airflow任务失败
    }
    
  • 如果是shell命令组合,在命令中加入错误检查逻辑:
    # 执行命令并捕获所有输出
    docker logs --follow shopping-PROJECT-scheduler-stage > /tmp/task_output 2>&1
    # 检查输出中是否包含错误关键词,存在则返回非0退出码
    if grep -q "ERR_GENERIC_CLIENT\|All configured authentication methods failed" /tmp/task_output; then
        exit 1
    fi
    

2. 调整SSHOperator的命令,添加输出检查

直接在Airflow的SSHOperator中修改command参数,让远程命令自动检查输出并返回对应退出码:

from airflow.providers.ssh.operators.ssh import SSHOperator

ssh_task = SSHOperator(
    task_id='ssh_operator_remote5',
    ssh_conn_id='arflw_ssh_remote_ec2',
    command="""
        docker logs --follow shopping-PROJECT-scheduler-stage > /tmp/task_log 2>&1
        grep -q "ERR_GENERIC_CLIENT\|All configured authentication methods failed" /tmp/task_log && exit 1 || exit 0
    """,
    dag=dag
)

这个命令会将日志输出到临时文件,检查是否包含指定错误关键词,存在则返回1,Airflow会标记任务失败。

3. 自定义Operator(进阶)

如果需要更灵活的错误判断逻辑,可以继承SSHOperator,重写execute方法,在获取远程输出后主动检查错误:

from airflow.providers.ssh.operators.ssh import SSHOperator
from airflow.exceptions import AirflowException

class CustomSSHOperator(SSHOperator):
    def execute(self, context):
        # 执行原SSHOperator的逻辑
        result = super().execute(context)
        # 检查输出中是否包含错误关键词
        error_keywords = ['ERR_GENERIC_CLIENT', 'All configured authentication methods failed']
        for keyword in error_keywords:
            if keyword in result.output:
                raise AirflowException(f"任务执行失败,检测到错误: {keyword}")
        return result

# 使用自定义Operator
ssh_task = CustomSSHOperator(
    task_id='ssh_operator_remote5',
    ssh_conn_id='arflw_ssh_remote_ec2',
    command="docker logs --follow shopping-PROJECT-scheduler-stage",
    dag=dag
)
验证方法

修改后重新触发DAG,当SFTP连接再次失败时:

  • 远程命令会返回非0退出码,或者自定义Operator抛出异常
  • Airflow会将ssh_operator_remote5任务标记为FAILED,整个DAG运行状态也会变为失败

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 11:59:59