当前位置:首页 > 文章列表 > 文章 > python教程 > Python multiprocessing shared_memory 管理共享缓冲区

Python multiprocessing shared_memory 管理共享缓冲区

来源:17golang原创 2026-10-03 23:45:44 0浏览 收藏

在一次本机多进程数据处理原型里,我真正想解决的不是“进程之间怎么传一个对象”,而是“同一大块字节怎样只保留一份,让多个进程按约定读取”。multiprocessing.shared_memory 很适合这个边界:创建者分配共享块,其他进程凭名称附加,数据通过 memoryview 直接访问;但容量、数据布局、同步和清理责任必须由应用自己定义。

官方文档:https://docs.python.org/3/library/multiprocessing.shared_memory.html

最实用的管理原则可以先记成四句话:共享块只有一个生命周期拥有者;每个进程只关闭自己的句柄;底层块只执行一次 unlink();并发读写必须另外使用锁、事件或分区协议。把这四点固定下来,共享缓冲区才不会从“少一次复制”变成难排查的资源泄漏和数据竞争。

什么时候值得从 Queue 或 Pipe 换到共享缓冲区

Queue、Pipe 和进程池参数的优势是语义简单:发送方交出一个对象,接收方得到一个独立对象。代价是复杂对象通常要经过序列化、传输和重新构造。当数据体很大、同一数据会被多个 worker 重复读取,或者生产者会频繁刷新固定大小的连续缓冲区时,这些复制更容易成为架构瓶颈。

我会在下面三个条件同时满足时考虑共享内存:

  • 数据只在同一台机器上的进程之间共享,不需要跨主机。
  • 数据能够表示为稳定的字节布局,例如图像、数组、模型中间张量或定长记录。
  • 团队愿意明确管理名称、容量、偏移、同步和清理,而不是把它当作普通 Python 列表。

小对象、低频消息、动态对象图和需要天然排队语义的任务,继续使用 Queue 往往更合适。共享内存优化的是数据面,不会自动替代任务调度、背压和错误传播。

原架构的瓶颈不只是复制次数

刚开始使用共享内存时,我最不适应的是:它只提供一块可跨进程访问的字节区域,不知道里面是图片、整数数组还是协议帧。SharedMemory.buf 返回的是 memoryview,而对象的形状、数据类型、已用长度和版本都要由应用保存。

问题消息传递通常怎么处理共享内存需要怎么处理
数据边界序列化格式携带长度自己定义长度头和 payload 区
对象类型反序列化恢复对象另传 dtype、shape 或协议版本
并发安全队列天然组织消息使用 Event、Lock 或固定分区
资源释放队列对象随进程管理区分 close 与 unlink

因此,不能只把 Queue.put(data) 替换成写 shm.buf。真正的迁移是把数据布局和生命周期从序列化库手里接回来。

新架构:数据留在共享区,进程只传名称

一个容易维护的最小结构是“长度头 + payload”。前四个字节保存有效数据长度,后面是固定容量的数据区。主进程是唯一拥有者,负责创建和最终删除;工作进程通过名称附加,只在完成后关闭自己的句柄。

主进程、共享内存名称、长度头、payload 和多个工作进程之间的静态结构关系
图1:大块数据保存在共享缓冲区,进程之间只需传递名称和控制信号;这是静态结构图。

下面的例子使用一个 Event 表示数据就绪,另一个 Event 表示消费完成。事件不是为了“让共享内存生效”,而是为了避免 worker 在长度头和 payload 写完之前读取。

import struct
from multiprocessing import Event, Process
from multiprocessing.shared_memory import SharedMemory

# 使用网络字节序的无符号整数保存有效 payload 长度。
HEADER = struct.Struct("!I")


def consume(name: str, total_size: int, ready: Event, done: Event) -> None:
    shm = None
    try:
        # 工作进程只按名称附加,不负责创建或删除共享块。
        shm = SharedMemory(name=name)
        ready.wait()

        (used,) = HEADER.unpack_from(shm.buf, 0)
        capacity = total_size - HEADER.size
        if used > capacity:
            raise ValueError(f"共享缓冲区长度头越界: {used} > {capacity}")

        # 切片仍是共享 memoryview;复制成 bytes 后即可安全脱离共享区。
        view = shm.buf[HEADER.size:HEADER.size + used]
        try:
            local_payload = bytes(view)
        finally:
            view.release()

        # 这里替换为实际计算;示例只保留可观察的本地结果。
        print(f"worker received {len(local_payload)} bytes")
    finally:
        if shm is not None:
            # 每个访问者只关闭自己的句柄。
            shm.close()
        done.set()


