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

MWAA中DAG运行间歇性出现FileNotFoundError问题咨询

间歇性FileNotFoundError排查与解决(Airflow DAG上传S3场景)

针对你遇到的间歇性文件找不到问题,结合Airflow的运行机制和你的代码,主要有几个常见原因和对应的解决思路:

1. 多Worker节点的本地文件隔离问题

如果你的Airflow集群用了多个Worker节点,每个Worker有独立的本地文件系统:

  • 当transform_json任务在Worker A生成CSV文件并将本地路径推送到XCom后,s3_file_upload任务可能被调度到Worker B,而B的本地目录中根本没有这个文件,就会触发FileNotFoundError。
  • 这种场景下的表现就是间歇性成功(两个任务被分配到同一个Worker时正常),失败时换Worker重试可能恢复。

解决办法:

  • 改用共享存储:给所有Worker挂载相同的NFS目录,让中间文件存储在共享路径下,确保所有Worker都能访问。
  • 跳过本地文件,直接通过XCom传递文件内容:transform_json任务生成CSV的字符串内容,推送到XCom,s3_file_upload任务从XCom拉取内容后直接写入S3(不需要本地文件),示例代码:
    # transform_json任务
    def transform_json(ti):
        # 生成CSV内容的逻辑
        csv_content = "currency,rate\nUSD,1.0\nEUR,0.92"
        ti.xcom_push(key='csv_content', value=csv_content)
        return csv_content
    
    # s3_file_upload任务
    def s3_file_upload(ti):
        csv_content = ti.xcom_pull(dag_id="exchange_rates_dag", task_ids='transform_json', key='csv_content')
        file_name = f'exchange_rates_{datetime.now().strftime("%Y%m%d%H%M%S")}.csv'
        source_session = boto3.Session(
            aws_access_key_id=Variable.get('access_key_id'),
            aws_secret_access_key=Variable.get('secret_access_key')
        )
        s3_client = source_session.client('s3')
        s3_client.put_object(
            Bucket=Variable.get('s3_bucket'),
            Key=f'exchange_rates/{file_name}',
            Body=csv_content.encode('utf-8')
        )
    
  • 或者直接在transform_json任务中完成S3上传,XCom只传递上传后的S3键,简化流程。

2. 文件未完全写入磁盘就推送XCom

如果transform_json任务中写入CSV文件时,没有确保文件完全flush到磁盘就推送了路径到XCom,可能导致s3_file_upload任务去读取时文件还不存在或不完整。

解决办法:

  • 写入文件时严格使用with上下文管理器,确保文件自动关闭并写入磁盘:
    def transform_json(ti):
        file_path = f'/usr/local/airflow/exchange/exchange_rates_{datetime.now().strftime("%Y%m%d%H%M%S")}.csv'
        with open(file_path, 'w', newline='', encoding='utf-8') as f:
            writer = csv.writer(f)
            # 写入表头和数据
            writer.writerow(['currency', 'rate'])
            writer.writerows([['USD', 1.0], ['EUR', 0.92]])
        # 确保文件存在后再推XCom
        if os.path.exists(file_path):
            ti.xcom_push(key='csv_file_path', value=file_path)
        else:
            raise Exception(f"生成文件失败: {file_path}")
    

3. XCom路径传递的可靠性问题

  • 如果transform_json任务出现部分失败,可能导致XCom中留存了旧的文件路径,后续任务拉取到无效路径报错。
  • 或者多个任务实例的XCom键冲突,拉取到了其他实例的路径。

解决办法:

  • 在s3_file_upload任务中先检查文件是否存在,不存在则抛出异常触发重试:
    def s3_file_upload(ti):
        file_path = ti.xcom_pull(dag_id="exchange_rates_dag", task_ids='transform_json', key='csv_file_path')
        if not os.path.exists(file_path):
            raise FileNotFoundError(f"目标文件不存在: {file_path}")
        # 后续上传逻辑...
    
  • 给DAG添加任务重试机制,比如给s3_file_upload任务设置retries=2,retry_delay=timedelta(seconds=10),让任务在文件不存在时自动重试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 20:12:08