↓快轉到主要內容
  1. 教學文章/

Python asyncio.Queue 實戰:背壓、Graceful Shutdown 與不丟資料測試

·8 分鐘· loading · loading · ·
Python Asyncio Queue Backpressure Graceful Shutdown Concurrency Testing
每日拍拍
作者
每日拍拍
科學家 X 科技宅宅
目錄
Python 學習 - 本文屬於一個選集。
§ 137: 本文

featured

非同步程式一開始常常很順。

把輸入包成 task,全部 create_task(),看起來每秒能吞好多工作。 直到上游比下游快十倍,記憶體一路漲,關閉服務時又發現還有幾百筆資料不知道處理到哪裡。

這不是 async 語法問題,而是流量與生命週期沒有被設計。

這篇拍拍君要用標準庫 asyncio.Queue,做一條有容量上限、能自然施加背壓、可以排空後關閉,並且能測出漏件與重複處理的工作管線。

如果還不熟 coroutine、await 與 task,可以先看 Python asyncio 非同步程式設計入門;本文不再重講基礎語法。

一. 先定義問題:快的 Producer 會淹沒慢的 Consumer
#

假設上游每 1 毫秒產生一筆工作,下游卻要 50 毫秒才能處理。

如果每筆工作一到就直接建立 task:

tasks = [asyncio.create_task(handle(item)) for item in huge_stream]
await asyncio.gather(*tasks)

輸入量一大,程式會同時保留大量:

  • task 物件;
  • coroutine frame;
  • request payload;
  • retry 與 log context;
  • 尚未釋放的 socket 或 buffer。

Semaphore 可以限制「同時進入某段程式」的數量,但如果你先建立十萬個 task,再讓它們卡在 Semaphore 前,十萬個 task 還是存在。

Queue 的角色不同:它在工作進入管線的邊界就限制庫存。

producer -- await queue.put() --> [ bounded buffer ] --> consumer
                   ^                    maxsize
                   | full 時暫停

當 Queue 滿了,await queue.put(item) 會等待空位。

上游因此自動降速,這就是背壓(backpressure)。

二. 環境與版本:Queue.shutdown 需要 Python 3.13+
#

asyncio 是標準庫,不用安裝套件。

本文使用 Python 3.13 新增的 Queue.shutdown() 與 QueueShutDown:

python --version
# Python 3.13+,本文實測為 Python 3.14.3

建立練習專案:

uv init queue-lab --python 3.14
cd queue-lab
uv add --dev pytest pytest-asyncio

如果專案還停在 Python 3.12,仍可用 sentinel 物件通知 worker 結束;但 sentinel 的數量、插入時機和 producer 失敗處理都要自己維護。

新版 shutdown API 讓「停止接收新工作」成為 Queue 自己的狀態,邊界更清楚。

三. 第一條有背壓的管線
#

先做一個最小範例。

import asyncio


async def producer(queue: asyncio.Queue[int]) -> None:
    for item in range(10):
        await queue.put(item)
        print("put", item, "qsize=", queue.qsize())


async def consumer(queue: asyncio.Queue[int]) -> None:
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            return

        try:
            await asyncio.sleep(0.05)
            print("done", item)
        finally:
            queue.task_done()


async def main() -> None:
    queue: asyncio.Queue[int] = asyncio.Queue(maxsize=3)

    async with asyncio.TaskGroup() as group:
        group.create_task(consumer(queue))
        await producer(queue)
        queue.shutdown()
        await queue.join()


asyncio.run(main())

Queue 最多只保留三筆等待中的工作。

Consumer 忙碌、Queue 又滿時,Producer 會停在 put();Consumer 取走一筆後,Producer 才能繼續。

這個等待不是錯誤,也不是效能退化。

它是系統在說:

下游目前只能處理這麼快,請上游不要再堆庫存。

四. maxsize 不是吞吐量旋鈕
#

maxsize=1000 不會讓 Consumer 變快。

它只決定可以吸收多少短暫波動。

估算時可以從 Little’s Law 的直覺出發:

排隊數量 ≈ 到達速率 × 可接受等待時間

例如下游每秒穩定處理 200 筆,而你願意容忍 0.5 秒的短暫突發:

200 × 0.5 = 100