def main() -> None:
    payload = b"shared-memory-payload"
    capacity = 4096
    ready = Event()
    done = Event()

    # 创建者为长度头和 payload 一次性分配固定容量。
    shm = SharedMemory(create=True, size=HEADER.size + capacity)
    worker = Process(
        target=consume,
        args=(shm.name, shm.size, ready, done),
    )
    worker.start()

    try:
        if len(payload) > capacity:
            raise ValueError("payload 超过共享缓冲区容量")

        shm.buf[HEADER.size:HEADER.size + len(payload)] = payload
        HEADER.pack_into(shm.buf, 0, len(payload))
        ready.set()

        if not done.wait(timeout=30):
            raise TimeoutError("worker 未在约定时间内完成")
        worker.join(timeout=5)
    finally:
        # 异常时也解除等待,避免子进程永久阻塞。
        ready.set()
        if worker.is_alive():
            worker.terminate()
            worker.join()

        shm.close()
        # 只有拥有者删除底层共享内存块,而且全局只调用一次。
        shm.unlink()


if __name__ == "__main__":
    # spawn 环境要求进程入口放在 main 保护中。
    main()

这里有一个看似多余但很关键的细节:worker 先把共享切片复制成 bytes,再释放切片视图并关闭句柄。如果后续算法能在共享区上原地工作,可以不做这次局部复制,但必须保证所有派生的 memoryview 在 close() 前不再使用。

多 worker 扩展时先选读写模型

把一个 worker 扩成多个 worker,不是简单地重复传同一个名称。需要先决定共享区属于哪一种模型:

  • 一次写、多次读:主进程写完后广播就绪事件,多个 worker 只读。这是最容易控制的模式。
  • 分区写:每个 worker 获得不重叠的 offset 和 length,避免写入同一字节区。
  • 共享更新:多个进程可能修改同一字段,必须用 Lock、信号量或版本协议保护。
  • 双缓冲:生产者写备用区,消费者读当前区,切换时只更新一个受保护的状态标记。

我更偏向“固定分区 + 小控制消息”。大数据仍留在共享区,队列只传 (name, offset, length, version) 这类元数据。这样既保留消息队列的任务边界,又避免把整个 payload 重复序列化。

from dataclasses import dataclass


@dataclass(frozen=True)
class SharedSlice:
    # 名称定位共享块,offset 和 length 定位当前任务的数据范围。
    name: str
    offset: int
    length: int
    version: int


def validate_slice(item: SharedSlice, total_size: int) -> None:
    # 同时检查负数和越界,禁止 worker 访问未分配区域。
    if item.offset  total_size:
        raise ValueError("共享切片超过缓冲区边界")

控制结构尽量保持不可变也很有帮助。worker 收到任务后先验证范围和版本,再创建视图;不要让一个过期任务继续读取已经被下一批数据覆盖的区域。

生命周期:close 与 unlink 不是一回事

close() 关闭当前 SharedMemory 实例持有的文件描述符或句柄,但不保证底层共享块立即消失。unlink() 删除底层共享块,在所有进程中只应调用一次。两者可以按任意顺序调用,但一旦执行 unlink(),其他进程继续访问该块在不同平台上可能产生错误。

创建者、工作进程句柄、close、unlink、resource tracker 与 track 参数之间的静态关系
图2:每个访问者关闭自己的句柄,拥有者只执行一次 unlink;这是静态关系图。

我会把责任写进代码结构,而不是留在注释里:

  • 创建共享块的进程是 owner,负责最终 unlink()。
  • 任何附加进程都是 borrower,只调用自己的 close()。
  • owner 必须等待所有 borrower 完成,或者明确接受强制回收造成的失败。
  • 异常路径与正常路径使用同一个 finally 清理出口。

如果共享块很多,或者生命周期完全跟随一个作业,可以使用 SharedMemoryManager。它启动专门的管理进程,通过上下文管理器退出时,会对其创建的共享块执行清理。代价是多一个管理进程和一层控制面,因此不一定适合每个短任务。

from multiprocessing.managers import SharedMemoryManager


def allocate_job_buffers() -> None:
    # 上下文退出时,管理器统一释放它创建的共享块。
    with SharedMemoryManager() as manager:
        input_block = manager.SharedMemory(size=1024 * 1024)
        output_block = manager.SharedMemory(size=1024 * 1024)

        # 实际任务只传名称和数据布局,不把 SharedMemory 当普通对象复制。
        print(input_block.name, output_block.name)

独立进程要特别注意 resource tracker

Python 3.13 为 SharedMemory 增加了 track 参数。在同一个 multiprocessing 进程家族中,各进程通常共享资源跟踪器,默认 track=True 能帮助清理异常遗留资源。

问题出现在互不相关的独立 Python 进程:它们可能各自拥有资源跟踪器。某个附加进程先退出时,其跟踪器可能删除仍被其他进程使用的共享块。若系统已经有一个明确的外部拥有者负责生命周期,Python 3.13 及以上的独立附加进程可以考虑设置 track=False。Windows 使用自身的句柄跟踪,文档说明该参数会被忽略。

from multiprocessing.shared_memory import SharedMemory


def attach_from_standalone_process(name: str) -> SharedMemory:
    # Python 3.13+:已有外部拥有者负责清理时,避免独立 tracker 抢先删除。
    return SharedMemory(name=name, track=False)

这不是通用推荐值。由同一 multiprocessing 程序创建的 worker 通常保持默认跟踪即可;只有确认进程彼此独立,并且另有可靠 owner 时,才应关闭跟踪。还要把最低 Python 版本写进项目要求,因为 3.12 及更早版本没有这个参数。

