非同步程式一開始常常很順。
把輸入包成 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 計數。
但這裡要分清楚兩件事:
task_done()表示 Queue 不再追蹤這筆工作;- 它不代表業務資料一定成功寫入。
若處理失敗後必須重試,應先把失敗記錄到 retry store、dead-letter queue 或持久化儲存,再決定是否結清原工作。
少呼叫一次,join() 可能永遠卡住。
多呼叫一次,Python 會丟出 ValueError。
所以最穩定的寫法是:每個成功 get() 對應一個、也只有一個 task_done()。
六. Graceful Shutdown 的正確順序 #
正常關閉不是先取消 Consumer。
那會把 Queue 裡尚未處理的工作留在原地。
排空優先(drain-first)的順序是:
- 停止接受新工作;
- 讓已在 Queue 裡的工作繼續被取出;
- 等待 unfinished-task counter 歸零;
- Consumer 看到 Queue 已關閉且為空,自然離開;
- 最後才結束整個 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 放進正式服務前,拍拍君會逐項確認:
- Queue 是否有合理的
maxsize,而不是預設無上限? - Producer 是否真的
await queue.put(),而不是忙碌重試put_nowait()? - 每次
get()是否恰好對應一次task_done()? - 正常關閉是否先停止收件,再
join()排空? - Consumer 失敗時,Producer 是否會被取消或解除阻塞?
immediate=True是否只用於中止路徑?CancelledError清理後是否重新拋出?- timeout 後的工作要丟棄、重試,還是寫入 durable store?
- 測試是否用小容量真正觸發背壓?
- 指標能否區分排隊、處理、失敗與強制丟棄?
結語 #
asyncio.Queue 不只是把資料從 Producer 傳給 Consumer 的容器。
設好容量後,它是流量控制閥;搭配 task_done() 與 join(),它是完成追蹤器;搭配 Python 3.13+ 的 shutdown(),它又能清楚表達「停止收件,但把手上的做完」。
最重要的原則只有一句:
正常關閉先排空,異常關閉才強制中止;任何「不丟資料」都要由可驗證的完成紀錄支持。
先把 Queue 設成小容量,寫出會觸發背壓的測試,再逐步量測與調整。這比先塞十萬個 task、出事後才找記憶體去哪裡,舒服多了。