可以先用 maxsize=100 做壓力測試,再觀察:

  • qsize() 是否長時間貼著上限;
  • put() 等待時間是否持續增加;
  • 單筆端到端 latency;
  • Consumer 的成功、失敗與重試率;
  • 關閉時排空 Queue 需要多久。

如果 Queue 永遠滿,問題通常不是「容量太小」,而是下游處理能力不足或輸入沒有節流。

盲目把容量放大,只是把過載改成晚一點爆。

五. task_done 與 join:追蹤的是未完成工作,不是 Queue 長度
#

每次成功 put(),Queue 內部的 unfinished-task counter 就加一。

Consumer 每完成一筆,必須呼叫一次 task_done();計數歸零後,join() 才會返回。

item = await queue.get()
try:
    await handle(item)
finally:
    queue.task_done()

finally 很重要。

如果 handle() 丟出例外、task 被取消,或程式提早 return,仍要替已經 get() 的那筆工作結清 Queue 計數。

但這裡要分清楚兩件事:

  1. task_done() 表示 Queue 不再追蹤這筆工作;
  2. 它不代表業務資料一定成功寫入。

若處理失敗後必須重試,應先把失敗記錄到 retry store、dead-letter queue 或持久化儲存,再決定是否結清原工作。

少呼叫一次,join() 可能永遠卡住。

多呼叫一次,Python 會丟出 ValueError。

所以最穩定的寫法是:每個成功 get() 對應一個、也只有一個 task_done()。

六. Graceful Shutdown 的正確順序
#

正常關閉不是先取消 Consumer。

那會把 Queue 裡尚未處理的工作留在原地。

排空優先(drain-first)的順序是:

  1. 停止接受新工作;
  2. 讓已在 Queue 裡的工作繼續被取出;
  3. 等待 unfinished-task counter 歸零;
  4. Consumer 看到 Queue 已關閉且為空,自然離開;
  5. 最後才結束整個 task scope。
await produce_all(queue)
queue.shutdown(immediate=False)
await queue.join()

shutdown() 預設的 immediate=False 會禁止新的 put(),但保留既有項目供 Consumer 排空。

Queue 清空後,再呼叫 get() 會得到 QueueShutDown。

這比「取消 worker,祈禱它剛好處理完」可靠得多。

七. 完整實作:多 Producer、多 Consumer
#

下面把資料模型、指標與關閉流程放在一起。

# pipeline.py
from __future__ import annotations

import asyncio
from collections import Counter
from dataclasses import dataclass
from typing import Iterable


@dataclass(frozen=True, slots=True)
class WorkItem:
    item_id: str
    payload: str


async def handle(item: WorkItem) -> str:
    await asyncio.sleep(0.01)
    return item.payload.upper()


async def producer(
    name: str,
    queue: asyncio.Queue[WorkItem],
    items: Iterable[WorkItem],
) -> None:
    for item in items:
        await queue.put(item)
        print(f"{name}: queued {item.item_id}; depth={queue.qsize()}")


async def consumer(
    name: str,
    queue: asyncio.Queue[WorkItem],
    completed: Counter[str],
) -> None:
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            print(f"{name}: drained")
            return

        try:
            result = await handle(item)
            completed[item.item_id] += 1
            print(f"{name}: {item.item_id} -> {result}")
        finally:
            queue.task_done()


async def run_pipeline(
    batches: list[list[WorkItem]],
    *,
    workers: int = 3,
    capacity: int = 8,
) -> Counter[str]:
    queue: asyncio.Queue[WorkItem] = asyncio.Queue(maxsize=capacity)
    completed: Counter[str] = Counter()

    try:
        async with asyncio.TaskGroup() as group:
            for number in range(workers):
                group.create_task(
                    consumer(f"worker-{number}", queue, completed)
                )

            producers = [
                group.create_task(producer(f"source-{n}", queue, batch))
                for n, batch in enumerate(batches)
            ]
            await asyncio.gather(*producers)

            queue.shutdown()
            await queue.join()
    finally:
        # 失敗或外部取消時,解除卡在 put/get 的 task。
        # immediate=True 是緊急出口,不代表剩餘工作已完成。
        queue.shutdown(immediate=True)

    return completed


