asyncio.Queue的maxsize不是背压而是阻塞信号,真正背压需上游显式await queue.put();端到端背压需组合queue、Semaphore及cancel-aware worker,并确保异常和取消时资源正确释放。

asyncio.Queue 的 maxsize 参数不是背压,而是阻塞信号
很多人以为给 asyncio.Queue(maxsize=100) 就算加了背压,其实这只是让 put() 在队列满时挂起协程——它不阻止上游生产者继续调用 put(),只是让调用卡在 await 上。一旦上游没做 await 或用了 create_task() 丢弃返回值,背压就完全失效,内存仍会暴涨。
真正起作用的是:上游必须显式 await queue.put(item),且不能跳过这个 await。常见错误包括:
- 用
asyncio.create_task(queue.put(item))异步提交,等于放任生产速度失控 - 在 for 循环里批量
put()却只在最后 await 一次(或根本不 await) - 把
queue当作“缓冲区”而非“流控关卡”,误以为 size 限制 = 流量限制
用 asyncio.Semaphore 实现端到端请求级背压
当流水线环节涉及外部调用(如 HTTP 请求、数据库写入),光靠队列不够,得把资源消耗也纳入控制。比如下游服务每秒最多处理 50 个请求,那上游每生成一个任务,就得先 acquire 一个许可。
实操建议:
立即学习“Python免费学习笔记(深入)”;
- 初始化一个
semaphore = asyncio.Semaphore(50),放在流水线共享作用域内 - 每个任务进入处理前,
await semaphore.acquire();处理完(无论成功失败)必须semaphore.release(),推荐用try/finally - 不要在
put()前 acquire —— 那只是控制入队,不是控制执行;要在实际执行逻辑开始前 acquire - 如果下游是多个并发 worker,确保 semaphore 是它们共用的同一实例,而不是每个 worker 自己 new 一个
这样,即使队列为空,新任务也无法挤占执行槽位,天然实现请求粒度的反压。
组合 queue + semaphore + cancel-aware worker 的最小可靠模式
单独用 queue 或 semaphore 都有盲区:queue 控制缓冲,semaphore 控制并发,但 worker 挂掉、超时、被 cancel 时,若不及时归还资源,会导致整个流水线死锁或饥饿。
一个健壮 worker 应该长这样:
async def worker(queue: asyncio.Queue, sem: asyncio.Semaphore):
while True:
try:
item = await queue.get()
try:
await sem.acquire()
await process_item(item) # 真正耗时操作
finally:
sem.release()
except asyncio.CancelledError:
# 必须在这里 release,否则 semaphore 永远少一个
if sem.locked():
sem.release()
raise
finally:
queue.task_done()
关键点:
-
queue.task_done()必须在 finally 中调用,否则queue.join()永远等不到完成 -
sem.release()要在 inner try/finally 里,确保即使process_item()抛异常或被 cancel,许可也能归还 - 不要在
process_item()外层包asyncio.wait_for()后直接 catch TimeoutError —— 这会吞掉 CancelledError,导致 semaphore 泄漏
警惕 aiohttp 和 aiomysql 默认无背压的陷阱
很多异步库默认不参与你的背压体系。例如:
-
aiohttp.ClientSession.post()返回ClientResponse对象,但你不 await.text()或.read(),响应体仍在内存中缓存,连接也不释放 -
aiomysql.Cursor.execute()只发 SQL,真正取结果要靠fetchone()或fetchall()—— 如果只 execute 不 fetch,连接会被占着,且结果集可能在服务端堆积 - 使用
async for row in cursor时,若中途 break 或 raise,需确保 cursor.close() 或连接 return 到池中
对策很简单:所有 IO 操作链必须完整 await,且在异常路径中显式 cleanup。比如用 aiohttp 时,习惯性写成:
async with session.post(url, json=payload) as resp:
await resp.text() # 强制读完,释放连接
而不是:
resp = await session.post(url, json=payload) # 忘了 await resp.text() → 连接泄漏 + 内存上涨
背压不是加个队列或限速器就能自动生效的机制,它是整个调用链上每个 await 点对资源获取与释放的精确配对。最容易被忽略的,永远是异常分支里的 cleanup 和 cancel 传播路径上的 semaphore 归还。


















