登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  python教程

multiprocessing 传输大对象为何变慢,如何减少序列化

来源:17golang原创

时间:2026-10-08 11:14:12 478浏览 收藏

multiprocessing 传输大对象变慢,通常不是进程数不够,而是每个任务都在支付序列化、管道传输、反序列化和内存复制成本。先测 pickle 体积与端到端时间;确认数据搬运占主导后,把“每次传完整对象”改成“数据只放一次,任务只传名称、索引和小结果”。

Python multiprocessing 官方文档:https://docs.python.org/3.14/library/multiprocessing.html

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

先看结论:优化数据通道而不是只加进程

进程池适合把 CPU 计算分散到多个解释器进程,但不同进程拥有各自的地址空间。把大 bytes、大列表或嵌套字典作为任务参数,并不等于把一个轻量引用交给 worker。Pool 底层要把对象变成可传输字节,再在另一端重建对象。

因此迁移方向可以浓缩成三条:

  • 减少跨进程的字节数:只传 ID、文件位置、共享内存名称和切片范围。
  • 减少任务数量:让一次任务处理一批数据,而不是为每个小元素排队。
  • 减少返回体积:在 worker 内完成聚合或过滤,只返回计数、索引或小结果。

如果单个任务计算很短,而 payload 很大,加进程反而可能更慢。优化前必须把“计算耗时”和“数据通道耗时”拆开,不应只比较进程数量。

为什么会慢:Pool 参数不是共享引用

Python 官方文档说明,放入 multiprocessing.Queue 的对象都会用 pickle 序列化,get() 得到的是重新创建的对象,并不与原对象共享内存。对象进入队列后,后台 feeder 线程还会把 pickle 字节写入底层管道。因此,只测一次 put() 调用的返回时间,可能看不到全部传输成本。

父进程大对象与 pickle、任务队列、进程管道、子进程副本和计算函数之间的静态依赖结构
图1:大对象传输结构图。Pool 任务参数跨过进程边界时依赖 pickle 与通信通道,子进程接收的是重新创建的对象,而不是父进程中的同一引用。

大对象的成本至少包含四部分:父进程执行 pickle、通信通道搬运字节、子进程执行 unpickle、两端分配和复制内存。若 worker 最后又返回一个大对象,这条成本链还会反向再走一遍。

可以先测对象自身的序列化体积与耗时。下面的代码不代表完整 Pool 成本,但能快速判断 payload 是否值得继续通过队列发送:

import pickle
import time


def inspect_pickle_cost(obj):
    # 使用最高协议测量当前解释器下的序列化体积与耗时。
    started = time.perf_counter()
    payload = pickle.dumps(obj, protocol=pickle.HIGHEST_PROTOCOL)
    elapsed = time.perf_counter() - started
    return len(payload), elapsed


if __name__ == "__main__":
    sample = bytes(range(256)) * 200_000
    size, seconds = inspect_pickle_cost(sample)
    # 输出由当前机器实测,避免把示例环境数字当成通用结论。
    print({"pickle_bytes": size, "pickle_seconds": seconds})

检查点是:序列化字节数是否接近原始数据大小,单次 pickle 是否已经接近任务计算时间,以及同一对象是否在多个任务中被重复序列化。命中其中一项,就应优先改数据布局。

旧代码受影响的位置

下面这种写法看起来已经把数据分块,实际同时制造了父进程切片副本,并把每个大切片作为独立任务 pickle:

import multiprocessing as mp


def sum_payload(chunk):
    # worker 接收到的是反序列化后的 chunk 副本。
    return sum(chunk)


def slow_parallel_sum(data, chunk_size):
    # bytes 切片会创建新对象,随后还要通过进程池序列化传输。
    chunks = [data[i:i + chunk_size] for i in range(0, len(data), chunk_size)]
    with mp.Pool() as pool:
        partials = pool.map(sum_payload, chunks)
    return sum(partials)


if __name__ == "__main__":
    source = bytes(range(256)) * 200_000
    print(slow_parallel_sum(source, 1_000_000))

问题不在 map 这个名字,而在任务消息的内容。即使设置更大的 chunksize,也只是减少任务调度批次,不会让大 payload 自动变成共享引用。真正要问的是:worker 是否必须拥有完整 Python 对象,还是只需读取一段只读字节。

另外,大结果回传同样会触发序列化。能在 worker 中做 sum、计数、哈希、过滤或局部聚合时,就不要把全部中间数据送回父进程。

迁移方案:从传对象改成传描述符

multiprocessing.shared_memory 允许不同进程访问同一块共享内存。父进程只需把大数据写入一次;Pool 任务传递共享内存名称以及 (start, end) 两个整数,worker 附加到同一数据块并创建切片视图。

import multiprocessing as mp
from multiprocessing import shared_memory


def sum_shared_slice(task):
    name, start, end = task
    # 子进程通过名称附加到已有共享内存,不接收完整大对象。
    shm = shared_memory.SharedMemory(name=name)
    view = shm.buf[start:end]
    try:
        return sum(view)  # 只把一个小整数返回父进程。
    finally:
        # 先释放 memoryview,再关闭当前进程持有的句柄。
        del view
        shm.close()


