【问题标题】:Python 3.4 multiprocessing Queue faster than Pipe, unexpectedPython 3.4 多处理队列比 Pipe 快,出乎意料
【发布时间】:2015-01-10 14:24:04
【问题描述】:

我正在做一个从 udp 套接字接收样本的音频播放器,一切正常。但是当我实现了一个 Lost Concealment 算法时,播放器未能以例外的速率保持沉默(每 10 毫秒发送一个包含多个 160 字节的列表)。

使用 pyaudio 播放音频时,使用阻塞调用 write 播放一些样本,我注意到它在样本持续时间内平均阻塞。所以我创建了一个新的专用流程来播放样本。

主进程处理音频的输出流,并使用 multiprocessing.Pipe 将结果发送到该进程。我决定使用 multiprocessing.Pipe 因为它应该比其他方式更快。

不幸的是,当我在虚拟机上运行程序时,比特率是我在快速 PC 上获得的一半,这并没有达到目标比特率。

经过一些测试,我得出结论,导致延迟的原因是 Pipe 的函数send

我做了一个简单的基准测试脚本(见下文)来查看传输到进程的各种方法之间的差异。该脚本不断发送[b'\x00'*160] 5 秒,并计算总共发送了多少字节对象。我测试了以下发送方法:“不发送”、multiprocessing.Pipe、multiprocessing.Queue、multiprocessing.Manager、multiprocessing.Listener/Client,最后是socket.socket:

我的“快速”PC 运行窗口 7 x64 的结果:

test_empty     :     1516076640
test_pipe      :       58155840
test_queue     :      233946880
test_manager   :        2853440
test_socket    :       55696160
test_named_pipe:       58363040

VirtualBox 的 VM 来宾运行 Windows 7 x64,主机运行 Windows 7 x64 的结果:

test_empty     :     1462706080
test_pipe      :       32444160
test_queue     :      204845600
test_manager   :         882560
test_socket    :       20549280
test_named_pipe:       35387840  

使用的脚本:

from multiprocessing import Process, Pipe, Queue, Manager
from multiprocessing.connection import Client, Listener
import time

FS = "{:<15}:{:>15}"


def test_empty():
    s = time.time()
    sent = 0
    while True:
        data = b'\x00'*160
        lst = [data]

        sent += len(data)
        if time.time()-s >= 5:
            break
    print(FS.format("test_empty", sent))


def pipe_void(pipe_in):
    while True:
        msg = pipe_in.recv()
        if msg == []:
            break


def test_pipe():
    pipe_out, pipe_in = Pipe()
    p = Process(target=pipe_void, args=(pipe_in,))
    p.start()
    s = time.time()
    sent = 0
    while True:
        data = b'\x00'*160
        lst = [data]
        pipe_out.send(lst)
        sent += len(data)
        if time.time()-s >= 5:
            break
    pipe_out.send([])
    p.join()
    print(FS.format("test_pipe", sent))


def queue_void(q):
    while True:
        msg = q.get()
        if msg == []:
            break


def test_queue():
    q = Queue()
    p = Process(target=queue_void, args=(q,))
    p.start()
    s = time.time()
    sent = 0
    while True:
        data = b'\x00'*160
        lst = [data]
        q.put(lst)
        sent += len(data)
        if time.time()-s >= 5:
            break
    q.put([])
    p.join()

    print(FS.format("test_queue", sent))


def manager_void(l, lock):
    msg = None
    while True:
        with lock:
            if len(l) > 0:
                msg = l.pop(0)
        if msg == []:
            break


def test_manager():
    with Manager() as manager:
        l = manager.list()
        lock = manager.Lock()
        p = Process(target=manager_void, args=(l, lock))
        p.start()
        s = time.time()
        sent = 0
        while True:
            data = b'\x00'*160
            lst = [data]
            with lock:
                l.append(lst)
            sent += len(data)
            if time.time()-s >= 5:
                break
        with lock:
            l.append([])
        p.join()

        print(FS.format("test_manager", sent))


