如何通过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
相关产品推荐
相关产品推荐

