Python中asyncio.Queue異步隊列的實現(xiàn)
一、概述
asyncio.Queue 是 Python 標(biāo)準(zhǔn)庫 asyncio 模塊中的異步隊列實現(xiàn),位于 Lib/asyncio/queues.py。它的設(shè)計思路與線程安全的 queue.Queue 類似,但專門為 async/await 協(xié)程模型 設(shè)計,不是線程安全的。
核心定位:在異步并發(fā)環(huán)境下,作為生產(chǎn)者-消費者模式的協(xié)調(diào)橋梁,實現(xiàn)協(xié)程之間的數(shù)據(jù)傳遞與流量控制。
import asyncio queue = asyncio.Queue(maxsize=10) # 有界隊列,容量10 queue = asyncio.Queue() # 無界隊列(maxsize=0)
二、內(nèi)部實現(xiàn)原理
從 CPython 源碼來看,asyncio.Queue 的核心數(shù)據(jù)結(jié)構(gòu)非常精巧:
class Queue(mixins._LoopBoundMixin):
def __init__(self, maxsize=0):
self._maxsize = maxsize
self._getters = collections.deque() # 等待獲取的 Future 隊列
self._putters = collections.deque() # 等待放入的 Future 隊列
self._unfinished_tasks = 0 # 未完成任務(wù)計數(shù)
self._finished = locks.Event() # join() 阻塞用的事件
self._finished.set()
self._init(maxsize)
self._is_shutdown = False
關(guān)鍵設(shè)計
| 內(nèi)部屬性 | 作用 |
|---|---|
| _queue | 實際存儲數(shù)據(jù)的 collections.deque(FIFO) |
| _getters | 當(dāng)隊列為空時,get() 的調(diào)用者會創(chuàng)建一個 Future 并掛入此 deque 等待 |
| _putters | 當(dāng)隊列已滿時,put() 的調(diào)用者會創(chuàng)建一個 Future 并掛入此 deque 等待 |
| _unfinished_tasks | 配合 task_done() / join() 實現(xiàn)任務(wù)追蹤 |
| _finished | 一個 asyncio.Event,當(dāng)未完成任務(wù)歸零時被 set,喚醒 join() |
協(xié)程掛起與喚醒機制
put() 和 get() 的阻塞并非真正的線程阻塞,而是通過 Future + await 實現(xiàn)的協(xié)程掛起:
# put() 核心邏輯(簡化)
async def put(self, item):
while self.full():
putter = self._get_loop().create_future()
self._putters.append(putter)
await putter # 協(xié)程在此掛起
self.put_nowait(item)
# get() 核心邏輯(簡化)
async def get(self):
while self.empty():
getter = self._get_loop().create_future()
self._getters.append(getter)
await getter # 協(xié)程在此掛起
return self.get_nowait()
當(dāng)對面操作發(fā)生時,通過 _wakeup_next() 喚醒等待者:
def _wakeup_next(self, waiters):
while waiters:
waiter = waiters.popleft()
if not waiter.done():
waiter.set_result(None) # 喚醒一個等待的協(xié)程
break
這種 逐個喚醒(one-by-one wakeup) 設(shè)計避免了"驚群效應(yīng)",保證公平性:先等待的協(xié)程先被喚醒。
三、完整 API 詳解
3.1 構(gòu)造函數(shù)
asyncio.Queue(maxsize=0)
maxsize <= 0:無界隊列,put()永遠(yuǎn)不會阻塞(內(nèi)存允許的前提下)maxsize > 0:有界隊列,隊列滿時put()會掛起等待
3.2 核心方法
async put(item)
將元素放入隊列。如果隊列已滿,掛起當(dāng)前協(xié)程直到有空間。
put_nowait(item)
非阻塞放入。隊列滿時立即拋出 QueueFull 異常。
async get()
從隊列取出并返回一個元素。如果隊列為空,掛起當(dāng)前協(xié)程直到有元素可用。
get_nowait()
非阻塞獲取。隊列空時立即拋出 QueueEmpty 異常。
3.3 任務(wù)追蹤方法
task_done()
標(biāo)記一個之前通過 get() 獲取的任務(wù)已完成。內(nèi)部將 _unfinished_tasks 減 1。當(dāng)計數(shù)歸零時,觸發(fā) _finished 事件,解除 join() 的阻塞。
注意:調(diào)用次數(shù)超過 put() 的次數(shù)會拋出 ValueError。
async join()
阻塞直到隊列中所有任務(wù)都被處理完畢(每個 put 進去的任務(wù)都收到了對應(yīng)的 task_done() 調(diào)用)。
# 典型用法 await queue.join() # 等待所有任務(wù)完成
3.4 狀態(tài)查詢
| 方法 | 說明 |
|---|---|
| qsize() | 返回隊列當(dāng)前元素數(shù)量 |
| empty() | 隊列是否為空 |
| full() | 隊列是否已滿(maxsize=0 時永遠(yuǎn)返回 False) |
| maxsize | 屬性,返回隊列容量上限 |
與 threading.Queue 不同,由于單線程事件循環(huán)的特性,qsize() 的返回值在 await 之前是可靠的——不會被其他線程中斷。
3.5 關(guān)閉方法(Python 3.13+)
shutdown(immediate=False)
queue.shutdown() # 優(yōu)雅關(guān)閉:允許消費完剩余元素 queue.shutdown(immediate=True) # 立即關(guān)閉:清空隊列,中斷所有等待
- 優(yōu)雅關(guān)閉(immediate=False):
- put() 立即拋出 QueueShutDown
- get() 可以繼續(xù)取出已有元素,隊列空后才拋出 QueueShutDown
- 立即關(guān)閉(immediate=True):
- 隊列被清空,_unfinished_tasks 相應(yīng)減少
- 所有阻塞的 get() 和 put() 立即被喚醒并拋出 QueueShutDown
- 如果 _unfinished_tasks 歸零,join() 也會被解除阻塞
四、三種隊列變體
asyncio 提供了三種隊列,它們通過重寫 _init、_get、_put 三個內(nèi)部方法實現(xiàn)不同的出隊策略:
4.1Queue(FIFO 先進先出)
def _init(self, maxsize):
self._queue = collections.deque()
def _put(self, item):
self._queue.append(item)
def _get(self):
return self._queue.popleft()
4.2PriorityQueue(優(yōu)先級隊列)
def _init(self, maxsize):
self._queue = []
def _put(self, item, heappush=heapq.heappush):
heappush(self._queue, item)
def _get(self, heappop=heapq.heappop):
return heappop(self._queue)
元素通常為 (priority, data) 元組,數(shù)值越小優(yōu)先級越高。
4.3LifoQueue(后進先出 / 棧)
def _init(self, maxsize):
self._queue = []
def _put(self, item):
self._queue.append(item)
def _get(self):
return self._queue.pop() # 從尾部取出
五、異常體系
| 異常 | 觸發(fā)時機 |
|---|---|
| asyncio.QueueEmpty | get_nowait() 在空隊列上調(diào)用 |
| asyncio.QueueFull | put_nowait() 在滿隊列上調(diào)用 |
| asyncio.QueueShutDown | 在已 shutdown() 的隊列上調(diào)用 put() 或 get()(3.13+) |
六、經(jīng)典模式:生產(chǎn)者-消費者
import asyncio
import random
async def producer(queue: asyncio.Queue, name: str):
for i in range(5):
item = f"{name}-item-{i}"
await queue.put(item)
print(f"[{name}] 生產(chǎn): {item}")
await asyncio.sleep(random.uniform(0.1, 0.5))
async def consumer(queue: asyncio.Queue, name: str):
while True:
item = await queue.get()
print(f" [{name}] 消費: {item}")
await asyncio.sleep(random.uniform(0.2, 0.8))
queue.task_done()
async def main():
queue = asyncio.Queue(maxsize=3) # 有界隊列,產(chǎn)生背壓
# 啟動生產(chǎn)者和消費者
producers = [asyncio.create_task(producer(queue, f"P{i}")) for i in range(2)]
consumers = [asyncio.create_task(consumer(queue, f"C{i}")) for i in range(3)]
# 等待所有生產(chǎn)完成
await asyncio.gather(*producers)
# 等待隊列中所有任務(wù)被消費完
await queue.join()
# 取消消費者(它們在 await queue.get() 處無限等待)
for c in consumers:
c.cancel()
asyncio.run(main())
關(guān)鍵要點
- maxsize 實現(xiàn)背壓(backpressure):當(dāng)消費者處理速度跟不上時,生產(chǎn)者自動被掛起,防止內(nèi)存無限增長
- task_done() + join() 配對使用:確保每個任務(wù)都被處理完畢后再收尾
- 消費者用 while True + cancel() 模式:消費者持續(xù)監(jiān)聽,生產(chǎn)結(jié)束后由外部取消
七、超時控制
asyncio.Queue 的方法本身不接受 timeout 參數(shù),需要配合 asyncio.wait_for() 使用:
try:
item = await asyncio.wait_for(queue.get(), timeout=5.0)
except asyncio.TimeoutError:
print("等待超時,隊列中沒有新數(shù)據(jù)")
try:
await asyncio.wait_for(queue.put(item), timeout=3.0)
except asyncio.TimeoutError:
print("隊列已滿,放入超時")
八、線程安全問題與跨線程使用
asyncio.Queue 不是線程安全的。 如果你需要從非 asyncio 線程向隊列中投遞任務(wù),必須使用 loop.call_soon_threadsafe():
# 從其他線程安全地放入數(shù)據(jù) loop.call_soon_threadsafe(queue.put_nowait, item)
或者使用 asyncio.run_coroutine_threadsafe() 調(diào)用異步方法:
future = asyncio.run_coroutine_threadsafe(queue.put(item), loop) future.result() # 阻塞等待完成
九、與queue.Queue的對比
| 特性 | queue.Queue | asyncio.Queue |
|---|---|---|
| 線程安全 | 是(內(nèi)置鎖) | 否(單線程事件循環(huán)) |
| 阻塞方式 | 線程阻塞(真正掛起 OS 線程) | 協(xié)程掛起(讓出事件循環(huán)) |
| timeout 參數(shù) | get(timeout=5) | 需配合 asyncio.wait_for() |
| qsize() 可靠性 | 不可靠(多線程競爭) | 在 await 點之間可靠 |
| join() / task_done() | 支持 | 支持 |
| shutdown() | 3.13+ 支持 | 3.13+ 支持 |
| 適用場景 | 多線程 | asyncio 協(xié)程 |
十、高級用法與最佳實踐
10.1 優(yōu)雅關(guān)閉(Python 3.13+)
async def main():
queue = asyncio.Queue()
# ... 啟動 workers ...
# 生產(chǎn)完畢后,優(yōu)雅關(guān)閉隊列
queue.shutdown()
await queue.join()
10.2 Python 3.12 及之前的優(yōu)雅關(guān)閉——哨兵值模式
SENTINEL = object()
async def consumer(queue):
while True:
item = await queue.get()
if item is SENTINEL:
queue.task_done()
break
process(item)
queue.task_done()
# 向每個 consumer 發(fā)送一個哨兵
for _ in range(num_consumers):
await queue.put(SENTINEL)
await queue.join()
10.3 扇出/扇入(Fan-out / Fan-in)
async def pipeline():
raw_queue = asyncio.Queue()
processed_queue = asyncio.Queue()
# Stage 1: 多個 fetcher 往 raw_queue 放數(shù)據(jù)
fetchers = [asyncio.create_task(fetch(raw_queue)) for _ in range(5)]
# Stage 2: 多個 processor 從 raw_queue 取出,處理后放入 processed_queue
processors = [asyncio.create_task(process(raw_queue, processed_queue)) for _ in range(3)]
# Stage 3: 單個 writer 從 processed_queue 取出并寫入
writer = asyncio.create_task(write(processed_queue))
10.4 注意事項
- 無界隊列的內(nèi)存風(fēng)險:
maxsize=0時生產(chǎn)速度遠(yuǎn)大于消費速度會導(dǎo)致內(nèi)存暴漲,生產(chǎn)環(huán)境應(yīng)始終設(shè)置合理的maxsize - 忘記調(diào)用
task_done():會導(dǎo)致join()永遠(yuǎn)阻塞 - 異常處理:消費者中如果處理邏輯拋出異常,必須確保
task_done()仍被調(diào)用(用try/finally):async def safe_consumer(queue): while True: item = await queue.get() try: await process(item) finally: queue.task_done() - 不要混用線程:除非通過
call_soon_threadsafe或run_coroutine_threadsafe
十一、總結(jié)
asyncio.Queue 是異步編程中不可或缺的協(xié)調(diào)原語。它的設(shè)計簡潔優(yōu)雅——底層通過 Future 實現(xiàn)掛起/喚醒,通過 deque 保證 FIFO 公平性,通過 _unfinished_tasks + Event 實現(xiàn) join 語義。理解它的內(nèi)部機制,能幫助你在構(gòu)建異步任務(wù)調(diào)度器、流水線處理、背壓控制等場景中做出更好的架構(gòu)決策。
到此這篇關(guān)于Python中asyncio.Queue異步隊列的實現(xiàn)的文章就介紹到這了,更多相關(guān)Python asyncio.Queue異步隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Python pass語句作用和Python assert斷言函數(shù)的用法
這篇文章主要介紹了Python pass語句作用和Python assert斷言函數(shù)的用法,文章內(nèi)容介紹詳細(xì)具有一定的參考價值,需要的小伙伴可以參考一下,希望對你有所幫助2022-03-03
Python tkinter之Bind(綁定事件)的使用示例
這篇文章主要介紹了Python tkinter之Bind(綁定事件)的使用詳解,幫助大家更好的理解和學(xué)習(xí)python的gui開發(fā),感興趣的朋友可以了解下2021-02-02
python實現(xiàn)文件名批量替換和內(nèi)容替換
這篇文章主要介紹了python實現(xiàn)文件名批量替換和內(nèi)容替換,第一個例子可以指定文件類型,需要的朋友可以參考下2014-03-03
Python tkinter之ComboBox(下拉框)的使用簡介
這篇文章主要介紹了Python tkinter之ComboBox(下拉框)的使用簡介,幫助大家更好的理解和使用python,感興趣的朋友可以了解下2021-02-02
幾種查看PyTorch、cuda 和 Python 版本方法小結(jié)
本文主要介紹了幾種查看PyTorch、cuda 和 Python 版本方法小結(jié),除了直接使用?torch.__version__?和?sys.version,我們還可以通過其他方式實現(xiàn)相同的功能,下面就一起來了解一下2025-04-04
Python爬蟲爬取新浪微博內(nèi)容示例【基于代理IP】
這篇文章主要介紹了Python爬蟲爬取新浪微博內(nèi)容,結(jié)合實例形式分析了Python基于代理IP實現(xiàn)的微博爬取與抓包分析相關(guān)操作技巧,需要的朋友可以參考下2018-08-08