async def main() -> None:
    batches = [
        [WorkItem(f"a-{n}", f"alpha-{n}") for n in range(5)],
        [WorkItem(f"b-{n}", f"beta-{n}") for n in range(5)],
    ]
    completed = await run_pipeline(batches)
    print(completed)


if __name__ == "__main__":
    asyncio.run(main())

這裡讓 Consumer 和 Producer 都屬於同一個 TaskGroup。

任一 task 失敗時,TaskGroup 會取消其他 task,避免 Producer 永遠卡在已無 Consumer 的滿 Queue。

成功路徑則先等所有 Producer 結束,再進入 drain-first shutdown。

八. immediate=True 是緊急煞車,不是成功完成
#

queue.shutdown(immediate=True) 會立即排空 Queue,並喚醒卡住的 put() 與 get()。

它很適合放在失敗清理路徑,避免整個程式死鎖。

但它會破壞 join() 平常代表的保證:join() 可能在工作沒有真正被處理時解除等待。

所以不要寫成:

queue.shutdown(immediate=True)
await queue.join()
print("全部成功")  # 這個結論不成立

比較好的語意是:

  • shutdown():停止收件,既有工作排空;
  • shutdown(immediate=True):系統正在中止,先解除等待;
  • 業務成功與否:由持久化結果、ack 或完成紀錄判斷。

如果工作不能遺失,只靠記憶體 Queue 本來就不夠。

服務 crash、機器斷電時,記憶體內的項目不會自己回來;這類需求要搭配資料庫 outbox、Redis Streams、NATS、RabbitMQ 或其他 durable broker。

九. Cancellation:清理後要繼續傳遞
#

Task 被取消時,下一個可取消的 await 會丟出 CancelledError。

Consumer 可以用 finally 清理,但通常不要吞掉取消:

async def cancellable_consumer(queue: asyncio.Queue[int]) -> None:
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            return

        try:
            await save(item)
        except asyncio.CancelledError:
            await record_interrupted(item)
            raise
        finally:
            queue.task_done()

重新 raise 很重要,因為 TaskGroup 與 asyncio.timeout() 都依賴 cancellation 運作。

若把 CancelledError 當普通錯誤吃掉,外層可能以為 task 已正常停止,實際上它還在跑。

另外,task_done() 應放哪裡取決於你的交付語意。

若 record_interrupted() 只是寫 log,不足以保證稍後重試,就不能宣稱資料已可靠保存。

十. Timeout 要放在邊界,不能假裝成功
#

Queue 本身的 put() 與 get() 沒有 timeout 參數。

可以在呼叫端設定等待上限:

try:
    async with asyncio.timeout(0.5):
        await queue.put(item)
except TimeoutError:
    await overflow_store.save(item)

這段的語意是:上游最多接受 0.5 秒背壓;超過後,把工作轉存到另一個可靠位置。

如果 timeout 後直接丟掉 item,系統當然比較快,但資料也真的不見了。

監控上至少要分開記錄:

  • enqueue 成功;
  • enqueue timeout;
  • shutdown 拒絕;
  • handler 成功;
  • handler 失敗;
  • forced shutdown 丟棄數量。

十一. 測試一:每筆剛好完成一次
#

非同步測試不要依賴「sleep 久一點應該會跑完」。

用確定的輸入與完成計數檢查不漏件、不重複:

# test_pipeline.py
import pytest

from pipeline import WorkItem, run_pipeline


@pytest.mark.asyncio
async def test_every_item_completes_exactly_once() -> None:
    batches = [
        [WorkItem(f"left-{n}", str(n)) for n in range(20)],
        [WorkItem(f"right-{n}", str(n)) for n in range(20)],
    ]

    completed = await run_pipeline(batches, workers=4, capacity=2)

    expected = {
        *(f"left-{n}" for n in range(20)),
        *(f"right-{n}" for n in range(20)),
    }
    assert set(completed) == expected
    assert all(count == 1 for count in completed.values())

刻意使用很小的 capacity=2,讓測試真的走過 Queue 滿載和 Producer 等待的路徑。

如果測試只用超大 Queue,背壓相關 bug 可能完全沒有被碰到。

十二. 測試二:證明 Producer 真的會被背壓
#

用 Event 控制 Consumer 何時開始,不靠時間猜測:

import asyncio
import pytest


