为何threading.Condition.notify_all无法唤醒等待线程?
问题
我编写了如下代码以演示线程同步逻辑:
- 启动独立线程更新图像;
- 基于这些图像实现异步生成器(async generator);
- 仅当异步生成器使用图像后,才更新新图像;
- 异步生成器需等待新图像创建完成。
但代码运行时卡在等待第一张图像的环节,控制台输出如下:
# Output create new image waiting for new image start waiter notify_all wait for someone to take it waiting for image_created
对应代码如下:
import asyncio import random import threading import time class ImageUpdater: def __init__(self): self.image = None self.image_used = threading.Event() self.image_created = threading.Condition() def update_image(self): while True: self.image_used.clear() with self.image_created: print("create new image") time.sleep(0.6) self.image = str(random.random()) print("notify_all") self.image_created.notify_all() print("wait for someone to take it") self.image_used.wait() print("someone took it") async def image_generator(self): def waiter(): print("start waiter") time.sleep(0.1) with self.image_created: print("waiting for image_created") self.image_created.wait() print("waiter finished") self.image_used.set() while True: print("waiting for new image") await asyncio.to_thread(waiter) yield self.image async def main(): updater = ImageUpdater() update_thread = threading.Thread(target=updater.update_image) update_thread.start() async for image in updater.image_generator(): print(f"Received new image: {image}") if __name__ == "__main__": loop = asyncio.run(main())
请问为何threading.Condition.notify_all无法释放image_created.wait()的阻塞?
原因与解决方案
问题出在时序颠倒:更新线程先调用了notify_all,之后等待线程才进入wait(),导致通知被“错过”了。
threading.Condition的通知是一次性的——如果在wait()调用前就触发了notify_all,后续的wait()会一直阻塞,因为它没接到任何有效通知。从输出也能验证这一点:notify_all在waiting for image_created之前打印,说明更新线程已经发完通知,等待线程才开始等待,自然收不到信号。
修复方案
需要给Condition搭配一个状态标记,用来判断是否已有可用图像,避免错过通知。修改步骤如下:
- 在
ImageUpdater中添加has_image布尔变量,标记是否有未被取用的图像; - 更新线程创建完图像后,设置
has_image = True再发通知; - 等待线程进入
wait()前,先检查has_image,如果已经为True就直接跳过等待,否则再等待通知。
修改后的完整代码:
import asyncio import random import threading import time class ImageUpdater: def __init__(self): self.image = None self.image_used = threading.Event() self.image_created = threading.Condition() self.has_image = False # 新增状态标记 def update_image(self): while True: self.image_used.clear() with self.image_created: print("create new image") time.sleep(0.6) self.image = str(random.random()) self.has_image = True # 标记图像已创建 print("notify_all") self.image_created.notify_all() print("wait for someone to take it") self.image_used.wait() print("someone took it") async def image_generator(self): def waiter(): print("start waiter") time.sleep(0.1) with self.image_created: print("waiting for image_created") # 先检查状态,没有图像才等待 while not self.has_image: self.image_created.wait() self.has_image = False # 标记图像已被取用 print("waiter finished") self.image_used.set() while True: print("waiting for new image") await asyncio.to_thread(waiter) yield self.image async def main(): updater = ImageUpdater() update_thread = threading.Thread(target=updater.update_image) update_thread.start() async for image in updater.image_generator(): print(f"Received new image: {image}") if __name__ == "__main__": asyncio.run(main())
关键说明
- 使用
while not self.has_image循环检查状态,而不是直接wait(),这是Condition的标准用法——既可以防止虚假唤醒,又能处理“通知先于等待”的场景; - 每次取用图像后重置
has_image,确保下一次等待能正确响应新的图像创建通知。
内容的提问来源于stack exchange,提问作者Dronakuul
相关产品推荐
相关产品推荐

