Python异步代码异常处理:生产者消费者模型代码挂起排查求助
生产者-消费者挂起问题分析与修复
根本原因
你的代码会挂起,完全是因为queue.join()在等一个永远不会发生的事——它会一直阻塞,直到队列里所有被取出的元素都调用了queue.task_done()来标记完成。但当consumer处理到元素3抛出异常时,虽然你在except块里调用了一次task_done(),但队列里还剩4、5、6这几个元素,以及后续main放入的None,这些元素要么没被consumer取出,要么取出后没机会调用task_done(),导致queue.join()无限等待,程序直接卡死。
具体流程拆解:
- 生产者把0到6共7个元素塞进队列,然后结束
- 消费者逐个取出元素处理,拿到3的时候抛出
TypeError - 消费者的except块打印错误、调用一次
task_done()后重新抛出异常,消费者任务直接终止 - 此时队列里还有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())
关键修改点
- 给每个队列元素加finally块:不管处理时有没有抛异常,只要取出了元素,就必须调用
task_done(),彻底避免队列任务残留 - 调整main的等待逻辑:先直接等待consumer任务完成(顺便捕获异常),再清理队列里剩下的元素,最后调用
queue.join(),确保不会无限等待 - 删掉无效的cancel操作:consumer要么正常结束要么异常终止,主动cancel完全没必要,反而可能搞出额外错误
内容的提问来源于stack exchange,提问作者Khachatur Sarkisyan
相关产品推荐
相关产品推荐

