使用Pipe在进程之间传输Python对象时的字节限制?

Col*_*n M 8 python parallel-processing struct multiprocessing

我有一个使用64位Python 3.3.0 CPython解释器在64位Linux(内核版本2.6.28.4)机器上运行的自定义模拟器(用于生物学).

因为模拟器依赖于许多独立实验来获得有效结果,所以我建立了并行处理来运行实验.线程之间的通信主要发生在具有托管multiprocessing Queues(doc)的生产者 - 消费者模式下 .该体系结构的破坏如下:

  • 处理产卵和管理Processes以及各种Queues 的主进程
  • N个工作进程进行模拟
  • 1结果消费者过程消耗模拟结果并对结果进行分类和分析

主进程和工作进程通过输入进行通信Queue.类似地,工作进程将结果放在Queue结果消费者进程从中消耗项目的输出中.最终的ResultConsumer对象通过multiprocessing Pipe (doc)传递回主进程.

一切正常,直到它试图通过以下方法将ResultConsumer对象传递回主进程Pipe:

Traceback (most recent call last):
  File "/home/cmccorma/.local/lib/python3.3/multiprocessing/process.py", line 258, in _bootstrap
    self.run()
  File "/home/cmccorma/.local/lib/python3.3/multiprocessing/process.py", line 95, in run
    self._target(*self._args, **self._kwargs)
  File "DomainArchitectureGenerator.py", line 93, in ResultsConsumerHandler
    pipeConn.send(resCon)
  File "/home/cmccorma/.local/lib/python3.3/multiprocessing/connection.py", line 207, in send
    self._send_bytes(buf.getbuffer())
  File "/home/cmccorma/.local/lib/python3.3/multiprocessing/connection.py", line 394, in _send_bytes
    self._send(struct.pack("!i", n))
struct.error: 'i' format requires -2147483648 <= number <= 2147483647
Run Code Online (Sandbox Code Playgroud)

我理解前两个跟踪(Process库中未处理的出口),第三个是我的代码行,用于将ResultConsumer对象发送 Pipe到主进程.最后两条痕迹是它变得有趣的地方.一个Pipepickle发送给它的任何对象,并将生成的字节传递给另一端(匹配连接),在运行时它被取消激活recv(). self._send_bytes(buf.getbuffer())正在尝试发送pickle对象的字节. self._send(struct.pack("!i", n))正在尝试打包一个长度为n的整数(network/big-endian)的结构,其中n是作为参数传入的缓冲区的长度(该struct 库处理Python值和表示为Python字符串的C结构之间的转换,请参阅文件).

只有在尝试大量实验时才会出现此错误,例如,10个实验不会导致它,但1000个将是有意义的(所有其他参数都是恒定的).到目前为止,我最好的假设struct.error是,尝试按下管道的字节数超过2 ^ 32-1(2147483647),或~2 GB.

所以我的问题是双重的:

  1. 我对我的调查感到困惑,因为struct.py基本上只是从中进口_struct而且我不知道它在哪里.

  2. 鉴于底层架构都是64位,字节限制似乎是任意的.那么,为什么我不能通过比这更大的东西?另外,如果我无法改变这个问题,那么这个问题是否有任何好的(阅读:简单)解决方法?

注意:我不认为用a Queue代替a Pipe会解决问题,因为我怀疑它Queue是使用类似的酸洗中间步骤. 编辑:正如abarnert的回答所指出的,这个说明是完全错误的.

aba*_*ert 10

我因为struct.py基本上只是从_struct导入而陷入我的调查,我不知道它在哪里.

在CPython中,_struct是一个C扩展模块,它是在源代码树_struct.cModules目录中构建的.您可以在这里找到在线代码.

每当foo.py一个import _foo,几乎总是一个C扩展模块,通常是由_foo.c.如果你根本找不到foo.py它,它可能是一个C扩展模块,由_foomodule.c.

即使你没有使用PyPy,它也经常值得查看等效的PyPy源代码.他们重新实现纯Python中的几乎所有扩展模块 - 对于其余的(包括本例),底层的"扩展语言"是RPython,而不是C.

但是,在这种情况下,您不需要知道如何struct工作超出文档中的内容.


鉴于底层架构都是64位,字节限制似乎是任意的.

看看它调用的代码:

self._send(struct.pack("!i", n))
Run Code Online (Sandbox Code Playgroud)

如果查看文档,'i'格式字符显式表示"4字节C整数",而不是"无论ssize_t是什么".为此,你必须使用'n'.或者您可能希望明确使用long long,with 'q'.

你可以multiprocessing使用monkeypatch struct.pack('!q', n).或者'!q'.或者以某种方式编码长度struct.当然,这将破坏与非补丁的兼容性,multiprocessing如果您尝试跨多台计算机或其他东西进行分布式处理,这可能是一个问题.但它应该很简单:

def _send_bytes(self, buf):
    # For wire compatibility with 3.2 and lower
    n = len(buf)
    self._send(struct.pack("!q", n)) # was !i
    # The condition is necessary to avoid "broken pipe" errors
    # when sending a 0-length buffer if the other end closed the pipe.
    if n > 0:
        self._send(buf)

def _recv_bytes(self, maxsize=None):
    buf = self._recv(8) # was 4
    size, = struct.unpack("!q", buf.getvalue()) # was !i
    if maxsize is not None and size > maxsize:
        return None
    return self._recv(size)
Run Code Online (Sandbox Code Playgroud)

当然,不能保证这种变化是充分的; 你会想要阅读周围代码的其余部分并测试它的地狱.


注意:我怀疑使用a Queue代替a Pipe不会解决问题,因为我怀疑它Queue使用类似的酸洗中间步骤.

嗯,问题与酸洗无关.Pipe不使用pickle发送长度,它正在使用struct.您可以验证pickle不会出现此问题:pickle.loads(pickle.dumps(1<<100)) == 1<<100将返回True.

(在早期版本中,pickle 存在大型物体的问题 - 例如,一个list2G元素 - 这可能导致问题的规模大约是目前正在击中的那么高的8倍.但是已经修正了3.3.)

与此同时......尝试看看它并不是更快,而不是通过挖掘源来试图弄清楚它是否会起作用?


另外,你确定你真的想通过隐式酸洗来传递2GB的数据结构吗?

如果我做了一些缓慢而且需要内存的东西,我宁愿将其显式化 - 例如,pickle到tempfile并发送路径或fd.(如果您正在使用numpy或者pandas使用其他内容,请使用其二进制文件格式,而不是pickle相同的想法.)

或者,更好的是,共享数据.是的,可变共享状态很糟糕......但共享不可变对象很好.无论你有2GB,你可以把它放在一个multiprocessing.Array,或者把它放在一个ctypes阵列或结构(阵列或结构......)中你可以分享multiprocessing.sharedctypes,或者ctypes它可以从file你那里分享mmap,或...... ?有一些额外的代码来定义和分离结构,但是当这些好处可能很大时,值得尝试.


最后,当您认为在Python中发现了一个错误/明显缺失的功能/不合理的限制时,值得查看错误跟踪器.看起来像问题17560:使用多处理真正大对象的问题?这正是你的问题,并有很多信息,包括建议的解决方法.