Airflow worker是否共享文件系统?并发同任务会不会访问同一文件?
结论
- 若两个任务实例确实分配到不同Worker节点,且Worker的
/temp目录为节点本地存储、未做跨节点共享挂载,不会出现引用同个文件的问题:不同Worker的本地文件系统完全隔离,各自的下载、执行、删除操作仅会操作本节点的文件,不会互相影响。 - 但该实现仍然存在严重的并发冲突风险:如果两个任务实例被调度到同一个Worker节点执行,必然会出现路径冲突:
- 先启动的实例下载完成后,后启动的实例可能直接覆盖
/temp/script.py的内容,导致先启动的实例运行spark任务时实际使用了被覆盖的异常文件 - 可能出现某一个实例执行完成后提前删除文件,另一个实例调用
spark-submit时找不到对应路径的问题 - 若文件下载耗时较长,还可能出现文件下载到一半就被另一个实例的spark任务读取执行的问题
- 先启动的实例下载完成后,后启动的实例可能直接覆盖
修复方案
不要使用固定文件名存储临时文件,为每个任务实例生成唯一的临时路径即可解决冲突,参考改造逻辑:
import tempfile import os def python_task_callback(**context): # 自动生成唯一临时文件夹,天然避免路径冲突 with tempfile.TemporaryDirectory() as temp_dir: file_path = os.path.join(temp_dir, 'script.py') download_file(file_name='script.py', save_path=file_path) spark_submit(path=file_path) # 无需手动删除,上下文退出时会自动清理整个临时目录
如果业务要求必须使用固定路径存储,可以在操作文件前加进程级独占文件锁避免读写冲突,但改造成本远高于直接使用唯一临时路径。
内容的提问来源于stack exchange,提问作者Shivansh Narayan
相关产品推荐
相关产品推荐

