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

多进程Pool间使用SharedMemory传递大数据时,第二个Pool出现SharedMemory找不到的FileNotFoundError问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.17 12:00:29