为何用asyncio.create_task启动的协程在time.sleep后无法运行?
关于asyncio协程与time.sleep的问题及多线程事件循环实现
一、time.sleep导致协程无法运行的原因
- asyncio的事件循环采用单线程调度模型,协程的切换、执行完全依赖事件循环获得CPU时间片,事件循环必须持续运行才能调度任务。
time.sleep(0.1)是同步阻塞调用,会让当前线程(也就是事件循环所在的线程)被强制挂起0.1秒,这段时间内事件循环无法执行任何调度操作,包括启动你创建的my_coroutine任务。- 移除
time.sleep后,main函数快速执行完毕,事件循环得以接管线程,调度已创建的协程任务,因此协程能正常运行。
修正方案(main保持非异步)
如果需要等待某个条件且不能阻塞事件循环,可以将阻塞的等待逻辑放到线程池中执行,避免占用事件循环线程:
import asyncio import time from concurrent.futures import ThreadPoolExecutor async def my_coroutine(): print("协程开始执行") await asyncio.sleep(1) print("协程执行结束") def not_ready(): # 模拟未就绪状态,3秒后变为就绪 return time.time() - start_time < 3 def main(): global start_time start_time = time.time() loop = asyncio.get_event_loop() # 创建协程任务 task = loop.create_task(my_coroutine()) # 用线程池处理阻塞等待,不占用事件循环线程 with ThreadPoolExecutor() as executor: # 循环等待直到就绪 while not_ready(): loop.run_in_executor(executor, time.sleep, 0.1) # 等待协程任务完成 loop.run_until_complete(task) main()
二、事件循环在多线程中的实现
事件循环可以在多线程中运行,但要注意asyncio事件循环不是线程安全的,不能直接在非事件循环线程调用其方法,以下是两种常见场景:
1. 在子线程中独立运行事件循环
将事件循环放到单独的线程启动,主线程可执行其他逻辑:
import asyncio import threading import time async def background_coro(): while True: print("子线程事件循环:协程运行中") await asyncio.sleep(1) def run_event_loop(): # 子线程中创建并设置新的事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) # 运行事件循环直到协程结束(这里是无限循环) loop.run_until_complete(background_coro()) def main(): # 创建并启动子线程 thread = threading.Thread(target=run_event_loop, daemon=True) thread.start() # 主线程执行阻塞操作 for _ in range(5): print("主线程:执行中") time.sleep(1) main()
2. 从其他线程向主线程事件循环提交任务
如果事件循环在主线程运行,需要从子线程提交任务时,必须用线程安全的方法asyncio.run_coroutine_threadsafe或loop.call_soon_threadsafe:
import asyncio import threading import time async def handle_task(msg): print(f"事件循环收到任务:{msg}") await asyncio.sleep(1) return f"任务处理完成:{msg}" def submit_task_from_thread(loop): # 线程安全地向事件循环提交协程任务 future = asyncio.run_coroutine_threadsafe(handle_task("来自子线程的请求"), loop) # 获取任务结果(会阻塞当前线程直到结果返回) result = future.result() print(result) def main(): loop = asyncio.get_event_loop() # 启动子线程提交任务 thread = threading.Thread(target=submit_task_from_thread, args=(loop,)) thread.start() # 主线程运行事件循环 loop.run_forever() main()
内容的提问来源于stack exchange,提问作者caeus
相关产品推荐
相关产品推荐

