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

如何通过Apache Airflow操作SFTP删除文件?自定义Operator实现方案

关于Apache Airflow删除SFTP文件的解决方案

首先明确:原生的SFTPOperator确实只支持put和get操作,你代码里写的operation="del"是无效的,官方并没有提供这个选项。下面给你两种可行的解决办法:

方案一:用PythonOperator结合SFTPHook(最简单)

不需要自定义Operator,直接用PythonOperator调用SFTPHook的方法来删除文件,这是最快捷的方式。示例代码如下:

from airflow.providers.ssh.hooks.ssh import SSHHook
from airflow.operators.python import PythonOperator
from airflow.exceptions import AirflowException

def delete_sftp_file(**context):
    # 获取XCom中的文件名
    del_td = context['task_instance'].xcom_pull(task_ids=None, key="del_td")
    remote_filepath = f'remote/path/.{del_td}.i01_1.spn'
    
    # 初始化SFTPHook
    ssh_hook = SSHHook(ssh_conn_id='nseix')
    sftp_client = ssh_hook.get_conn().open_sftp()
    
    try:
        # 删除远程文件
        sftp_client.remove(remote_filepath)
        print(f"成功删除文件: {remote_filepath}")
    except FileNotFoundError:
        raise AirflowException(f"要删除的文件不存在: {remote_filepath}")
    except Exception as e:
        raise AirflowException(f"删除文件失败: {str(e)}")
    finally:
        # 关闭SFTP客户端和SSH连接
        sftp_client.close()
        ssh_hook.close_conn()

# 定义任务
delete_old_span_to_ftp = PythonOperator(
    task_id="delete_old_span_to_ftp",
    python_callable=delete_sftp_file,
    provide_context=True,
    dag=dag
)

方案二:自定义SFTP删除Operator

如果需要在多个DAG中复用删除逻辑,可以自定义一个Operator,继承BaseOperator:

from airflow.models.baseoperator import BaseOperator
from airflow.providers.ssh.hooks.ssh import SSHHook
from airflow.utils.decorators import apply_defaults
from airflow.exceptions import AirflowException

class SFTPDeleteOperator(BaseOperator):
    @apply_defaults
    def __init__(
        self,
        ssh_conn_id: str,
        remote_filepath: str,
        *args, **kwargs
    ):
        super().__init__(*args, **kwargs)
        self.ssh_conn_id = ssh_conn_id
        self.remote_filepath = remote_filepath

    def execute(self, context):
        # 渲染远程文件路径(处理XCom模板)
        rendered_remote_path = self.render_template(self.remote_filepath, context)
        
        ssh_hook = SSHHook(ssh_conn_id=self.ssh_conn_id)
        sftp_client = ssh_hook.get_conn().open_sftp()
        
        try:
            sftp_client.remove(rendered_remote_path)
            self.log.info(f"已成功删除SFTP文件: {rendered_remote_path}")
        except FileNotFoundError:
            raise AirflowException(f"SFTP文件不存在: {rendered_remote_path}")
        except Exception as e:
            raise AirflowException(f"删除SFTP文件出错: {str(e)}")
        finally:
            sftp_client.close()
            ssh_hook.close_conn()

# 使用自定义Operator
delete_old_span_to_ftp = SFTPDeleteOperator(
    task_id="delete_old_span_to_ftp",
    ssh_conn_id='nseix',
    remote_filepath='remote/path/.{{ task_instance.xcom_pull(task_ids=None, key="del_td") }}.i01_1.spn',
    dag=dag
)

关于FTPHook的说明

如果你使用的是FTP服务器而非SFTP,可以用FTPHook;但你这里用的是SFTP(基于SSH的文件传输),FTPHook并不适用,必须用SFTPHook(即上面代码中的SSHHook,因为Airflow的SFTP功能是基于SSHHook实现的)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 17:52:37