Python异步开发中实现When式循环等待线程状态的方案咨询
问题核心痛点与优化方案
原有写法的问题
你当前的写法最大的问题是忙等待(busy waiting):空转的while pass循环会占满单个CPU核心,造成不必要的资源浪费。
核心优化思路
Python标准库的threading.Event是专门用于线程间事件通知的工具,等待事件触发时会进入系统级阻塞状态,完全不占用CPU资源,完美适配你的需求。我们可以基于Event实现你理想中的when condition:语法。
第一步:封装你要的when语法
用上下文管理器实现接近你预期的调用方式:
import threading from typing import Callable, Any class when: """等待事件触发后执行代码块的上下文管理器""" def __init__(self, event: threading.Event, timeout: float | None = None): self.event = event self.timeout = timeout def __enter__(self): self.event.wait(self.timeout) return self def __exit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None: pass
使用方式完全符合你的预期,仅多了一个with关键字:
with when(class_containing_threads.stopped): do_something() do_something_else() do_some_other_thing()
第二步:改造原有线程管理类,替换布尔标记为Event
把原来的_stopping、_stopped布尔变量全部替换为threading.Event,同时删除所有忙等待循环,改用内置阻塞方法:
import threading import random import time class ClassContainingThreads: def __init__(self): self._coordinating_thread = None self._workers: list[threading.Thread] = [] # 用Event替换布尔标记 self._stopping = threading.Event() self.stopped = threading.Event() self.all_workers_running = threading.Event() def run(self): self._coordinating_thread = threading.Thread(target=self._run) self._coordinating_thread.start() def stop(self): self._stopping.set() def _run(self): # 批量创建工作线程,避免重复代码 for i in range(3): worker = threading.Thread(target=self._important_function, args=(i+1,)) self._workers.append(worker) worker.start() # 标记所有工作线程已启动 self.all_workers_running.set() # 等待停止信号,阻塞不占CPU self._stopping.wait() # 批量等待所有工作线程退出,不用自己写循环判断 for worker in self._workers: worker.join() # 标记所有线程已停止 self.stopped.set() def _important_function(self, thread_id: int): print(f'Thread {thread_id} started') # 等待停止信号,阻塞不占CPU self._stopping.wait() print(f'Thread {thread_id} received stop signal') # 直接用sleep替代手动写的延迟循环 delay = random.random() * 5 time.sleep(delay) print(f'Thread {thread_id} stopped')
第三步:优化主逻辑,删除所有忙等待
if __name__ == '__main__': print('Program started') class_containing_threads = ClassContainingThreads() class_containing_threads.run() # 等待所有工作线程启动 with when(class_containing_threads.all_workers_running): while input("Type 'exit' to exit > ") != "exit": pass class_containing_threads.stop() # 等待所有线程停止后执行后续逻辑 with when(class_containing_threads.stopped): print('Program stopped')
额外说明(异步场景适配)
如果你写的是asyncio异步代码,只需要把threading.Event替换为asyncio.Event,将when改为异步上下文管理器即可,逻辑完全一致:
import asyncio from typing import Any class async_when: def __init__(self, event: asyncio.Event, timeout: float | None = None): self.event = event self.timeout = timeout async def __aenter__(self): await asyncio.wait_for(self.event.wait(), self.timeout) return self async def __aexit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None: pass
使用时搭配async with关键字即可。
内容的提问来源于stack exchange,提问作者tionichm
相关产品推荐
相关产品推荐

