如何设置concurrent.future.Future对象的running属性?
问题原因与解决方法
为什么fut.running()始终返回False?
你通过asyncio.run_coroutine_threadsafe获取的Future并非asyncio原生的asyncio.Future,而是concurrent.futures.Future实例。这个Future的running()方法判断的是跨线程提交协程的动作是否在执行,而非后台事件循环里协程的运行状态。
实际场景中,“把协程提交到后台事件循环”这个动作会瞬间完成,之后这个Future就进入等待结果的状态,不再处于running,所以你看到的一直是False,但后台的协程其实一直在运行。
如何正确获取协程的运行状态?
要追踪事件循环中协程的真实运行状态,需要拿到事件循环内的原生asyncio.Future,并通过线程安全的方式查询它的状态(因为asyncio的对象不支持跨线程直接操作)。
修改后的代码示例
import asyncio import time from threading import Thread import multiprocessing as mp class MultiProcessor: def __init__(self, n) -> None: self._loop = asyncio.new_event_loop() self._bg_thread = Thread(target=loop_forever, args=(self._loop, ), daemon=True) self._bg_thread.start() def submit(self, f, *args, **kwargs): # 在事件循环线程内创建原生asyncio Future def _create_task(): coro = run_one_process(f, args, kwargs) return asyncio.create_task(coro) # 线程安全地获取原生Future native_fut = self._loop.call_threadsafe(_create_task) # 获取用于等待结果的threadsafe Future threadsafe_fut = asyncio.run_coroutine_threadsafe(run_one_process(f, args, kwargs), self._loop) return native_fut, threadsafe_fut def is_coroutine_running(self, native_fut): # 线程安全地查询原生Future的running状态 return self._loop.call_threadsafe(lambda: native_fut.running()) async def run_one_process(f, args, kwargs): q = mp.Queue() kwargs['q'] = q p = mp.Process(target=f, args=args, kwargs=kwargs, daemon=True) p.start() while p.is_alive(): await asyncio.sleep(0.1) res = q.get() p.close() return res def test_func(a, b, **kwargs): print('start running') time.sleep(2) res = a + b kwargs['q'].put(res) def loop_forever(loop): asyncio.set_event_loop(loop) loop.run_forever() if __name__ == "__main__": pool = MultiProcessor(4) start = time.perf_counter() native_fut, fut = pool.submit(test_func, 1, 2) while not fut.done(): # 调用线程安全的方法查询状态 print(pool.is_coroutine_running(native_fut)) time.sleep(0.1)
关键改动说明
- 在
submit方法中,通过loop.call_threadsafe在后台事件循环线程内创建原生asyncio.Task(本质是Future的子类),确保拿到的是协程对应的真实任务对象。 - 新增
is_coroutine_running方法,通过call_threadsafe线程安全地调用原生Future的running()方法,避免跨线程操作asyncio对象引发的问题。 - 主程序中通过这个安全方法查询状态,就能正确获取协程是否在运行。
内容的提问来源于stack exchange,提问作者Lionnel104
相关产品推荐
相关产品推荐

