如何将阻塞函数转为可在无限循环中运行的非阻塞函数?
问题描述
我有一个运行在同步无限循环中的程序,需要调用外部库的阻塞函数且无法修改该函数。我希望将该函数改造为非阻塞形式,让无限循环不被阻塞,持续打印变量n,仅在get_data()执行完成后打印data。
示例代码
基础示例
import time SECONDS_TO_WAIT = 3 # 模拟外部库的阻塞函数 def get_data(): time.sleep(SECONDS_TO_WAIT) return "Data received." def main(): n = 0 while True: data = get_data() print(n) print(data) n += 1 if __name__ == "__main__": main()
贴近实际场景的示例
import time class LibraryWrapper(): SECONDS_TO_WAIT = 3 def __init__(self): # 库初始化 self.data = None # 模拟外部库的阻塞函数 def get_data(self): time.sleep(self.SECONDS_TO_WAIT) return "Data received." class Program(): def __init__(self, library_wrapper): self.library_wrapper = library_wrapper def run(self, user_input): # 执行一些操作 print(user_input) if library_wrapper.data: print(library_wrapper.data) def main(): library_wrapper = LibraryWrapper() program = Program(library_wrapper) user_input = 0 while True: # 模拟非阻塞用户输入读取 user_input += 1 library_wrapper.data = library_wrapper.get_data() program.run(user_input) if __name__ == "__main__": main()
尝试的asyncio代码(仍阻塞)
import asyncio SECONDS_TO_WAIT = 3 async def get_data(): await asyncio.sleep(3) return "Data received." async def main(): n = 0 while True: task = asyncio.create_task(get_data()) print(n) print(await task) n += 1 asyncio.run(main()) # 输出: # 0 # 3秒后... # Data received. # 1 # 3秒后... # Data received. # 2 # 3秒后... # Data received. # 3
我认为asyncio库可能是最佳解决方案,已了解协程、async/await、任务和Future的基础概念,但尝试的代码仍会阻塞循环。也查阅了相关问题,但仍不清楚如何封装阻塞函数为非阻塞形式。
解决方案
核心思路是用asyncio.run_in_executor(或Python3.9+的asyncio.to_thread())将阻塞函数放到线程池中执行,避免阻塞asyncio事件循环。这样事件循环可以继续处理其他逻辑(比如持续打印n),同时等待阻塞函数的结果。
针对基础示例的改造
import asyncio import time SECONDS_TO_WAIT = 3 # 不可修改的外部阻塞函数 def get_data(): time.sleep(SECONDS_TO_WAIT) return "Data received." async def main(): n = 0 # 启动第一次数据获取任务,不立即等待结果 data_task = asyncio.create_task(asyncio.to_thread(get_data)) while True: print(n) n += 1 # 检查任务是否完成 if data_task.done(): try: data = data_task.result() print(data) # 完成后启动下一次数据获取 data_task = asyncio.create_task(asyncio.to_thread(get_data)) except Exception as e: print(f"获取数据出错: {e}") # 出错也重新启动任务 data_task = asyncio.create_task(asyncio.to_thread(get_data)) # 让出CPU时间,避免无限循环占用100%资源 await asyncio.sleep(0.1) asyncio.run(main())
针对实际场景示例的改造
import asyncio import time class LibraryWrapper(): SECONDS_TO_WAIT = 3 def __init__(self): self.data = None self.data_task = None # 不可修改的外部阻塞函数 def get_data(self): time.sleep(self.SECONDS_TO_WAIT) return "Data received." # 后台持续非阻塞获取数据 async def start_continuous_fetch(self): while True: # 用线程执行阻塞函数 self.data = await asyncio.to_thread(self.get_data) # 可以在这里添加获取间隔,或者直接连续获取 class Program(): def __init__(self, library_wrapper): self.library_wrapper = library_wrapper def run(self, user_input): print(user_input) if self.library_wrapper.data: print(self.library_wrapper.data) # 打印后清空数据,避免重复输出 self.library_wrapper.data = None async def main(): library_wrapper = LibraryWrapper() # 启动后台数据获取任务 asyncio.create_task(library_wrapper.start_continuous_fetch()) program = Program(library_wrapper) user_input = 0 while True: user_input += 1 program.run(user_input) # 让出CPU时间 await asyncio.sleep(0.1) asyncio.run(main())
关键说明
asyncio.to_thread()是Python3.9+的便捷方法,自动将阻塞函数放入线程池执行,替代手动创建ThreadPoolExecutor的繁琐操作。- 不要直接
await阻塞任务,而是通过done()方法检查任务状态,这样事件循环可以继续执行其他逻辑。 - 添加
await asyncio.sleep(0.1)是为了避免无限循环占用100%CPU,实际场景中可根据需求调整间隔。 - 任务完成后需重新启动下一次获取,保证数据的持续更新。
内容的提问来源于stack exchange,提问作者someApp
相关产品推荐
相关产品推荐

