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

使用asyncio.run_coroutine_threadsafe()跨线程调用异步函数时程序挂起问题

问题描述

我的应用全程使用async/await异步编程,经常需要调用同步函数,而这些同步函数有时又得调用异步函数。我尝试通过子线程利用主线程的事件循环,用asyncio.run_coroutine_threadsafe()实现需求,但调用该方法返回的Future对象的.result()时,程序出现挂起。

以下是可在Python 3.12和3.13中复现问题的代码:

from asyncio import run_coroutine_threadsafe, get_running_loop, run, AbstractEventLoop
from collections.abc import Callable, Awaitable
from concurrent.futures.thread import ThreadPoolExecutor
from typing import TypeVar

_T = TypeVar("_T")


def _await_in_thread(loop: AbstractEventLoop, f: Callable[[], Awaitable[_T]]) -> _T:
    """
    Await something inside a thread, using the main thread's loop.
    """

    print("THIS IS PRINTED")
    return run_coroutine_threadsafe(f(), loop).result()


def _await_to_thread(pool: ThreadPoolExecutor, f: Callable[[], Awaitable[_T]]) -> _T:
    """
    Await something by moving it to a thread.
    """
    return pool.submit(_await_in_thread, get_running_loop(), f).result()


async def _async_main(pool: ThreadPoolExecutor) -> None:
    """
    Run the main application, which is asynchronous.
    """

    # Eventually, the application calls a function that due to its nature (maybe a third-party API) is synchronous.
    _some_sync_function(pool)


def _sync_main() -> None:
    with ThreadPoolExecutor() as pool:
        run(_async_main(pool))


def _some_sync_function(pool: ThreadPoolExecutor) -> None:
    # This synchronous function then has to call a function that is asynchronous.
    result = _await_to_thread(pool, _some_async_function)
    assert result == 123


async def _some_async_function() -> int:
    print("BUT THIS IS NOT PRINTED")
    return 123


if __name__ == "__main__":
    _sync_main()

需求说明:

  • 必须保留ThreadPoolExecutor的使用和返回值传递
  • 要求多线程共用同一事件循环
  • 代码需兼容Python 3.11+,必要时可升级至3.12+

问题根源

代码挂起的核心是主线程事件循环被阻塞,形成死锁:

  1. 异步函数_async_main直接调用同步的_some_sync_function,此时主线程事件循环被同步逻辑完全占用,无法处理新的异步任务。
  2. _some_sync_function调用_await_to_thread时,pool.submit(...).result()会阻塞主线程,等待子线程执行完成。
  3. 子线程中run_coroutine_threadsafe(f(), loop).result()需要主线程事件循环调度执行_some_async_function,但主线程正被阻塞卡死,无法处理该异步任务,最终两边互相等待,程序挂起。

正确实现方式

解决核心是保证主线程事件循环不被阻塞,同时正确处理跨线程异步调用逻辑。以下是符合需求的实现方案:

方案一:自定义线程池兼容版

from asyncio import run_coroutine_threadsafe, get_running_loop, run, AbstractEventLoop
from collections.abc import Callable, Awaitable
from concurrent.futures.thread import ThreadPoolExecutor
from typing import TypeVar
import asyncio

_T = TypeVar("_T")


def _await_in_thread(loop: AbstractEventLoop, f: Callable[[], Awaitable[_T]]) -> _T:
    print("THIS IS PRINTED")
    future = run_coroutine_threadsafe(f(), loop)
    return future.result()


async def _await_to_thread(pool: ThreadPoolExecutor, f: Callable[[], Awaitable[_T]]) -> _T:
    # 改用异步方式等待线程任务,避免阻塞事件循环
    return await asyncio.wrap_future(pool.submit(_await_in_thread, get_running_loop(), f))


async def _async_main(pool: ThreadPoolExecutor) -> None:
    # 用asyncio.to_thread将同步函数丢到线程池,主线程事件循环保持活跃
    await asyncio.to_thread(_some_sync_function, pool)


def _sync_main() -> None:
    with ThreadPoolExecutor() as pool:
        run(_async_main(pool))


def _some_sync_function(pool: ThreadPoolExecutor) -> None:
    loop = get_running_loop()
    # 此时主线程事件循环未被阻塞,能处理run_coroutine_threadsafe提交的任务
    future = run_coroutine_threadsafe(_some_async_function(), loop)
    result = future.result()
    assert result == 123
    print(f"Got result: {result}")


async def _some_async_function() -> int:
    print("NOW THIS IS PRINTED")
    return 123


if __name__ == "__main__":
    _sync_main()

方案二:简化版(利用事件循环自带线程池)

如果不需要自定义线程池,可直接用asyncio内置的线程池管理:

import asyncio

async def _some_async_function() -> int:
    print("NOW THIS IS PRINTED")
    return 123

def _some_sync_function(loop: asyncio.AbstractEventLoop) -> None:
    # 同步函数中跨线程调用异步任务
    future = asyncio.run_coroutine_threadsafe(_some_async_function(), loop)
    result = future.result()
    assert result == 123
    print(f"Result received: {result}")

async def _async_main() -> None:
    loop = asyncio.get_running_loop()
    # 将同步函数丢到线程池执行,不阻塞主线程事件循环
    await loop.run_in_executor(None, _some_sync_function, loop)

def _sync_main() -> None:
    asyncio.run(_async_main())

if __name__ == "__main__":
    _sync_main()

关键总结

  • 主线程事件循环绝对不能被同步调用阻塞:只要事件循环卡死,run_coroutine_threadsafe提交的任务永远无法被处理,必然死锁。
  • 同步函数需放到子线程执行:用asyncio.to_thread或loop.run_in_executor把同步逻辑转移到线程池,让主线程事件循环保持活跃。
  • 跨线程调用异步任务必须用run_coroutine_threadsafe:但前提是目标事件循环处于运行状态,能处理新任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 15:55:01