【问题标题】:Implementation of multiprocessing.Queue and queue.Queuemultiprocessing.Queue和queue.Queue的实现
【发布时间】:2017-12-22 05:29:50
【问题描述】:

我正在寻找关于 Python 中队列实现的更多见解,而不是在文档中找到的。

根据我的理解,如果我在这方面错了,请原谅我的无知:

queue.Queue():通过内存中的基本数组实现,因此不能在多个进程之间共享,但可以在线程之间共享。到目前为止,一切顺利。

multiprocessing.Queue():通过具有大小限制的管道 (man 2 pipes) 实现(相当小:在 Linux 上,man 7 pipe 表示 65536 未调整):

从Linux 2.6.35开始,默认管道容量为65536字节,但可以使用fcntl(2)F_GETPIPE_SZF_SETPIPE_SZ操作查询和设置容量

但是,在 Python 中,每当我尝试将大于 65536 字节的数据写入管道时,它都会毫无例外地工作 - 我可能会以这种方式淹没我的内存:

import multiprocessing
from time import sleep

def big():
    result = ""
    for i in range(1,70000):
        result += ","+str(i)
    return result # 408888 bytes string

def writequeue(q):
    while True:
        q.put(big())
        sleep(0.1)

if __name__ == '__main__':
    q = multiprocessing.Queue()
    p = multiprocessing.Process(target=writequeue, args=(q,))
    p.start()
    while True:
        sleep(1) # No pipe consumption, we just want to flood the pipe

所以这是我的问题:

  • Python 是否调整了管道限制?如果是,多少?欢迎使用 Python 源代码。

  • Python 管道通信是否可以与其他非 Python 进程互操作?如果是,欢迎使用工作示例(最好是 JS)和资源链接。

【问题讨论】:

  • 您可以查看模块的代码 :) 了解事物如何工作的最佳方式,并查看好的 python 代码。
  • Louis 绝对是一个彻底的回答,关于最后一部分,如果您对 2 个程序之间的通信(无论它们的语言)感兴趣,您可能想看看实现 AMQP protocol 的经纪人(rabbitmqzeromq...等)。
  • 当您需要每秒执行数千次调用时,消息代理比直接文件描述符消耗慢很多......因此对python的数据格式互操作性感兴趣(我怀疑它“腌制”对象...)

标签: python linux queue pipe


【解决方案1】:

为什么 q.put() 没有阻塞??

mutiprocessing.Queue 创建一个管道,如果管道已满则阻塞。当然,写入超过管道容量会导致write 调用阻塞,直到读取端清除了足够的数据。好的,如果管道在达到其容量时阻塞,为什么q.put()在管道满时阻塞?即使是示例中对q.put() 的第一次调用也应该填满管道,并且所有内容都应该阻塞在那里,不是吗?

不,它不会阻塞,因为multiprocessing.Queue 实现将.put() 方法与对管道的写入分离。 .put() 方法将传递给它的数据排入内部缓冲区,并且有一个单独的线程负责从这个缓冲区读取并写入管道。当管道已满时,该线程将阻塞,但不会阻止.put() 将更多数据排入内部缓冲区。

.put() 的实现将数据保存到self._buffer 并注意如果没有线程正在运行,它是如何启动线程的:

def put(self, obj, block=True, timeout=None):
    assert not self._closed
    if not self._sem.acquire(block, timeout):
        raise Full

    with self._notempty:
        if self._thread is None:
            self._start_thread()
        self._buffer.append(obj)
        self._notempty.notify()

._feed() 方法是从self._buffer 读取数据并将数据提供给管道。而._start_thread() 是设置一个运行._feed() 的线程。

如何限制队列大小?

如果您想限制可以将多少数据写入队列,我看不到通过指定字节数来做到这一点的方法,但您可以限制存储在内部缓冲区中的项目数任何时候通过将数字传递给multiprocessing.Queue

q = multiprocessing.Queue(2)

当我使用上面的参数并使用您的代码时,q.put() 将排队两个项目,并在第三次尝试时阻塞。

Python 管道通信是否可以与其他非 Python 进程互操作?

这取决于。 multiprocessing 模块提供的功能不容易与其他语言互操作。我希望可能使multiprocessing 与其他语言互操作,但实现这一目标将是一项重大事业。编写该模块时期望所涉及的进程正在运行 Python 代码。

如果您查看更通用的方法,那么答案是肯定的。您可以使用套接字作为两个不同进程之间的通信管道。例如,一个从命名套接字读取的 JavaScript 进程:

var net = require("net");
var fs = require("fs");

sockPath = "/tmp/test.sock"
try {
    fs.unlinkSync(sockPath);
}
catch (ex) {
    // Don't care if the path does not exist, but rethrow if we get
    // another error.
    if (ex.code !== "ENOENT") {
        throw ex;
    }
}

var server = net.createServer(function(stream) {
  stream.on("data", function(c) {
    console.log("received:", c.toString());
  });

  stream.on("end", function() {
    server.close();
  });
});

server.listen(sockPath);

还有一个写入它的 Python 进程:

import socket
import time

sockfile = "/tmp/test.sock"

conn = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
conn.connect(sockfile)

count = 0
while True:
    count += 1
    conn.sendall(bytes(str(count), "utf-8"))
    time.sleep(1)

如果你想尝试上面的方法,你需要先启动 JavaScript 端,以便 Python 端有东西可以写入。这是一个概念验证。一个完整的解决方案需要更多的润色。

为了将复杂的结构从 Python 传递到其他语言,您必须找到一种方法,以一种双方都可以读取的格式序列化您的数据。不幸的是,泡菜是 Python 特有的。每当我需要在语言之间进行序列化时,我通常会选择 JSON,或者如果 JSON 不会这样做,则使用 ad-hoc 格式。

【讨论】:

  • 您的回答很有见地。我赞成这个原因。问题的第二部分是关于与基于其他语言的其他进程互操作队列。例如:使用类似兼容的put()get() 操作将字符串从/到Python 传递到Javascript 进程。或者 Python 是否使用了一种隐蔽的格式来传递对象?不要犹豫,更新您对另一部分的答案。
  • 对不起。我确实错过了那一点。我已经编辑了我的答案来解决这个问题。
  • 这听起来不错,谢谢。如果有人愿意分享他自己的观点作为答案,我会花一些时间并稍后批准答案。
  • 我错过了赏金截止日期,否则我会验证的。很抱歉给您带来不便,我正忙着 IRL。你的回答非常好,所以我相信它比几个代表点更有价值:-)
  • 这会将队列限制为两个项目,但仍然不限制队列的大小。
猜你喜欢
  • 1970-01-01
  • 2019-12-04
  • 1970-01-01
  • 2018-04-13
  • 2011-03-08
  • 1970-01-01
  • 1970-01-01
  • 2011-07-02
  • 1970-01-01
相关资源
最近更新 更多