
本文揭示了 asyncio.Queue 行为异常的真正原因:并非 get() 不等待,而是主协程过早退出导致后台任务被强制取消;核心在于正确管理异步生命周期,确保长期运行的任务不被提前终止。
本文揭示了 asyncio.queue 行为异常的真正原因:并非 `get()` 不等待,而是主协程过早退出导致后台任务被强制取消;核心在于正确管理异步生命周期,确保长期运行的任务不被提前终止。
asyncio.Queue.get() 的行为完全符合文档规范——当队列为空时,它会自动挂起协程并等待新元素入队,绝不会立即报错或返回。你遇到的 Task was destroyed but it is pending! 错误,根本不是队列逻辑的问题,而是异步程序生命周期失控所致。
问题根源在于 launch() 函数的执行流程:
async def launch():
queue = MessageQueue() # 创建实例
await asyncio.gather(queue.start()) # 启动后台任务后立即返回 → 程序结束!queue.start() 内部仅调用 asyncio.create_task(self.__send()) 并立即返回(它本身不 await 任何耗时操作),因此 await asyncio.gather(...) 在任务刚创建、甚至还没来得及进入 __send() 的 while True 循环前就完成了。此时 asyncio.run(launch()) 认为整个程序已“执行完毕”,开始清理资源——而仍在运行的 __send() 任务被强制取消,触发警告。
✅ 正确做法是:让主协程保持活跃,直到明确需要停止队列。以下是修复后的完整示例:
import asyncio
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class MessageQueue:
def __init__(self):
self.proc = None
self.queue = asyncio.Queue()
self._stop_event = asyncio.Event() # 用于优雅停止
async def start(self):
self.proc = asyncio.create_task(self.__send())
logger.info("MessageQueue started.")
async def stop(self):
if self.proc is not None:
self._stop_event.set() # 通知循环退出
await self.proc # 等待任务自然结束
self.proc = None
logger.info("MessageQueue stopped.")
async def put(self, mobj):
await self.queue.put(mobj)
async def __send(self):
logger.info("Starting message processor.")
while not self._stop_event.is_set():
try:
msg = await self.queue.get() # ✅ 正确:自动等待,无需 .empty() 检查
logger.info(f"Processing: {msg}")
self.queue.task_done() # 若需 task_done 语义(如 join)
except asyncio.CancelledError:
logger.info("Message processor cancelled.")
break
except Exception as e:
logger.error(f"Error processing message: {e}")
# 主函数:保持主协程活跃,模拟真实工作流
async def launch():
queue = MessageQueue()
await queue.start()
# 模拟生产消息
for i in range(3):
await queue.put(f"Message #{i}")
await asyncio.sleep(0.5)
# 停止前可再发一条
await queue.put("Final message")
await asyncio.sleep(0.3)
# 优雅停止
await queue.stop()
# ✅ 关键:asyncio.run 会等待 launch() 完全结束,从而保证后台任务有足够时间执行
if __name__ == "__main__":
asyncio.run(launch())? 关键要点总结:
- ❌ 错误认知:
queue.get()不等待 → 实际它永远等待,除非任务被取消; - ✅ 核心原则:
asyncio.run()的入口协程必须显式控制所有子任务的生命周期; - ?
asyncio.gather()在此场景中冗余且误导——它适用于并发等待多个有明确返回值的任务,而非启动长期守护任务; - ? 推荐模式:使用
asyncio.Event或asyncio.Condition实现可控停止,避免task.cancel()引发的异常中断; - ⚠️ 注意:永远不要在
while True循环中忽略CancelledError捕获,否则无法响应取消信号。
通过以上重构,__send() 将稳定运行,get() 会按预期阻塞等待,整个系统具备健壮性与可维护性。

















