如何将依赖同步可迭代对象的第三方同步生成器转为异步生成器?
异步上下文使用同步第三方生成器的解决方案
有一个无法修改的第三方同步生成器函数my_gen,仅接受同步可迭代对象作为输入。需要在异步上下文环境中传入异步可迭代对象使用它,且不能阻塞事件循环。
同步使用示例
sync_input_iterable = range(0, 10000) sync_output_iterable = my_gen(input_iterable) for v in sync_output_iterable: print(v)
期望的异步使用框架
async def main(): # 示例异步可迭代生成器 async def arange(start, end): for i in range(start, end): yield(i) await asyncio.sleep(0) async_input_iterable = arange(0, 10000) async_output_iterable = # 此处需要实现适配逻辑 async for v in async_output_iterable: print(v) asyncio.run(main())
解决方案
核心思路是用线程隔离同步生成器的运行,通过队列实现异步与同步的数据交互:
- 将异步可迭代对象的元素通过异步队列传递到后台线程,转换成同步可迭代对象供
my_gen使用 my_gen在后台线程运行,产出的结果再通过队列传回异步上下文,包装成异步生成器供async for遍历
完整适配实现
import asyncio from concurrent.futures import ThreadPoolExecutor def async_gen_adapter(sync_gen_func, async_iterable): # 输入队列:异步→同步;输出队列:同步→异步 input_queue = asyncio.Queue(maxsize=1) output_queue = asyncio.Queue(maxsize=1) # 同步迭代器:从异步队列取元素,供同步生成器使用 def sync_iterable(): while True: # 线程安全地获取异步队列元素 get_future = asyncio.run_coroutine_threadsafe(input_queue.get(), asyncio.get_running_loop()) try: item = get_future.result() if item is None: # 终止信号 break yield item except Exception: break finally: input_queue.task_done() # 后台线程任务:运行同步生成器,将结果放入输出队列 def run_sync_gen(): try: for output_item in sync_gen_func(sync_iterable()): put_future = asyncio.run_coroutine_threadsafe(output_queue.put(output_item), asyncio.get_running_loop()) put_future.result() except Exception as e: # 将同步侧异常传递到异步侧 asyncio.run_coroutine_threadsafe(output_queue.put(Exception(f"同步生成器错误: {e}")), asyncio.get_running_loop()).result() finally: # 发送结束信号 asyncio.run_coroutine_threadsafe(output_queue.put(None), asyncio.get_running_loop()).result() # 异步任务:将异步可迭代对象的元素喂</think_never_used_51bce0c785ca2f68081bfa7d91973934>输入队列 async def feed_input(): try: async for item in async_iterable: await input_queue.put(item) except Exception as e: # 将异步侧异常传递到同步侧 await input_queue.put(Exception(f"异步迭代器错误: {e}")) finally: await input_queue.put(None) # 启动线程和异步任务 executor = ThreadPoolExecutor(max_workers=1) executor.submit(run_sync_gen) asyncio.create_task(feed_input()) # 异步产出结果 while True: result = await output_queue.get() if result is None: break if isinstance(result, Exception): raise result yield result output_queue.task_done() executor.shutdown(wait=True)
使用方式
在main函数中直接调用适配器:
async_output_iterable = async_gen_adapter(my_gen, async_input_iterable)
关键说明
- 线程隔离:同步生成器在单独线程中运行,完全不会阻塞异步事件循环
- 线程安全交互:通过
asyncio.run_coroutine_threadsafe实现异步队列在不同线程中的安全操作 - 终止与异常处理:用
None作为迭代终止信号,同时捕获并传递异步/同步两侧的异常,保证逻辑完整性 - 兼容性:适配Python 3.7+,若使用Python 3.9+,可将
ThreadPoolExecutor替换为asyncio.to_thread简化线程管理
内容的提问来源于stack exchange,提问作者Michal Charemza
相关产品推荐
相关产品推荐

