You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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()  # 子进程结束后关闭管道发送端,避免资源泄漏

关键改进点

  1. 事件监听替代轮询:
    通过loop.add_reader监听管道的文件描述符,当子进程发送数据时,事件循环会立即触发回调函数,无需频繁轮询rx.poll()和调用asyncio.sleep,完全消除无效开销和等待延迟。
  2. 异步处理进程等待:
    用asyncio.to_thread把同步的process.join()放到线程中执行,避免阻塞异步事件循环,保证其他异步任务正常运行。
  3. 完善资源管理:
    在子进程的finally块中关闭管道发送端,避免资源泄漏。

方案优势

  • 零轮询开销:仅在管道有数据时触发处理,无无效循环和sleep开销;
  • 响应无延迟:数据就绪后立即处理,不会因大sleep值产生不必要等待;
  • 不阻塞事件循环:异步处理进程等待,保障整个异步程序的并发性能。

内容的提问来源于stack exchange,提问作者DarkMath

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 08:35:36