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

call_soon_threadsafe在异步函数内无法触发回调的问题求助

异步函数嵌套调用时call_soon_threadsafe失效导致程序卡住的问题

我在用一个第三方库,它会在其他线程里随机调用我提供的函数,可建模为其他线程中的延迟函数调用。我希望该函数能调用异步函数,且不想用asyncio.run这类方式阻塞代码。此前找到的方案几乎完美,但遇到嵌套调用问题:当异步函数尝试以相同方式调用另一个异步函数时,call_soon_threadsafe不会执行回调,导致程序永久卡住。

复现代码

import asyncio
import threading
import time
from typing import Coroutine, Any
import logging

logging.basicConfig(format='%(asctime)s %(message)s', level=logging.INFO, datefmt='%H:%M:%S')

coro_queue: asyncio.Queue[Coroutine[Any, Any, Any]] = asyncio.Queue()
task: asyncio.Task[Any] | None = None
task_ready = asyncio.Event()
loop: asyncio.AbstractEventLoop | None = None

async def start_coro_queue() -> None:
    global task, loop
    loop = asyncio.get_event_loop()
    while True:
        coro = await coro_queue.get()
        task = asyncio.create_task(coro)
        task_ready.set()
        
def put_coro_in_queue(coro: Coroutine[Any, Any, Any]) -> None:
    logging.info("put_coro_in_queue called.")
    coro_queue.put_nowait(coro)
    logging.info("put_coro_in_queue finished.")
        
def run_coro_in_background(coro: Coroutine[Any, Any, Any]) -> asyncio.Task[Any]:
    logging.info("run_coro_in_background called")
    global task
    assert loop is not None
    task_ready.clear()
    future = asyncio.run_coroutine_threadsafe(task_ready.wait(), loop)
    loop.call_soon_threadsafe(put_coro_in_queue, coro)
    future.result()
    assert task is not None
    output = task
    task = None
    logging.info("run_coro_in_background finished")
    return output

async def async_func() -> None:
    logging.info("async_func called.")
    await asyncio.sleep(2)
    logging.info("async_func finished.")
    
def delayed_async_func() -> None:
    logging.info("delayed_async_func called")
    time.sleep(5)
    run_coro_in_background(async_func())
    logging.info("delayed_async_func finished.")
    
async def nested_async_func() -> None:
    logging.info("nested_async_func called")
    await asyncio.sleep(3)
    run_coro_in_background(async_func())
    logging.info("nested_async_func finished")
    
def delayed_nested_async_func() -> None:
    logging.info("delayed_nested_async_func called")
    time.sleep(4)
    run_coro_in_background(nested_async_func())
    logging.info("delayed_nested_async_func finished")

async def main() -> None:
    t = threading.Thread(target=delayed_nested_async_func)
    t.start()
    await start_coro_queue()
    
asyncio.run(main())

实际输出

19:25:51 delayed_nested_async_func called
19:25:55 run_coro_in_background called
19:25:55 put_coro_in_queue called.
19:25:55 put_coro_in_queue finished.
19:25:55 nested_async_func called
19:25:55 run_coro_in_background finished
19:25:55 delayed_nested_async_func finished
19:25:58 run_coro_in_background called

程序随后永久卡住。

期望输出

应包含以下内容:

19:25:58 put_coro_in_queue called
19:25:58 put_coro_in_queue finished
19:25:58 async_func called
19:26:00 async_func finished

非嵌套场景运行正常

修改main函数调用delayed_async_func后,输出如下:

19:34:00 delayed_async_func called
19:34:05 run_coro_in_background called
19:34:05 put_coro_in_queue called.
19:34:05 put_coro_in_queue finished.
19:34:05 async_func called.
19:34:05 run_coro_in_background finished
19:34:05 delayed_async_func finished.
19:34:07 async_func finished.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 16:22:57