如何在Python asyncio中高效等待同步条件?
问题描述
我现在需要等待一个同步条件,目前用主动轮询的方式实现:
while my_syncronous_condition_is_not_fulfilled(): asyncio.sleep(0.001)
这种方式能工作,但每次调用asyncio.sleep存在开销,性能不佳;如果把sleep值设大比如asyncio.sleep(1),又会导致不必要的等待。请问最优实现方式是什么?
完整代码
import asyncio import inspect from multiprocessing import Process, Pipe async def calculate_in_subprocess(func, *args, **kwargs): rx, tx = Pipe(duplex=False) # receiver & transmitter ; Pipe is one-way only process = Process(target=_inner, args=(tx, func, *args), kwargs=kwargs) process.start() while not rx.poll(): # do not use process.is_alive() as condition here await asyncio.sleep(0.001) result = rx.recv() process.join() # this blocks synchronously! make sure that process is terminated before you call join() rx.close() if isinstance(result, Exception): raise result return result def _inner(tx, fun, *a, **kw_args) -> None: """ This runs in another process. """ event_loop = None if inspect.iscoroutinefunction(fun): event_loop = asyncio.new_event_loop() asyncio.set_event_loop(event_loop) try: if event_loop is not None: res = event_loop.run_until_complete(fun(*a, **kw_args)) else: res = fun(*a, **kw_args) except Exception as ex: tx.send(ex) else: tx.send(res)
示例用法
import time import asyncio def f(value: int) -> int: time.sleep(10) # 耗时较长的同步阻塞计算 return 2 * value asyncio.run(calculate_in_subprocess(func=f, value=42))
最优实现方案
核心思路是用异步IO的事件监听替代轮询,彻底消除asyncio.sleep的开销和等待延迟,同时解决同步process.join()阻塞事件循环的问题。
修改后的完整代码
import asyncio import inspect from multiprocessing import Process, Pipe from functools import partial async def calculate_in_subprocess(func, *args, **kwargs): rx, tx = Pipe(duplex=False) process = Process(target=_inner, args=(tx, func, *args), kwargs=kwargs) process.start() # 用Future等待管道可读事件 ready_future = asyncio.Future() def on_pipe_ready(fut, pipe): if not fut.done(): fut.set_result(True) # 移除监听,避免重复触发 asyncio.get_running_loop().remove_reader(pipe.fileno()) # 注册管道可读事件监听 loop = asyncio.get_running_loop() loop.add_reader(rx.fileno(), partial(on_pipe_ready, ready_future, rx)) await ready_future # 等待数据就绪,无轮询无延迟 result = rx.recv() rx.close() # 异步等待进程结束,不阻塞事件循环 await asyncio.to_thread(process.join) if isinstance(result, Exception): raise result return result def _inner(tx, fun, *a, **kw_args) -> None: """ Runs in another process. """ event_loop = None if inspect.iscoroutinefunction(fun): event_loop = asyncio.new_event_loop() asyncio.set_event_loop(event_loop) try: if event_loop is not None: res = event_loop.run_until_complete(fun(*a, **kw_args)) else: res = fun(*a, **kw_args) except Exception as ex: tx.send(ex) else: tx.send(res) finally: tx.close() # 子进程结束后关闭管道发送端,避免资源泄漏
关键改进点
- 事件监听替代轮询:
通过loop.add_reader监听管道的文件描述符,当子进程发送数据时,事件循环会立即触发回调函数,无需频繁轮询rx.poll()和调用asyncio.sleep,完全消除无效开销和等待延迟。 - 异步处理进程等待:
用asyncio.to_thread把同步的process.join()放到线程中执行,避免阻塞异步事件循环,保证其他异步任务正常运行。 - 完善资源管理:
在子进程的finally块中关闭管道发送端,避免资源泄漏。
方案优势
- 零轮询开销:仅在管道有数据时触发处理,无无效循环和sleep开销;
- 响应无延迟:数据就绪后立即处理,不会因大sleep值产生不必要等待;
- 不阻塞事件循环:异步处理进程等待,保障整个异步程序的并发性能。
内容的提问来源于stack exchange,提问作者DarkMath
相关产品推荐
相关产品推荐

