我正在嘗試實作一個佇列。這是舊代碼,要么取自我前一段時間所做的某種教程,要么取自我閱讀檔案時所做的某種實驗,或者兩者的混合。問題是我不確定代碼是否是我的,但我試圖用它作為一個例子來學習。該腳本有一個生產者在串列中生成數字,兩個消費者競爭獲取這些數字并將它們相加,總和最高的人獲勝。
所以,這是我的問題:在“consume_numbers”函式的以下代碼中,我有一個 time.sleep(0.01) 行,它使代碼運行。沒有它,代碼就會掛起,但它運行起來很順暢。有人可以解釋為什么會發生這種情況以及如何在沒有這個問題的情況下實作佇列嗎?
import concurrent.futures
import time
import random
import threading
import queue
class MyQueue(queue.Queue):
def __init__(self, maxsize=10):
super().__init__()
self.maxsize = maxsize
self.numbers = []
def set_number(self, number):
self.put(number)
self.numbers.append(number)
def get_number(self):
return self.get()
def produce_random_numbers(q: MyQueue, maxcount: int, evnt: threading.Event):
count = 0
while not evnt.is_set():
num = random.randint(1, 5)
q.set_number(num)
count = 1
if count > maxcount:
event.set()
def consume_numbers(q: MyQueue, consumed: list, evnt: threading.Event):
while not q.empty() or not evnt.is_set():
num = q.get_number()
time.sleep(0.01)
consumed.append(num)
if __name__ == "__main__":
q = MyQueue(maxsize=10)
event = threading.Event()
cons1 = []
cons2 = []
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as ex:
ex.submit(produce_random_numbers, q, 50, event)
ex.submit(consume_numbers, q, cons1, event)
ex.submit(consume_numbers, q, cons2, event)
event.set()
print(f'Generated Numbers: {q.numbers}')
print(f'Numbers Consumed by Thread1 which summed up to {sum(cons1)} are: {cons1}')
print(f'Numbers Consumed by Thread2 which summed up to {sum(cons2)} are: {cons2}')
if sum(cons1) > sum(cons2):
print("Thread1 Wins!")
elif sum(cons1) < sum(cons2):
print("Thread2 Wins!")
else:
print("It's a tie!")
謝謝!
uj5u.com熱心網友回復:
該代碼沒有從頭開始實作佇列,而是擴展queue.Queue以添加記憶體。有一個事件物件用于向消費者發出生產者執行緒已完成的信號。當佇列中只有一項時,消費者中存在隱藏的競爭條件。
not q.empty() or not evnt.is_set()如果佇列中有東西或尚未設定事件,檢查將運行回圈代碼。可能會發生這樣的情況:
- 一個執行緒看到佇列不為空,進入回圈
- 發生執行緒切換,另一個執行緒消耗最后一項
- 切換發生在第一個執行緒上,該執行緒呼叫
get_number()并阻塞
檢查時會發生類似的競爭條件evnt.is_set():
- 生產者將最后一項添加到佇列中,并發生執行緒切換
- 一個執行緒消耗最后一項,一個開關
- 發生執行緒切換,消費者獲取最后一項并回傳回圈條件。由于尚未設定事件,因此執行回圈并
get_number()阻塞
讓執行緒等待可以最大限度地減少這些情況發生的機會。沒有等待,很可能單個消費者執行緒將消耗所有佇列專案,而另一個仍在進入其回圈。
使用超時很麻煩。避免使用事件的一個有用的習慣用法是使用iter和使用一個不可能的值作為哨兵:
# --- snip ---
def produce_random_numbers(q: MyQueue, maxcount: int, n_consumers: int):
for _ in range(maxcount):
num = random.randint(1, 5)
q.set_number(num)
for _ in range(n_consumers):
q.put(None) # <--- I use put to put one sentinel per consumer
def consume_numbers(q: MyQueue, consumed: list):
for num in iter(q.get_number, None):
consumed.append(num)
if __name__ == "__main__":
q = MyQueue(maxsize=10)
cons1 = []
cons2 = []
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as ex:
ex.submit(produce_random_numbers, q, 500000, 2)
ex.submit(consume_numbers, q, cons1)
ex.submit(consume_numbers, q, cons2)
print(f'Generated Numbers: {q.numbers}')
# --- snip ---
還有一些其他問題和我會做不同的事情:
event.set()后面的塊with...沒用:事件已經被生產者設定了- 生產者中有一個錯字,使用全域
event變數而不是區域evnt引數。幸運的是,它們指的是同一個物件。 - 因為只有一個生產者,所以不會有問題。否則 MyQueue.numbers 的順序可能與將專案添加到佇列中的順序不同:
put在一個執行緒上呼叫- 發生執行緒切換
- a
putappend發生在新執行緒中 - 發生執行緒切換,第一個值是
appended
- 而不是定義
MyQueue.set_number我會覆寫put
轉載請註明出處,本文鏈接:https://www.uj5u.com/qukuanlian/417311.html
標籤:
