如何将db_pipe的处理结果返回给对应的exec_in_db协程?
实现协程间的请求结果返回
你可以通过给每个请求附加一个Future对象来实现结果的返回,Future是asyncio中用于异步结果传递的核心组件。下面是修改后的完整代码:
import asyncio db_queue = asyncio.Queue() async def db_pipe(): while True: query, future = await db_queue.get() print("DB got", query) # 模拟数据库处理逻辑,替换为实际数据库操作即可 processed_result = f"Result for: {query}" # 将处理结果返回给对应的exec_in_db协程 future.set_result(processed_result) # 标记队列任务完成(可选但推荐,用于后续队列join操作) db_queue.task_done() async def exec_in_db(query, timeout): await asyncio.sleep(timeout) # 创建Future对象用于接收结果 result_future = asyncio.Future() # 将查询和Future绑定后放入队列 await db_queue.put((query, result_future)) # 等待Future返回处理结果 result = await result_future print(f"Got result for '{query}': {result}") return result async def main(): asyncio.create_task(db_pipe()) # 收集两个协程的执行结果 results = await asyncio.gather(exec_in_db("Long query", 4), exec_in_db("Fast query", 1)) print("All queries completed, results:", results) if __name__ == "__main__": asyncio.run(main())
关键逻辑说明:
- Future对象的核心作用:每个
exec_in_db协程创建专属的Future,与查询一起传入队列。db_pipe处理完查询后,通过future.set_result()将结果写入Future,对应的exec_in_db通过await result_future即可获取结果。 - 队列任务标记:
db_queue.task_done()是良好实践,若后续需要用await db_queue.join()等待所有队列任务完成,这个标记是必需的。 - 异常处理扩展(可选):如果数据库操作可能抛出异常,可在
db_pipe中捕获异常后调用future.set_exception(exc),这样exec_in_db在await时会抛出对应异常,便于统一处理错误。
内容的提问来源于stack exchange,提问作者Gulaev Valentin
相关产品推荐
相关产品推荐