@pytest.mark.asyncio
async def test_full_queue_blocks_the_next_put() -> None:
    queue: asyncio.Queue[str] = asyncio.Queue(maxsize=1)
    await queue.put("first")

    second_put = asyncio.create_task(queue.put("second"))
    await asyncio.sleep(0)

    assert not second_put.done()

    assert await queue.get() == "first"
    queue.task_done()

    await second_put
    assert queue.qsize() == 1

    assert await queue.get() == "second"
    queue.task_done()
    await queue.join()

await asyncio.sleep(0) 只是把控制權交回 event loop 一輪,不是等待任意秒數。

這個測試直接驗證:容量為一時,第二次 put() 在 Consumer 釋出空位前不會完成。

十三. 測試三:關閉後仍會排空既有工作
#

最後驗證 graceful shutdown 的核心契約:

import asyncio
import pytest


@pytest.mark.asyncio
async def test_shutdown_drains_existing_items() -> None:
    queue: asyncio.Queue[int] = asyncio.Queue(maxsize=3)
    for item in range(3):
        await queue.put(item)

    queue.shutdown()
    seen: list[int] = []

    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            break
        else:
            seen.append(item)
            queue.task_done()

    await queue.join()
    assert seen == [0, 1, 2]

    with pytest.raises(asyncio.QueueShutDown):
        await queue.put(99)

這同時確認兩件事:

  • shutdown 前已入列的三筆工作仍可被取出;
  • shutdown 後的新工作會被拒絕。

十四. 上線前的檢查清單
#

把 Queue 放進正式服務前,拍拍君會逐項確認:

  1. Queue 是否有合理的 maxsize,而不是預設無上限?
  2. Producer 是否真的 await queue.put(),而不是忙碌重試 put_nowait()?
  3. 每次 get() 是否恰好對應一次 task_done()?
  4. 正常關閉是否先停止收件,再 join() 排空?
  5. Consumer 失敗時,Producer 是否會被取消或解除阻塞?
  6. immediate=True 是否只用於中止路徑?
  7. CancelledError 清理後是否重新拋出?
  8. timeout 後的工作要丟棄、重試,還是寫入 durable store?
  9. 測試是否用小容量真正觸發背壓?
  10. 指標能否區分排隊、處理、失敗與強制丟棄?

結語
#

asyncio.Queue 不只是把資料從 Producer 傳給 Consumer 的容器。

設好容量後,它是流量控制閥;搭配 task_done() 與 join(),它是完成追蹤器;搭配 Python 3.13+ 的 shutdown(),它又能清楚表達「停止收件,但把手上的做完」。

最重要的原則只有一句:

正常關閉先排空,異常關閉才強制中止;任何「不丟資料」都要由可驗證的完成紀錄支持。

先把 Queue 設成小容量,寫出會觸發背壓的測試,再逐步量測與調整。這比先塞十萬個 task、出事後才找記憶體去哪裡,舒服多了。

延伸閱讀
#

Python 學習 - 本文屬於一個選集。
§ 137: 本文

相關文章

Python contextvars 實戰:Async Context、結構化 Logging 與隔離測試
·8 分鐘· loading · loading
Python Contextvars Asyncio Logging Concurrency Testing Standard-Library
Textual DirectoryTree 檔案瀏覽器:安全路徑、非同步預覽與互動測試
·8 分鐘· loading · loading
Python Textual DirectoryTree TUI Pathlib Asyncio Testing
Polars Rolling / Dynamic Windows:時間窗聚合、邊界語意與 QA
·6 分鐘· loading · loading
Python Polars Time-Series Rolling-Window Dynamic-Window Data-Engineering Testing
Python concurrent.futures 實戰:Thread、Process 與 Interpreter Pool
·8 分鐘· loading · loading
Python Concurrent.futures ThreadPoolExecutor ProcessPoolExecutor InterpreterPoolExecutor Concurrency Parallelism
Python AnyIO 實戰:TaskGroup、取消管理與跨框架非同步工具
·6 分鐘· loading · loading
Python AnyIO Async Asyncio Trio Developer-Tools
Python pytest fixtures 進階:conftest、factory 與測試資料管理
·8 分鐘· loading · loading
Python Pytest Fixtures Testing Conftest Monkeypatch Developer-Tools