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

多进程写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 01:35:25