基于threading.Condition的队列添加async get方法后取消超时问题
问题根本原因
asyncio.wait_for 触发超时取消协程时,只会取消asyncio层面的协程任务,不会终止已经提交到线程池的同步阻塞任务,你遇到的问题本质是每次超时都会在后台遗留一个卡在self.condition.wait()的get()调用,后续生产者放入的新元素都会被这些遗留的线程提前取走,导致新提交的get_async永远拿不到元素,持续触发超时。
具体触发流程:
- 消费者取完
foo、bar后队列空,第一次调用get_async,向线程池提交self.get()任务 - 后台线程进入
get()方法,拿到condition锁,发现队列为空,调用condition.wait()释放锁并阻塞 - 5秒超时后,
wait_for取消get_async协程,抛出TimeoutError,但后台的get()调用仍处于阻塞状态,没有被终止 - 等到第7秒(producer sleep 2秒叠加前面的5秒),producer放入
baz并调用notify(),那个后台遗留的get()被唤醒,拿到锁取走baz,但对应的协程早已被取消,这个baz直接被丢弃 - 后续每次
get_async超时都会多遗留一个阻塞的get()线程,所有新放入队列的元素都会被这些遗留线程抢先取走,永远到不了当前正在等待的协程里,所以会一直输出超时。
解决思路
核心要解决的是协程被取消时,同步终止掉线程池里对应的阻塞get()调用,避免遗留线程偷取后续元素,可选方案:
- 如果不需要跨线程生产消费,直接用标准库的
asyncio.Queue或者基于asyncio.Condition实现异步队列,不要混用threading同步原语和asyncio逻辑 - 如果必须要支持跨线程生产+异步消费,可以给每个
get_async调用生成专属取消标记:get()方法的等待循环同时检查队列是否为空、以及对应取消标记是否被触发- 协程被取消时,触发对应标记,再调用一次
condition.notify_all(),让遗留的get()检测到取消后主动退出,不要取走队列元素
- 也可以把
condition.wait()换成带超时的condition.wait(timeout=1),定期唤醒检查是否要终止,避免永久阻塞。
内容的提问来源于stack exchange,提问作者user4344
相关产品推荐
相关产品推荐

