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

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
  • 更稳妥的实现方式:将下载+发邮件的逻辑合并为一个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 18:09:00