【发布时间】:2018-04-15 13:50:44
【问题描述】:
我需要帮助来理解multiprocessing.Queue。我面临的问题是,与对queue.put(...) 和队列缓冲区(双端队列)的调用相比,从queue.get(...) 获取结果非常落后。
这种泄漏的抽象使我研究了队列的内部结构。它直截了当的source code 只是将我指向deque implementation,这似乎也很简单,以至于我无法用它来解释我所看到的行为。我还读到 Queue 使用管道,但我似乎无法在源代码中找到它。
我已将其归结为重现问题的最小示例,并在其下方指定可能的输出。
import threading
import multiprocessing
import queue
q = None
def enqueue(item):
global q
if q is None:
q = multiprocessing.Queue()
process = threading.Thread(target=worker, args=(q,)) # or multiprocessing.Process Doesn't matter
process.start()
q.put(item)
print(f'len putted item: {len(item)}. qsize: {q.qsize()}. buffer len: {len(q._buffer)}')
def worker(local_queue):
while True:
try:
while True: # get all items
item = local_queue.get(block=False)
print(f'len got item: {len(item)}. qsize: {q.qsize()}. buffer len: {len(q._buffer)}')
except queue.Empty:
print('empty')
if __name__ == '__main__':
for i in range(1, 100000, 1000):
enqueue(list(range(i)))
输出:
empty
empty
empty
len putted item: 1. qsize: 1. buffer len: 1
len putted item: 1001. qsize: 2. buffer len: 2
len putted item: 2001. qsize: 3. buffer len: 1
len putted item: 3001. qsize: 4. buffer len: 2
len putted item: 4001. qsize: 5. buffer len: 3
len putted item: 5001. qsize: 6. buffer len: 4
len putted item: 6001. qsize: 7. buffer len: 5
len putted item: 7001. qsize: 8. buffer len: 6
len putted item: 8001. qsize: 9. buffer len: 7
len putted item: 9001. qsize: 10. buffer len: 8
len putted item: 10001. qsize: 11. buffer len: 9
len putted item: 11001. qsize: 12. buffer len: 10
len putted item: 12001. qsize: 13. buffer len: 11
len putted item: 13001. qsize: 14. buffer len: 12
len putted item: 14001. qsize: 15. buffer len: 13
len putted item: 15001. qsize: 16. buffer len: 14
len got item: 1. qsize: 15. buffer len: 14
len putted item: 16001. qsize: 16. buffer len: 15
len putted item: 17001. qsize: 17. buffer len: 16
len putted item: 18001. qsize: 18. buffer len: 17
len putted item: 19001. qsize: 19. buffer len: 18
len putted item: 20001. qsize: 20. buffer len: 19
len putted item: 21001. qsize: 21. buffer len: 20
len putted item: 22001. qsize: 22. buffer len: 21
len putted item: 23001. qsize: 23. buffer len: 22
len putted item: 24001. qsize: 24. buffer len: 23
len putted item: 25001. qsize: 25. buffer len: 24
len putted item: 26001. qsize: 26. buffer len: 25
len putted item: 27001. qsize: 27. buffer len: 26
len putted item: 28001. qsize: 28. buffer len: 27
len got item: 1001. qsize: 27. buffer len: 27
empty
len putted item: 29001. qsize: 28. buffer len: 28
empty
empty
empty
len got item: 2001. qsize: 27. buffer len: 27
empty
len putted item: 30001. qsize: 28. buffer len: 28
我希望您注意以下结果:插入元素 28001 后,worker 发现队列中没有剩余元素,而还有几十个元素。由于同步,我可以只获取除少数之外的所有内容。但它只能找到两个!
这种模式还在继续。
这似乎与我放入队列的对象的大小有关。对于小对象,比如i 而不是list(range(i)),则不会出现此问题。但是所讨论的对象的大小仍然是千字节,不足以让如此严重的延迟有尊严(在我的真实世界非最小示例中,这很容易花费几分钟)
我的具体问题是:如何在 Python 中的进程之间共享(不是这样)大量数据? 另外,我想知道在 Queue 的内部实现中,这种迟缓是从哪里来的
【问题讨论】:
-
另外我是 Python 新手,欢迎评论
-
你找到解决办法了吗
标签: python-3.x message-queue python-multiprocessing