def parallel_sum(data, chunk_size):
    shm = shared_memory.SharedMemory(create=True, size=len(data))
    try:
        shm.buf[:len(data)] = data  # 大数据只复制进共享区一次。
        tasks = [
            (shm.name, start, min(start + chunk_size, len(data)))
            for start in range(0, len(data), chunk_size)
        ]

        # 使用当前平台默认上下文;任务函数必须定义在可导入模块顶层。
        ctx = mp.get_context()
        with ctx.Pool() as pool:
            partials = pool.map(sum_shared_slice, tasks)
        return sum(partials)
    finally:
        # 父进程负责关闭并最终删除共享内存名称。
        shm.close()
        shm.unlink()


if __name__ == "__main__":
    source = bytes(range(256)) * 200_000
    print(parallel_sum(source, 1_000_000))
SharedMemory、共享内存名称、切片索引、Pool 任务、worker 视图和小结果之间的静态数据结构
图2:共享内存描述符结构图。大数据保留在 SharedMemory 中,Pool 任务只携带名称和切片索引,worker 通过 memoryview 读取并仅返回小结果。

这个改写的关键不是“SharedMemory 一定最快”,而是把每个任务的消息从大字节块缩小成一个短名称和两个整数。父进程仍有一次写入共享内存的成本,worker 仍有附加与读取成本,所以它更适合一份大数据被多个任务复用,或计算量足以覆盖共享资源管理开销的场景。

close() 表示当前进程不再使用该句柄,unlink() 表示整块共享内存不再需要。官方文档要求在所有进程不再使用时最终调用 unlink();把清理放进父进程的 finally,可以减少异常路径上的泄漏。

最小验证:分别测计算与传输

不要只跑一次总耗时就下结论。准备同一份输入,至少比较三个基线:串行计算、原始 Pool 传大对象、SharedMemory 传描述符。每种方案重复多次并核对结果完全一致。

import time


def timed(label, func, *args):
    # 只记录端到端耗时,包含任务调度、传输和结果收集。
    started = time.perf_counter()
    result = func(*args)
    elapsed = time.perf_counter() - started
    print({"label": label, "seconds": elapsed})
    return result


if __name__ == "__main__":
    data = bytes(range(256)) * 200_000
    expected = timed("serial", sum, data)
    actual = timed("shared_memory", parallel_sum, data, 1_000_000)
    # 性能结论之前先验证业务结果,防止切片遗漏或重复。
    assert actual == expected

检查时同时记录数据大小、任务数、块大小、进程数、启动方法和平台。Python 3.14 起,fork 已不再是任何平台的默认启动方法;POSIX 的默认选择也发生了变化。不要把“Linux 上 fork 后好像不复制”的偶然表现当成跨平台数据协议,尤其不要在库代码中偷偷强制全局启动方法。

如果提高块大小后明显改善,说明任务调度和序列化次数占比较高;如果 SharedMemory 改写仍没有收益,可能是数据只用一次、计算太短、内存带宽已成为瓶颈,或创建与附加共享内存的固定成本超过节省量。

边界选择:共享内存、线程还是 Manager

方案适合场景主要代价
小参数 + Pool任务计算重,输入输出都小仍需 pickle,但成本可控
SharedMemory本机进程反复读取大块规则数据需要切片协议、同步和生命周期清理
文件映射数据本来就在文件中,需要分段读取受文件格式与 I/O 模式影响
线程池I/O 密集,或底层扩展在计算时释放 GIL纯 Python CPU 代码通常受 GIL 限制
Manager 代理需要共享复杂 Python 对象和较低访问频率官方文档明确其通常比共享内存慢

共享内存不意味着可以忽略并发写入。如果多个 worker 修改重叠区域,仍要设计不重叠分区、锁或单写者协议。只读大数据最容易获得收益,也最容易验证正确性。

最终采用建议是:默认保持任务消息小而明确;大对象只读且反复复用时考虑 SharedMemory;数据天然在磁盘时考虑内存映射;底层工作会释放 GIL 时先评估线程。不要为了避开 pickle 把所有状态都改成 Manager,也不要靠固定启动方法掩盖数据通道设计问题。

相关问题

调大 Pool.map 的 chunksize 能解决大对象序列化吗?

它能减少任务调度批次,但不会让任务参数自动共享内存。若每个任务仍包含完整大对象,pickle 与复制成本依然存在。

把大对象放到全局变量就一定不会复制吗?

不一定。效果依赖启动方法和平台;spawn、forkserver 与 fork 的继承行为不同。可移植代码应使用明确的数据协议,而不是假设子进程总能继承父进程内存。

SharedMemory 为什么还可能没有串行快?

共享内存只减少数据序列化与复制,不会消除建池、调度、附加句柄、内存访问和合并结果的成本。计算量太小或数据只读一次时,串行通常更简单。

如何避免共享内存泄漏?

每个进程结束访问时调用 close(),最后由明确的所有者调用一次 unlink()。父进程应把清理放入 finally,并避免在 worker 尚未结束时提前删除。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>