You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何设置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)

关键改动说明

  1. 在submit方法中,通过loop.call_threadsafe在后台事件循环线程内创建原生asyncio.Task(本质是Future的子类),确保拿到的是协程对应的真实任务对象。
  2. 新增is_coroutine_running方法,通过call_threadsafe线程安全地调用原生Future的running()方法,避免跨线程操作asyncio对象引发的问题。
  3. 主程序中通过这个安全方法查询状态,就能正确获取协程是否在运行。

内容的提问来源于stack exchange,提问作者Lionnel104

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 01:32:18