为何loop.run_in_executor仍会阻塞事件循环?如何从同步代码触发异步任务?
嗨,我来帮你拆解这个问题~你的预期输出和实际结果不符,核心原因其实是跨线程操作事件循环的线程安全问题,以及事件循环没有被及时唤醒处理新任务。
问题根源分析
你写的c()函数是在ThreadPoolExecutor的子线程里执行的,但事件循环(loop)是跑在主线程的。当你在子线程里直接调用loop.create_task(a()),这个操作并不线程安全——虽然表面上把任务加到了队列,但此时主线程的事件循环正处于run_forever()的等待状态,它完全不知道有新任务进来了,自然不会立刻切换去执行a()。只有等子线程里的c()执行完(也就是sleep(1)结束、打印d之后),run_in_executor的future完成,事件循环才会被唤醒,这时候才会去处理之前添加的a()任务,这就导致了你看到的延迟。
解决方案:用线程安全的方式提交任务
要从同步代码(尤其是非事件循环线程的同步代码)触发异步任务,必须用事件循环提供的线程安全方法来提交任务,比如loop.call_soon_threadsafe()或者loop.call_async_threadsafe()。这类方法不仅能安全地把任务加入队列,还会立刻唤醒事件循环,让它马上处理新任务。
修改你的代码试试:
把c()里的x = loop.create_task(a())改成:
# 用线程安全的方式提交异步任务 loop.call_soon_threadsafe(asyncio.create_task, a())
如果需要获取a()任务的返回结果,可以用call_async_threadsafe:
# 获取异步任务的future对象 future_a = loop.call_async_threadsafe(asyncio.create_task, a())
调整后的完整代码如下(顺便修正了原代码的缩进问题):
import asyncio import concurrent.futures from time import sleep loop = asyncio.get_event_loop() async def a(): print('a') # 模拟IO-bound异步代码 await asyncio.sleep(0.1) print('b') def c(): print('c') # 线程安全提交异步任务 loop.call_soon_threadsafe(asyncio.create_task, a()) # 模拟CPU-bound代码 sleep(1) print('d') executor = concurrent.futures.ThreadPoolExecutor(max_workers=1) future = loop.run_in_executor(executor, c) loop.run_forever()
运行后输出就会和你的预期一致:
c a [等待0.1秒] b [等待0.9秒] d
额外提醒
如果你的同步代码本身就在事件循环所在的线程里,直接用asyncio.create_task()完全没问题;但只要是跨线程操作事件循环,一定要用线程安全的方法,不然很容易出现各种难以排查的同步问题。
备注:内容来源于stack exchange,提问作者iaalm

