如何将基于轮询的异步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
关键注意事项
- 取消逻辑绑定:必须在封装层同步调用原任务的
cancel(),确保异步任务被取消时,远程任务也能被终止。 - 轮询间隔:根据业务场景调整睡眠时长,避免过短占用CPU,或过长导致响应延迟。
- 异常传递:原
get_result()抛出的异常会自动传递到异步调用层,无需额外处理。 - 线程安全:若原
RemoteTask的方法非线程安全,需在轮询时加锁,避免异步调度引发的并发问题。
内容的提问来源于stack exchange,提问作者shadowtalker
相关产品推荐
相关产品推荐

