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

如何将基于轮询的异步API封装为Python Awaitable接口?

如何将同步RemoteTask封装为Awaitable接口

核心思路是利用异步库的非阻塞睡眠替代同步sleep,在异步轮询中等待任务完成,同时保留原任务的取消、结果获取逻辑,且完全不修改原RemoteTask类。以下分通用方案和各主流异步库的具体实现:

通用封装思路(适配多数异步库)

实现一个符合Python Awaitable协议的类,内部通过异步睡眠轮询任务状态,同时绑定原任务的取消逻辑:

class AsyncRemoteTask:
    def __init__(self, remote_task):
        self.remote_task = remote_task
        self.remote_task.start()  # 启动远程任务

    def __await__(self):
        # 需替换为对应异步库的sleep方法,比如anyio.sleep/asyncio.sleep
        from anyio import sleep
        while not self.remote_task.is_done():
            await sleep(2)
        return self.remote_task.get_result()

    def cancel(self):
        self.remote_task.cancel()

用法:

async def foo():
    task = AsyncRemoteTask(RemoteTask())
    return await task

针对Asyncio的实现

方案1:自定义可await类

直接基于asyncio的异步睡眠实现,逻辑简洁:

import asyncio

class AsyncioRemoteTask:
    def __init__(self, remote_task):
        self.remote_task = remote_task
        self.remote_task.start()
        self._cancelled = False

    async def _wait_for_completion(self):
        while not self.remote_task.is_done() and not self._cancelled:
            await asyncio.sleep(2)
        if self._cancelled:
            raise asyncio.CancelledError
        return self.remote_task.get_result()

    def __await__(self):
        return self._wait_for_completion().__await__()

    def cancel(self):
        self._cancelled = True
        if not self.remote_task.is_done():
            self.remote_task.cancel()

方案2:包装为asyncio.Future

贴合asyncio生态,可兼容asyncio的所有调度工具:

import asyncio

def wrap_to_asyncio_future(remote_task):
    loop = asyncio.get_running_loop()
    future = loop.create_future()

    def _check_task_status():
        if remote_task.is_done():
            try:
                result = remote_task.get_result()
                loop.call_soon_threadsafe(future.set_result, result)
            except Exception as e:
                loop.call_soon_threadsafe(future.set_exception, e)
            return
        # 未完成则继续调度检查
        loop.call_later(2, _check_task_status)

    remote_task.start()
    loop.call_soon(_check_task_status)

    # 绑定异步取消逻辑
    def _on_future_cancel(fut):
        if not remote_task.is_done():
            remote_task.cancel()
    future.add_done_callback(_on_future_cancel)

    return future

用法:

async def foo():
    remote_task = RemoteTask()
    async_future = wrap_to_asyncio_future(remote_task)
    return await async_future

针对Trio的实现

结合Trio的取消作用域,实现更贴合Trio风格的封装:

import trio

async def run_trio_remote_task(remote_task):
    remote_task.start()
    try:
        while not remote_task.is_done():
            await trio.sleep(2)
        return remote_task.get_result()
    except trio.Cancelled:
        if not remote_task.is_done():
            remote_task.cancel()
        raise

用法:

async def foo():
    task = RemoteTask()
    return await run_trio_remote_task(task)

针对AnyIO的实现

AnyIO支持多后端(asyncio/Trio),可直接写通用异步函数:

import anyio

async def run_anyio_remote_task(remote_task):
    remote_task.start()
    try:
        while not remote_task.is_done():
            await anyio.sleep(2)
        return remote_task.get_result()
    except anyio.Cancelled:
        if not remote_task.is_done():
            remote_task.cancel()
        raise

关键注意事项

  1. 取消逻辑绑定:必须在封装层同步调用原任务的cancel(),确保异步任务被取消时,远程任务也能被终止。
  2. 轮询间隔:根据业务场景调整睡眠时长,避免过短占用CPU,或过长导致响应延迟。
  3. 异常传递:原get_result()抛出的异常会自动传递到异步调用层,无需额外处理。
  4. 线程安全:若原RemoteTask的方法非线程安全,需在轮询时加锁,避免异步调度引发的并发问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 22:15:26