多进程Pool间使用SharedMemory传递大数据时,第二个Pool出现SharedMemory找不到的FileNotFoundError问题
这个问题我之前也踩过坑!核心原因是SharedMemory的生命周期管理漏洞,尤其是在多进程Pool的场景下,子进程退出时会自动清理它创建的SharedMemory,而大数据生成耗时更长,主进程来不及在清理前建立引用,导致后续找不到对应的共享内存文件。
问题根源拆解
- 在
generate_data函数里,子进程创建SharedMemory后调用了shared_memory.close(),随后子进程完成任务退出。这时候如果主进程还没打开这个SharedMemory,操作系统会因为该SharedMemory没有任何活跃的文件描述符引用,自动将其销毁(执行unlink操作)。 - 小数据能正常运行是因为生成速度快,子进程退出后,操作系统还没来得及清理SharedMemory文件,主进程就已经打开它并建立了引用;但大数据生成慢,子进程退出后,操作系统有足够时间完成清理,等主进程通过
p.map()拿到结果再去打开时,共享内存文件已经不存在了。而且p.map()是阻塞式的,要等所有子进程都完成才返回结果列表,这时候所有子进程都已经退出,对应的SharedMemory早就被清理干净了。
两种可行解决方案
方案一:改用imap_unordered实时处理结果,提前建立引用
把第一个Pool的map换成imap_unordered,这样主进程可以在子进程刚返回SharedMemory名字时就立即打开它,增加引用计数,避免子进程退出后SharedMemory被销毁:
from io import BytesIO from multiprocessing import Pool from multiprocessing.shared_memory import SharedMemory def generate_data(_dummy_arg) -> str: data = BytesIO(bytearray(1024 * 1024 * 38)) buf = data.getbuffer() shared_memory = SharedMemory(create=True, size=buf.nbytes) shared_memory.buf[:] = buf shared_memory.close() # 子进程关闭自己的引用,但不主动unlink return shared_memory.name def process_data(data_name: str) -> str: data = SharedMemory(data_name) # 这里添加你的数据处理逻辑 data.close() return "some_result" datas: list[SharedMemory] = [] with Pool(5) as p: # 用imap_unordered实时获取子进程返回的结果,立即打开SharedMemory建立引用 for data_name in p.imap_unordered(generate_data, range(5)): data = SharedMemory(data_name) datas.append(data) # 处理数据并最后清理共享内存 try: with Pool(5) as p: for returned in p.map(process_data, (data.name for data in datas)): print(returned) finally: for data in datas: data.close() data.unlink() # 手动销毁SharedMemory,避免系统资源残留
方案二:主进程提前创建SharedMemory,子进程仅负责写入
这种方式更可控,主进程作为SharedMemory的创建者持有核心引用,子进程只是打开并写入数据,子进程退出后SharedMemory不会被销毁:
from io import BytesIO from multiprocessing import Pool from multiprocessing.shared_memory import SharedMemory def generate_data(args) -> str: shm_name, _dummy_arg = args # 打开主进程预先创建的SharedMemory shm = SharedMemory(shm_name) data = BytesIO(bytearray(1024 * 1024 * 38)) buf = data.getbuffer() shm.buf[:buf.nbytes] = buf shm.close() # 子进程关闭自己的引用 return shm_name def process_data(data_name: str) -> str: data = SharedMemory(data_name) # 这里添加你的数据处理逻辑 data.close() return "some_result" # 主进程提前创建所有需要的SharedMemory,持有核心引用 datas: list[SharedMemory] = [] data_size = 1024 * 1024 * 38 for _ in range(5): shm = SharedMemory(create=True, size=data_size) datas.append(shm) # 第一个Pool负责写入数据 with Pool(5) as p: p.map(generate_data, [(shm.name, i) for i, shm in enumerate(datas)]) # 第二个Pool处理数据 try: with Pool(5) as p: for returned in p.map(process_data, (data.name for data in datas)): print(returned) finally: # 最后统一清理共享内存 for data in datas: data.close() data.unlink()
重要注意事项
- 无论哪种方案,处理完数据后一定要记得
close()并unlink()SharedMemory,否则会在系统中残留共享内存文件,持续占用资源。 - 尽量避免在大数据场景下使用
p.map()这种同步等待所有子进程完成的方法,改用迭代式的imap系列方法可以更早建立引用,避免共享内存被提前清理。
备注:内容来源于stack exchange,提问作者Jordi Torrents
相关产品推荐
相关产品推荐

