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
相关产品推荐
相关产品推荐

