PyTorch多进程共享内存耗尽,如何调用torch_shm_manager回收?
PyTorch多进程处理数据时共享内存文件堆积问题解决
问题背景
使用PyTorch多进程并行处理20000+数据点,采用file_system共享策略(因文件描述符超限)和spawn启动方式,进程运行一段时间后崩溃,提示无法分配共享内存,查看/dev/shm发现存在80000+未清理的临时文件。尝试过分块创建Pool、删除返回值等方法无效,改用多线程又受GIL限制导致效率低下。
原始未分块代码:
import torch #I have to use file_system because the number of file descriptors exceeds OS limits torch.multiprocessing.set_sharing_strategy("file_system") from glob import glob def processData(fname): #Tensor magic goes here, a list of tensors is created from the file at fname and returned if __name__ == "__main__": datadir = "/path/to/dir" processedData = [] #I hope to eventually use this with CUDA tensors, but for now I'm having issues with CPU torch.multiprocessing.set_start_method("spawn") p = torch.multiprocessing.Pool(8) for result in p.imap_unordered(processData, glob(datadir+"*.txt")): processedData += result p.close() p.join()
解决方案
1. 显式释放张量的共享内存资源
PyTorch的file_system策略会将张量数据存储在/dev/shm的临时文件中,默认依赖Python垃圾回收,但多进程环境下GC可能不及时。可以在处理完每个结果后,手动调用张量的storage().unlink_shared()方法释放对应文件,再强制触发垃圾回收:
import gc # ... 其他代码 ... for result in p.imap_unordered(processData, glob(datadir+"*.txt")): processedData += result # 逐个释放张量的共享内存文件 for tensor in result: tensor.storage().unlink_shared() del tensor # 强制触发GC回收资源 gc.collect()
2. 改用异步回调处理结果,避免内存驻留
imap_unordered会持续持有结果引用,导致共享内存文件无法被及时删除。可以用apply_async配合回调函数,处理完结果后立即释放资源:
import gc processedData = [] def collect_result(result): global processedData processedData += result # 释放当前结果里的张量资源 for tensor in result: tensor.storage().unlink_shared() del tensor gc.collect() if __name__ == "__main__": datadir = "/path/to/dir" torch.multiprocessing.set_sharing_strategy("file_system") torch.multiprocessing.set_start_method("spawn") p = torch.multiprocessing.Pool(8) for fname in glob(datadir+"*.txt"): p.apply_async(processData, args=(fname,), callback=collect_result) p.close() p.join()
3. 避免跨进程共享张量,转成numpy数组传输
如果不需要在子进程和主进程间共享张量内存,可以在子进程中将张量转为numpy数组返回,主进程再转回张量,这样PyTorch不会创建共享内存文件:
def processData(fname): # 生成张量列表的逻辑 tensor_list = [...] # 转为numpy数组返回 return [t.numpy() for t in tensor_list] # 主进程中处理结果 for result in p.imap_unordered(processData, glob(datadir+"*.txt")): # 转回PyTorch张量 processedData += [torch.from_numpy(arr) for arr in result]
4. 子进程内及时清理中间张量
在processData函数中,处理完单个文件后,立即删除中间生成的临时张量,并触发子进程内的垃圾回收,减少每个子进程生成的共享内存文件数量:
import gc def processData(fname): # 中间处理步骤,生成临时张量 temp_tensor = ... # 生成最终需要返回的张量列表 final_tensors = [...] # 删除中间张量并回收资源 del temp_tensor gc.collect() return final_tensors
5. 临时调整系统共享内存大小(治标方案)
如果以上代码层面的优化仍无法解决,可以临时增大/dev/shm的容量(需root权限):
mount -o remount,size=16G /dev/shm
注意这只是临时解决方案,重启后会恢复默认值,建议优先从代码层面解决内存泄漏问题。
内容的提问来源于stack exchange,提问作者Eric Bell
相关产品推荐
相关产品推荐

