最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Python中asyncio.Queue異步隊列的實現(xiàn)

 更新時間:2026年05月15日 09:12:07   作者:無風(fēng)聽海  
本文主要介紹了Python中asyncio.Queue異步隊列的實現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧

一、概述

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.QueueEmptyget_nowait() 在空隊列上調(diào)用
asyncio.QueueFullput_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)鍵要點

  1. maxsize 實現(xiàn)背壓(backpressure):當(dāng)消費者處理速度跟不上時,生產(chǎn)者自動被掛起,防止內(nèi)存無限增長
  2. task_done() + join() 配對使用:確保每個任務(wù)都被處理完畢后再收尾
  3. 消費者用 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.Queueasyncio.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 注意事項

  1. 無界隊列的內(nèi)存風(fēng)險maxsize=0 時生產(chǎn)速度遠(yuǎn)大于消費速度會導(dǎo)致內(nèi)存暴漲,生產(chǎn)環(huán)境應(yīng)始終設(shè)置合理的 maxsize
  2. 忘記調(diào)用 task_done():會導(dǎo)致 join() 永遠(yuǎn)阻塞
  3. 異常處理:消費者中如果處理邏輯拋出異常,必須確保 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()
    
  4. 不要混用線程:除非通過 call_soon_threadsaferun_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ù)的用法

    這篇文章主要介紹了Python pass語句作用和Python assert斷言函數(shù)的用法,文章內(nèi)容介紹詳細(xì)具有一定的參考價值,需要的小伙伴可以參考一下,希望對你有所幫助
    2022-03-03
  • Python tkinter之Bind(綁定事件)的使用示例

    Python tkinter之Bind(綁定事件)的使用示例

    這篇文章主要介紹了Python tkinter之Bind(綁定事件)的使用詳解,幫助大家更好的理解和學(xué)習(xí)python的gui開發(fā),感興趣的朋友可以了解下
    2021-02-02
  • python實現(xiàn)文件名批量替換和內(nèi)容替換

    python實現(xiàn)文件名批量替換和內(nèi)容替換

    這篇文章主要介紹了python實現(xiàn)文件名批量替換和內(nèi)容替換,第一個例子可以指定文件類型,需要的朋友可以參考下
    2014-03-03
  • Pygame鼠標(biāo)進行圖片的移動與縮放案例詳解

    Pygame鼠標(biāo)進行圖片的移動與縮放案例詳解

    pygame是Python的第三方庫,里面提供了使用Python開發(fā)游戲的基礎(chǔ)包。本文將介紹如何通過Pygame實現(xiàn)鼠標(biāo)進行圖片的移動與縮放,感興趣的可以關(guān)注一下
    2021-12-12
  • python自動下載圖片的方法示例

    python自動下載圖片的方法示例

    這篇文章主要介紹了python自動下載圖片的方法示例,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-03-03
  • python微信公眾號之關(guān)鍵詞自動回復(fù)

    python微信公眾號之關(guān)鍵詞自動回復(fù)

    這篇文章主要為大家詳細(xì)介紹了python微信公眾號之關(guān)鍵詞自動回復(fù),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2018-06-06
  • Python tkinter之ComboBox(下拉框)的使用簡介

    Python tkinter之ComboBox(下拉框)的使用簡介

    這篇文章主要介紹了Python tkinter之ComboBox(下拉框)的使用簡介,幫助大家更好的理解和使用python,感興趣的朋友可以了解下
    2021-02-02
  • 幾種查看PyTorch、cuda 和 Python 版本方法小結(jié)

    幾種查看PyTorch、cuda 和 Python 版本方法小結(jié)

    本文主要介紹了幾種查看PyTorch、cuda 和 Python 版本方法小結(jié),除了直接使用?torch.__version__?和?sys.version,我們還可以通過其他方式實現(xiàn)相同的功能,下面就一起來了解一下
    2025-04-04
  • Pysvn 使用指南

    Pysvn 使用指南

    本文主要介紹了Pysvn 使用指南,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2023-02-02
  • Python爬蟲爬取新浪微博內(nèi)容示例【基于代理IP】

    Python爬蟲爬取新浪微博內(nèi)容示例【基于代理IP】

    這篇文章主要介紹了Python爬蟲爬取新浪微博內(nèi)容,結(jié)合實例形式分析了Python基于代理IP實現(xiàn)的微博爬取與抓包分析相關(guān)操作技巧,需要的朋友可以參考下
    2018-08-08

最新評論

盐亭县| 嘉善县| 乌拉特后旗| 精河县| 酒泉市| 正宁县| 新建县| 德昌县| 青浦区| 建湖县| 徐水县| 遂宁市| 壤塘县| 平湖市| 临洮县| 措勤县| 通河县| 晋城| 明水县| 屯门区| 衡东县| 博爱县| 永春县| 湘阴县| 洪泽县| 泉州市| 沙田区| 越西县| 浮梁县| 英超| 犍为县| 伊川县| 凉城县| 册亨县| 亳州市| 湾仔区| 徐闻县| 和龙市| 天气| 炎陵县| 南投市|