Python Asyncio使用队列时同步代码阻塞问题求助
同步代码通过asyncio.Queue向异步代码发送数据时的阻塞问题
我在Python中尝试结合同步代码块与Asyncio异步代码块,让同步代码通过asyncio.Queue向异步块发送数据。不使用队列时一切运行正常,但引入队列后出现阻塞问题,尝试get_nowait等方法仍未解决。
初始代码
import asyncio import time queue = asyncio.Queue() async def processor() -> None: print("Started proc") while True: print("waiting for quee") msg = await queue.get() print(f"Got command from queue: {msg}") # do something await asyncio.sleep(5) def run_sync(url: str) -> int: while 1: print("Sending HTTP request") input("enter to send message to queue\n") queue.put_nowait(url) #do other work time.sleep(10) async def run_sync_threaded( url: str) -> int: return await asyncio.to_thread(run_sync, url) async def main() -> None: await asyncio.gather( processor(), run_sync_threaded("https://www.example.com"), ) asyncio.run(main())
临时修改后的代码
我已经实现了功能,但这更像权宜之计,感觉不够稳定。
import asyncio import time queue = asyncio.Queue() async def processor() -> None: print("Started proc") while True: print("waiting for quee") msg = await queue.get() print(f"Got command from queue: {msg}") # do something await asyncio.sleep(5) async def async_send(url): print(f'Adding {url} to queue') queue.put_nowait(url) def send(url, loop): asyncio.run_coroutine_threadsafe(async_send(url), loop) def run_sync(url: str, loop) -> int: while 1: input("enter to send message to queue\n") send(url, loop) #do other work time.sleep(3) async def run_sync_threaded( url: str, loop) -> int: return await asyncio.to_thread(run_sync, url, loop) async def main() -> None: loop = asyncio.get_event_loop() t = asyncio.create_task( processor()) t2 = asyncio.create_task(run_sync_threaded("https://www.example.com", loop)) asyncio.gather( await t, await t2 ) # This does not work # asyncio.gather( # await processor(), # await run_sync_threaded("https://www.example.com", loop) # ) asyncio.run(main())
问题根源与正确解法
初始代码问题
asyncio.Queue并非线程安全结构,直接在同步线程中调用put_nowait会破坏队列内部状态,导致异步代码阻塞。该队列设计用于异步协程间通信,跨线程操作必须通过线程安全方式与事件循环交互。
修改后代码问题
修改版用asyncio.run_coroutine_threadsafe是正确的跨线程调用异步方法的方式,但main函数中asyncio.gather(await t, await t2)写法错误——await不能直接放在gather参数中,应直接传入任务对象后await gather(...),否则会变成顺序执行而非并发。
正确实现代码
import asyncio import time queue = asyncio.Queue() async def processor() -> None: print("Started proc") while True: print("waiting for queue") msg = await queue.get() print(f"Got command from queue: {msg}") # 标记任务完成,避免队列积压(可选,按需使用) queue.task_done() await asyncio.sleep(5) def run_sync(url: str, loop) -> None: while True: input("enter to send message to queue\n") # 用run_coroutine_threadsafe安全向队列添加元素 asyncio.run_coroutine_threadsafe(queue.put(url), loop) print(f"Added {url} to queue") # 模拟同步工作 time.sleep(3) async def main() -> None: loop = asyncio.get_running_loop() # 创建异步任务 processor_task = asyncio.create_task(processor()) # 将同步函数放到线程中运行 sync_task = asyncio.to_thread(run_sync, "https://www.example.com", loop) # 并发运行所有任务 await asyncio.gather(processor_task, sync_task) asyncio.run(main())
关键要点
- 线程安全操作队列:跨线程时必须通过
asyncio.run_coroutine_threadsafe调用队列的put方法(或其他异步方法),禁止直接调用put_nowait。 - 正确使用asyncio.gather:直接传递任务对象给
gather,再await gather,才能实现并发执行。 - 可选的task_done:若需跟踪队列任务完成情况,调用
queue.task_done()可配合queue.join()使用,等待所有队列任务处理完毕。
内容的提问来源于stack exchange,提问作者will.mendil
相关产品推荐
相关产品推荐

