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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 11:55:46