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
相关产品推荐
相关产品推荐

