多进程写入npy文件时的锁机制咨询及异常排查
多进程读写npy文件的并发问题与解决办法
我编写了一个多进程并行计算矩阵的程序:每个进程独立计算大型矩阵后,将其追加到从npy文件加载的数组中,最终通过np.save将追加后的数组保存回磁盘。但当多个进程同时读写该文件时,偶尔会抛出异常:
ValueError: cannot reshape array of size 0 into shape (1,10,10)。请问当一个进程写入文件时,是否存在锁机制让其他进程等待锁释放后再进行读写操作?最小可复现示例代码:
import os import time import numpy as np from multiprocessing import Pool def main(): p = Pool(processes=8) files = os.listdir("matrices/dummy") for f in files: try: os.remove(f) except FileNotFoundError: print("File not found while deleting") a = np.random.randn(1, 10, 10) np.save("matrices/dummy/saved_mat.npy", a) p.starmap(save_mat, enumerate(np.random.rand(8))) def save_mat(index, sleep_dur): # time.sleep(sleep_dur) np.random.seed(index) try: a = np.load("matrices/dummy/saved_mat.npy") new_val = np.random.randn(1, 10, 10) a = np.concatenate([a, new_val], axis=0) except FileNotFoundError: a = np.random.randn(1, 10, 10) np.save("matrices/dummy/saved_mat.npy", a) if __name__ == "__main__": main() b = np.load("matrices/dummy/saved_mat.npy") c = 0
核心问题解析
numpy的np.load和np.save没有内置的并发锁机制,当多个进程同时对同一个npy文件进行读写操作时,会出现文件写入被打断的情况:比如进程A正在执行np.save写入文件的过程中,进程B调用np.load读取,此时文件内容可能还未完全写入,导致加载到的是不完整甚至空的数组,进而触发ValueError。
解决办法
方法1:使用进程间锁控制文件读写
通过multiprocessing.Manager创建跨进程的锁,确保同一时间只有一个进程能执行文件的读写操作。修改后的代码如下:
import os import numpy as np from multiprocessing import Pool, Manager def main(): with Manager() as manager: file_lock = manager.Lock() p = Pool(processes=8) # 确保目录存在 os.makedirs("matrices/dummy", exist_ok=True) # 清理目标目录文件 files = os.listdir("matrices/dummy") for f in files: try: os.remove(f"matrices/dummy/{f}") except FileNotFoundError: print("File not found while deleting") # 初始化基础数组 a = np.random.randn(1, 10, 10) np.save("matrices/dummy/saved_mat.npy", a) # 传递锁给每个进程 tasks = [(index, file_lock) for index in range(8)] p.starmap(save_mat, tasks) def save_mat(index, file_lock): np.random.seed(index) new_val = np.random.randn(1, 10, 10) # 加锁后执行读写操作 with file_lock: try: a = np.load("matrices/dummy/saved_mat.npy") a = np.concatenate([a, new_val], axis=0) except FileNotFoundError: a = new_val np.save("matrices/dummy/saved_mat.npy", a) if __name__ == "__main__": main() b = np.load("matrices/dummy/saved_mat.npy") print(f"最终数组形状: {b.shape}") # 应为(9,10,10)
方法2:避免共享文件,用临时文件后合并(更高效)
多进程共享单个文件会导致大量等待,效率低下。更优的方式是让每个进程将计算结果写入独立的临时文件,最后由主进程统一合并所有临时文件:
import os import numpy as np from multiprocessing import Pool def main(): # 初始化目录 os.makedirs("matrices/dummy", exist_ok=True) os.makedirs("matrices/temp", exist_ok=True) # 清空临时文件目录 for f in os.listdir("matrices/temp"): try: os.remove(f"matrices/temp/{f}") except FileNotFoundError: pass # 初始化基础数组 base_arr = np.random.randn(1, 10, 10) np.save("matrices/dummy/saved_mat.npy", base_arr) p = Pool(processes=8) p.map(save_temp_mat, range(8)) # 主进程合并所有临时文件 all_arrs = [base_arr] for f in os.listdir("matrices/temp"): temp_arr = np.load(f"matrices/temp/{f}") all_arrs.append(temp_arr) final_arr = np.concatenate(all_arrs, axis=0) np.save("matrices/dummy/saved_mat.npy", final_arr) # 清理临时文件 for f in os.listdir("matrices/temp"): os.remove(f"matrices/temp/{f}") print(f"最终数组形状: {final_arr.shape}") # 应为(9,10,10) def save_temp_mat(index): np.random.seed(index) new_val = np.random.randn(1, 10, 10) # 写入独立临时文件 np.save(f"matrices/temp/mat_{index}.npy", new_val) if __name__ == "__main__": main()
总结
- 直接共享npy文件会因无内置锁引发并发读写错误,必须手动加锁或调整架构。
- 临时文件合并的方式避免了进程间等待,在计算任务较重时性能更优。
内容的提问来源于stack exchange,提问作者learner
相关产品推荐
相关产品推荐

