You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

asyncio队列异常问题:队列为空时get()未等待却触发任务销毁错误

asyncio队列异常问题:队列为空时get()未等待却触发任务销毁错误

嘿,我来帮你捋清楚这个问题~其实你遇到的不是asyncio.Queue.get()的问题,它确实会乖乖等待队列有消息再返回,真正的问题出在你的事件循环过早关闭了!

咱们来拆解一下你的代码逻辑:

  • 你调用asyncio.run(launch()),这个函数会启动事件循环,直到launch()协程执行完毕,然后立刻关闭事件循环。
  • 在launch()里,你只做了await asyncio.gather(queue.start())——而queue.start()只是创建了一个__send任务就返回了,所以launch()很快就执行完了。
  • 这时候事件循环马上关闭,但__send任务还卡在await self.queue.get()这一步(处于pending状态),事件循环关闭时会销毁所有pending的任务,于是就出现了你看到的Task was destroyed but it is pending!错误。

至于你说加了if not self.queue.empty()就“能工作”,其实这只是个假象:当队列为空时,代码跳过了await queue.get(),__send会进入无限空循环,这时候任务没有处于await的pending状态,但事件循环关闭时它依然会被销毁,只是可能没触发明显的错误提示而已——而且这种写法本身就有问题:empty()只是瞬时状态,在你判断完empty()和调用get()之间,可能有其他协程往队列里放了消息,导致你漏掉这些消息,完全不符合队列的使用规范。

那正确的做法应该是怎样的?核心就是让事件循环保持运行,直到你主动结束任务,下面给你修改后的完整代码,我还加了一些最佳实践的细节:

import asyncio
import logging

# 先配置好日志,不然你的logger会报错
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class MessageQueue:
    def __init__(self):
        self.proc = None
        self.queue = asyncio.Queue()
    
    async def start(self):
        # 给任务起个名字,方便调试
        self.proc = asyncio.create_task(self.__send(), name='message_queue')
    
    async def stop(self):
        if self.proc is not None:
            self.proc.cancel()
            try:
                # 等待任务处理完CancelledError,确保任务优雅结束
                await self.proc
            except asyncio.CancelledError:
                logger.info('队列任务已成功取消')
    
    async def put(self, mobj):
        await self.queue.put(mobj)
    
    async def __send(self):
        logger.info('启动队列处理任务')
        try:
            while True:
                # 这里直接await get()就好,它会自动等待消息
                msg = await self.queue.get()
                logger.info(f'收到消息: {msg}')
                # 记得调用task_done,后续如果需要用queue.join()等待所有消息处理完会用到
                self.queue.task_done()
        except asyncio.CancelledError:
            logger.info('正在停止队列处理任务')
            raise  # 重新抛出异常,让stop方法里的await捕获
        except Exception as e:
            logger.error(f'处理消息时出错: {e}')


async def launch():
    queue = MessageQueue()
    await queue.start()
    
    # 模拟一个定时发消息的任务,验证队列能正常接收
    async def test_send_messages():
        await asyncio.sleep(2)
        await queue.put("Hello, asyncio队列!")
        await asyncio.sleep(1)
        await queue.put("这是另一条测试消息")
    
    asyncio.create_task(test_send_messages())
    
    try:
        # 用一个永远不会触发的Event让事件循环保持运行,直到你按Ctrl+C中断
        await asyncio.Event().wait()
    except KeyboardInterrupt:
        logger.info('收到中断信号,正在停止队列...')
        await queue.stop()

asyncio.run(launch())

关键修改点说明:

  1. 补上了日志配置,避免原代码中logger未初始化的错误
  2. 在stop()方法里,不仅cancel任务,还await self.proc,确保任务能优雅处理取消信号
  3. __send()里正确捕获CancelledError,保证任务能正常退出
  4. launch()里添加了测试发消息的逻辑,同时用await asyncio.Event().wait()让事件循环一直运行,直到你按下Ctrl+C触发中断,再主动停止队列任务
  5. 加上了queue.task_done(),这是asyncio队列的最佳实践,配合queue.join()可以等待所有已入队的消息都被处理完

这样修改后,你运行代码就能看到队列正常接收消息,按Ctrl+C也能优雅停止,不会再出现任务销毁的错误啦~

备注:内容来源于stack exchange,提问作者Ramesses III

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 17:04:35