def socket_void():
    addr = ('127.0.0.1', 20000)
    conn = Client(addr)
    while True:
        msg = conn.recv()
        if msg == []:
            break


def test_socket():
    addr = ('127.0.0.1', 20000)
    listener = Listener(addr, "AF_INET")
    p = Process(target=socket_void)
    p.start()
    conn = listener.accept()
    s = time.time()
    sent = 0
    while True:
        data = b'\x00'*160
        lst = [data]
        conn.send(lst)
        sent += len(data)
        if time.time()-s >= 5:
            break
    conn.send([])
    p.join()

    print(FS.format("test_socket", sent))


def named_pipe_void():
    addr = '\\\\.\\pipe\\Test'
    conn = Client(addr)
    while True:
        msg = conn.recv()
        if msg == []:
            break


def test_named_pipe():
    addr = '\\\\.\\pipe\\Test'
    listener = Listener(addr, "AF_PIPE")
    p = Process(target=named_pipe_void)
    p.start()
    conn = listener.accept()
    s = time.time()
    sent = 0
    while True:
        data = b'\x00'*160
        lst = [data]
        conn.send(lst)
        sent += len(data)
        if time.time()-s >= 5:
            break
    conn.send([])
    p.join()

    print(FS.format("test_named_pipe", sent))


if __name__ == "__main__":
    test_empty()
    test_pipe()
    test_queue()
    test_manager()
    test_socket()
    test_named_pipe()

问题

  • 如果 Queue 使用 Pipe 在这种情况下它比 Pipe 快吗? 这与Python multiprocessing - Pipe vs Queue 的问题相矛盾
  • 如何保证从一个进程到另一个进程的恒定比特率流,同时具有低发送延迟?

更新 1

在我的程序中,在尝试使用队列而不是管道之后。 我得到了巨大的提升

在我的计算机上,使用 Pipes 我得到 +- 16000 B/s,使用 Queues 我得到 +-750 万 B/s。在虚拟机上,我的速度从 +-13000 B/s 提高到了 650 万 B/s。使用 Queue instread 的 Pipe 字节数增加了大约 500 倍。

当然我不会每秒播放数百万字节,我只会播放正常的声音速率。 (在我的情况下为 16000 B/s,与上述值一致)。
但关键是,我可以将速率限制在我想要的范围内,同时仍有时间完成其他计算(如从套接字接收、应用声音算法等)

【问题讨论】:

  • 您的 Queue 测试无效,因为 5 秒用于填充本地队列 (q._buffer),而工作线程提取并发送对象。然后你 join 等待队列清空,这需要超过 5 秒的时间。相反,您应该计算发送给定字节数所需的时间。
  • 另外,test_pipe 和 test_named_pipe 本质上是同一个测试。您只是在后者中使用了明确的名称,而不是 Pipe 使用的 connection.arbitrary_address('AF_PIPE')
  • @eryksun 感谢您关注细节。我将阅读队列实现并进行不同的测试
  • 也许 Python3 只是简单地改进了multiprocessing 模块,所指的比较不再相关?也可能是特定于操作系统的原因,Pipe 实际上使用管道,那么 Windows 实现的性能将与 Linux 或 OSX 完全不同。同时Queue 必须使用共享内存(我认为),因此跨操作系统的性能将是相似的。
  • 我正在对 circuits 进行基准测试,因为这个问题激起了我对异步 I/O 可以 实现什么的兴趣;请参阅:gist.github.com/prologic/fea809047ba61847ebb5——我以 10 毫秒的间隔使用 circuits.node 获得 ~577KB/s。

标签: python windows sockets python-3.x python-internals


【解决方案1】:

我不能肯定地说,但我认为您要处理的问题是同步 I/O 与异步 I/O。我的猜测是 Pipe 以某种方式以同步方式结束,而 Queue 以异步方式结束。为什么一个是默认一种方式而另一种是另一种方式可能会更好地通过这个问题和答案来回答:

Synchronous/Asynchronous behaviour of python Pipes

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-01-12
    • 2013-06-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-02-10
    • 1970-01-01
    相关资源
    最近更新 更多