Airflow使用SubDag拉取SFTP文件后EmailOperator找不到附件路径如何解决
问题解决方案
1 路径问题修复
文件找不到问题核心是相对路径不可靠,以及任务运行环境的工作目录不固定:
- Airflow任务的默认工作目录会随部署方式、Worker节点、任务实例变化,
./my_path/my_file.txt这种相对路径无法保证两个任务指向同一个位置 - 修复方案如下:
- 改用绝对路径存储文件,比如使用Airflow工作目录下的固定路径:
/opt/airflow/runtime/my_file.txt,需要提前确保所有对应Worker节点上该路径存在,且Airflow运行用户有读写权限 - 可以临时在下载任务后加一段调试代码确认路径:
check_path = PythonOperator( task_id="check_path", python_callable=lambda: print(f"当前工作目录:{os.getcwd()},目录下文件:{os.listdir('./my_path/')}"), dag=child_dag ) - 额外注意你现有代码里
ssh_conn_id="=my_conn"多了一个等号,属于配置错误,需要修正为你实际的SFTP连接ID
- 改用绝对路径存储文件,比如使用Airflow工作目录下的固定路径:
- 更稳妥的实现方式:将下载+发邮件的逻辑合并为一个PythonOperator任务,完全避免跨Worker路径不一致的问题,也不需要用到已经被Airflow官方废弃的SubDagOperator(Airflow 2.x起推荐用TaskGroup替代SubDag做任务分组)。如果需要保留两个任务的拆分结构,可以给两个任务指定同一个专属队列,保证任务调度到相同的Worker节点执行。
2 基于GCS的更优实现方案
使用GCS共享存储可以完全规避同Worker绑定的要求,稳定性更高,流程如下:
- 第一步用
SFTPToGCSOperator直接将SFTP文件同步到GCS路径,例如gs://your_bucket/temp/my_file.txt - 第二步编写自定义Python任务:先将GCS文件下载到当前任务的本地临时目录,再调用邮件能力发送附件,发送完成后可以按需清理GCS和本地的临时文件
- 示例核心代码参考:
from airflow.providers.google.cloud.hooks.gcs import GCSHook import os import smtplib from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from email.mime.base import MIMEBase from email import encoders def send_email_with_gcs_attachment(**context): # 下载GCS文件到本地临时目录 gcs_hook = GCSHook(gcp_conn_id="your_gcp_conn") local_path = "/tmp/my_file.txt" gcs_hook.download(bucket_name="your_bucket", object_name="temp/my_file.txt", filename=local_path) # 构造邮件 msg = MIMEMultipart() msg['From'] = "sender@example.com" msg['To'] = "user_email@example.com" msg['Subject'] = "Testing onlyyy" msg.attach(MIMEText("blahblah falafel wooooo", "html")) # 添加附件 with open(local_path, "rb") as f: part = MIMEBase('application', 'octet-stream') part.set_payload(f.read()) encoders.encode_base64(part) part.add_header('Content-Disposition', f"attachment; filename= my_file.txt") msg.attach(part) # 也可以直接调用Airflow内置的send_email工具方法简化逻辑 with smtplib.SMTP('your_smtp_host', 587) as server: server.starttls() server.login("your_smtp_user", "your_smtp_pass") server.send_message(msg) # 清理临时文件 os.remove(local_path)
这种方案不需要绑定任务到同一Worker,也不需要在Worker节点留存持久化文件,更适合分布式部署的Airflow集群。
内容的提问来源于stack exchange,提问作者user6308605
相关产品推荐
相关产品推荐

