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

Python异步代码异常处理:生产者消费者模型代码挂起排查求助

生产者-消费者挂起问题分析与修复

根本原因

你的代码会挂起,完全是因为queue.join()在等一个永远不会发生的事——它会一直阻塞,直到队列里所有被取出的元素都调用了queue.task_done()来标记完成。但当consumer处理到元素3抛出异常时,虽然你在except块里调用了一次task_done(),但队列里还剩4、5、6这几个元素,以及后续main放入的None,这些元素要么没被consumer取出,要么取出后没机会调用task_done(),导致queue.join()无限等待,程序直接卡死。

具体流程拆解:

  1. 生产者把0到6共7个元素塞进队列,然后结束
  2. 消费者逐个取出元素处理,拿到3的时候抛出TypeError
  3. 消费者的except块打印错误、调用一次task_done()后重新抛出异常,消费者任务直接终止
  4. 此时队列里还有4、5、6和后来的None没被处理,没人给这些元素发task_done()信号,queue.join()就一直卡着不动

修复方案

核心思路是不管消费者有没有抛出异常,每个被取出的队列元素都必须调用task_done(),同时调整main函数的逻辑,避免无效的等待。

修复后的完整代码

import asyncio


async def producer(queue):
    for i in range(7):
        await queue.put(i)
        await asyncio.sleep(0.2)
        print('Producer Queue size:', queue.qsize())


async def consumer(queue):
    try:
        while True:
            item = await queue.get()
            # 用finally块确保不管处理成功还是抛异常,都标记任务完成
            try:
                await asyncio.sleep(0.3)
                print(item)

                if item is None:
                    print('Queue is empty')
                    break

                if item == 3:
                    raise TypeError("Producer got exception")

                print('Queue size:', queue.qsize())
            finally:
                queue.task_done()
    except TypeError as e:
        print(f'Error happened asf {e}')
        # 把异常抛出去让main捕获
        raise e


async def main():
    queue = asyncio.Queue()
    producer_task = asyncio.create_task(producer(queue))
    consumer_task = asyncio.create_task(consumer(queue))

    # 等生产者把所有元素都放进队列
    await producer_task
    # 给消费者发结束信号
    await queue.put(None)
    
    # 直接等消费者任务完成,同时捕获它抛出的异常
    try:
        await consumer_task
    except TypeError as e:
        print(f"Main caught exception: {e}")
    
    # 清理队列里剩下的元素(防止消费者提前崩溃导致残留)
    while not queue.empty():
        item = await queue.get()
        queue.task_done()
    
    # 现在所有任务都标记完成了,join可以正常结束
    await queue.join()
    print("All tasks completed")


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

关键修改点

  1. 给每个队列元素加finally块:不管处理时有没有抛异常,只要取出了元素,就必须调用task_done(),彻底避免队列任务残留
  2. 调整main的等待逻辑:先直接等待consumer任务完成(顺便捕获异常),再清理队列里剩下的元素,最后调用queue.join(),确保不会无限等待
  3. 删掉无效的cancel操作:consumer要么正常结束要么异常终止,主动cancel完全没必要,反而可能搞出额外错误

内容的提问来源于stack exchange,提问作者Khachatur Sarkisyan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 07:18:14