为何第二次调用queue.get()时会陷入无限循环?
问题
尝试同时运行两个函数并获取它们的输出,编写了如下代码,但仅能成功获取q1或q2的输出一次,第二次调用q1.get()时会陷入无限等待。
from threading import Thread from multiprocessing import Queue import multiprocessing import time def functionA(A, B): dictA=[] for i in list(range(A, B)): print(i, "from A") time.sleep(1) for p in list(range(0, 10)): dictA.append({i:p}) return(dictA) def functionB(C, D): dictB=[] for i in list(range(C, D)): print(i, "from B") time.sleep(1) for p in list(range(0, 10)): dictB.append({i:p}) return(dictB) # 包装函数,将函数返回值放入队列 def wrapper(func, args, queue): queue.put(func(*args)) q1, q2 = multiprocessing.Manager().Queue(), multiprocessing.Manager().Queue() Thread(target=wrapper, args=(functionA, [0, 10], q1)).start() Thread(target=wrapper, args=(functionB, [10, 20], q2)).start()
现象
第一次调用q1.get()可正常获取结果:
print(q1.get()) # 耗时约0.9s # 输出: [{0: 0}, {0: 1}, {0: 2}, {0: 3}, {0: 4}, {0: 5}, {0: 6}, {0: 7}, {0: 8}, {0: 9}, {1: 0}, {1: 1}, {1: 2}, {1: 3}, {1: 4}, {1: 5}, {1: 6}, {1: 7}, {1: 8}, {1: 9}, {2: 0}, {2: 1}, {2: 2}, {2: 3}, {2: 4}, {2: 5}, {2: 6}, {2: 7}, {2: 8}, {2: 9}, {3: 0}, {3: 1}, {3: 2}, {3: 3}, {3: 4}, {3: 5}, {3: 6}, {3: 7}, {3: 8}, {3: 9}, {4: 0}, {4: 1}, {4: 2}, {4: 3}, {4: 4}, {4: 5}, {4: 6}, {4: 7}, {4: 8}, {4: 9}, {5: 0}, {5: 1}, {5: 2}, {5: 3}, {5: 4}, {5: 5}, {5: 6}, {5: 7}, {5: 8}, {5: 9}, {6: 0}, {6: 1}, {6: 2}, {6: 3}, {6: 4}, {6: 5}, {6: 6}, {6: 7}, {6: 8}, {6: 9}, {7: 0}, {7: 1}, {7: 2}, {7: 3}, {7: 4}, {7: 5}, {7: 6}, {7: 7}, {7: 8}, {7: 9}, {8: 0}, {8: 1}, {8: 2}, {8: 3}, {8: 4}, {8: 5}, {8: 6}, {8: 7}, {8: 8}, {8: 9}, {9: 0}, {9: 1}, {9: 2}, {9: 3}, {9: 4}, {9: 5}, {9: 6}, {9: 7}, {9: 8}, {9: 9}]
第二次调用q1.get()则会无限卡住:
print(q1.get()) # 程序无输出,持续阻塞
原因分析
multiprocessing.Manager().Queue()是进程间通信队列,队列中的每个元素只能被get()一次。第一次调用q1.get()时,已经把队列里唯一的元素(functionA的返回值)取走,队列变为空状态。
第二次调用q1.get()时,因为队列没有新元素,get()方法默认会阻塞等待新数据,但你的线程已经执行完毕(wrapper只往队列里放了一次数据),没有后续元素加入,所以程序会一直等待下去,表现为“无限循环”。
解决方案
根据需求不同,有两种处理方式:
1. 仅获取一次结果(匹配原代码逻辑)
如果只需要获取一次函数返回值,确保每个队列只调用一次get()即可:
# 分别获取两个函数的结果 resultA = q1.get() resultB = q2.get() print(resultA) print(resultB)
2. 多次获取结果(若需要函数持续输出)
如果需要函数持续产生数据并多次获取,需修改业务函数和包装逻辑,让函数循环生成数据并放入队列:
def functionA(A, B, queue): # 改为循环产生多轮数据 for round in range(3): dictA = [] for i in list(range(A, B)): print(i, "from A round", round) time.sleep(0.1) for p in list(range(0, 2)): dictA.append({i:p}) queue.put(dictA) def wrapper(func, args, queue): func(*args, queue) q1 = multiprocessing.Manager().Queue() Thread(target=wrapper, args=(functionA, [0, 3], q1)).start() # 可以三次获取结果 for _ in range(3): print(q1.get())
内容的提问来源于stack exchange,提问作者Rivered
相关产品推荐
相关产品推荐

