使用asyncio的run_coroutine_threadsafe报错CancelledError,求排查
问题:跨线程调用
run_coroutine_threadsafe抛出CancelledError 我尝试在另一个线程中调用run_coroutine_threadsafe,却出现如下报错:
Exception in thread my_thread: Traceback (most recent call last): File "/.pyenv/versions/3.9.0/lib/python3.9/threading.py", line 950, in _bootstrap_inner self.run() File "/.pyenv/versions/3.9.0/lib/python3.9/threading.py", line 888, in run self._target(*self._args, **self._kwargs) File "/sample1.py", line 17, in hello y = x.result() File "/.pyenv/versions/3.9.0/lib/python3.9/concurrent/futures/_base.py", line 438, in result raise CancelledError() concurrent.futures._base.CancelledError
对应的测试代码:
import asyncio import threading async def hel(): return 4 class Hello: def __init__(self): self.loop = asyncio.get_running_loop() self.my_thread = threading.Thread(name='my_thread', target=self.hello) self.my_thread.start() def hello(self): x = asyncio.run_coroutine_threadsafe(hel(), self.loop) y = x.result() print(y) async def h(): Hello() asyncio.run(h())
原因分析
问题出在主线程的事件循环提前退出。asyncio.run(h())会在h()协程执行完毕后立刻关闭事件循环,但此时Hello类初始化时启动的子线程还在执行run_coroutine_threadsafe提交任务的操作。事件循环关闭后,所有未完成的任务会被取消,因此调用x.result()时抛出CancelledError。
解决办法
要让主线程的事件循环等待子线程执行完毕后再关闭,这里提供一种简单的实现方式:
修改h()协程,通过asyncio.to_thread等待子线程完成,避免阻塞事件循环:
import asyncio import threading async def hel(): return 4 class Hello: def __init__(self): self.loop = asyncio.get_running_loop() self.my_thread = threading.Thread(name='my_thread', target=self.hello) self.my_thread.start() def hello(self): x = asyncio.run_coroutine_threadsafe(hel(), self.loop) y = x.result() print(y) async def h(): hello_inst = Hello() # 异步等待子线程结束,确保事件循环在子线程任务完成后再关闭 await asyncio.to_thread(hello_inst.my_thread.join) asyncio.run(h())
运行修改后的代码,就能正常输出4,不会再抛出取消错误。核心是保证run_coroutine_threadsafe提交的协程在事件循环运行期间完成执行。
内容的提问来源于stack exchange,提问作者the_tech_guy
相关产品推荐
相关产品推荐