上线后的信号应该看什么

共享内存架构上线后,我不会只看 CPU 时间。更有价值的是观察资源和协议是否稳定:

  • 创建中的共享块数量是否有上限,任务完成后能否回落。
  • worker 超时、异常退出后,owner 是否仍能完成清理。
  • 长度头、offset、length 和 version 越界是否被拒绝。
  • 是否出现 resource tracker 警告、重复 unlink 或找不到名称。
  • 多写者场景是否存在没有锁或没有分区的重叠区域。
  • 缓冲区容量是否按峰值无限增长,而没有复用和回收策略。

性能结果必须用自己的数据规模、进程数和平台测量。共享内存减少了某些数据复制,但 worker 如果立刻执行 bytes(view) 或 NumPy 拷贝,仍然会产生本地复制;同步等待、缓存局部性和内存带宽也可能成为新的瓶颈。

下一轮改进的检查清单

检查项推荐做法
拥有者只有创建者或管理器负责 unlink
访问者每个进程关闭自己的句柄和派生视图
容量写入前校验 length,不依赖切片静默截断
布局固定头部、版本、offset 和 payload 约定
同步使用 Event、Lock、Semaphore 或不重叠分区
异常超时、终止、finally 清理路径都可达
独立进程Python 3.13+ 按实际拥有者评估 track
适用范围只用于同机进程,不当作分布式共享内存

对我来说,shared_memory 最有价值的不是“更底层”,而是把大数据的存放位置和任务控制解耦。只要先定清楚谁创建、谁读写、谁关闭、谁删除,再决定用锁、事件还是分区,SharedMemory 就能成为可复用的进程间数据层,而不是一块没人敢回收的匿名缓冲区。

常见延伸问题

  • SharedMemory 能跨机器吗? 不能,它面向同一台机器上的进程;跨主机需要网络协议或分布式存储。
  • 多个进程同时写会自动加锁吗? 不会,需要应用自己提供锁、信号量或不重叠分区。
  • 只调用 close 会释放内存吗? 不一定。非 Windows 平台通常还需要由拥有者调用一次 unlink。
  • 什么时候使用 SharedMemoryManager? 当一组共享块的生命周期跟随同一个作业,并希望统一清理时更合适。
版本声明
本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
动漫岛“全站免费无广”是真的吗?宣传语、服务范围与核验边界动漫岛“全站免费无广”是真的吗?宣传语、服务范围与核验边界
上一篇
动漫岛“全站免费无广”是真的吗?宣传语、服务范围与核验边界
Go archive/tar 流式写入文件元数据
下一篇
Go archive/tar 流式写入文件元数据
查看更多
最新文章
查看更多
课程推荐
  • 前端进阶之JavaScript设计模式
    前端进阶之JavaScript设计模式
    设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
    543次学习
  • GO语言核心编程课程
    GO语言核心编程课程
    本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
    516次学习
  • 简单聊聊mysql8与网络通信
    简单聊聊mysql8与网络通信
    如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
    500次学习
  • JavaScript正则表达式基础与实战
    JavaScript正则表达式基础与实战
    在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
    487次学习
  • 从零制作响应式网站—Grid布局
    从零制作响应式网站—Grid布局
    本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
    485次学习
查看更多
AI推荐
  • PubMedQA数据集详解:生物医学问答基准、功能与应用指南
    PubMedQA
    深入了解PubMedQA生物医学问答数据集,涵盖其核心功能、使用方法及在临床决策、药物研发等场景的应用,助力提升NLP模型性能。
    318次使用
  • H2O EvalGPT:开源LLM大模型评估与排行榜工具
    H2O EvalGPT
    H2O EvalGPT是H2O.ai推出的开源LLM评估平台,提供详细的大模型性能排行榜、行业特定基准测试及A/B测试功能,助您快速选择最适合项目的高性能大语言模型。
    374次使用
  • LMArena是什么?伯克利AI模型评估平台使用指南与功能解析
    LMArena
    LMArena是加州大学伯克利分校推出的AI模型匿名评测平台。通过盲测投票机制,用户可对比不同大模型回答并生成实时排行榜,助力开发者优化模型及用户选择最佳AI工具。
    371次使用
  • 斯坦福HELM:大语言模型Holistic Evaluation整体评估框架详解
    HELM
    深入了解斯坦福推出的HELM(Holistic Evaluation of Language Models)大模型评测体系。本文解析其核心功能、安装配置步骤及应用场景,涵盖准确性、公平性、鲁棒性等多维度指标,助力开发者全面优化语言模型性能。
    337次使用
  • MMBench详解:多模态大模型基准测试、功能特点与使用指南
    MMBench
    MMBench是由上海人工智能实验室等机构联合推出的多模态基准测试平台,提供细粒度能力评估、大规模数据集及VLMEvalKit工具。本文详细介绍其核心功能、安装使用方法及应用场景,助力开发者全面评估多模态模型性能。
    162次使用
微信登录更方便
  • 密码登录
  • 注册账号
登录即同意 用户协议 和 隐私政策
返回登录
  • 重置密码