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

异步协程队列通信报错:Future绑定不同事件循环修复方案

问题修复:异步协程队列操作的循环不匹配错误

问题描述

需要实现两个异步协程:协程A向LifoQueue存入数据,协程B从队列获取并打印数据,队列为空时协程B等待直到有数据存入。但运行代码时出现got Future <Future pending> attached to a different loop错误。

错误原因

全局初始化的LifoQueue会在模块加载时绑定到当时的默认事件循环,而asyncio.run()会创建一个全新的独立事件循环。当协程在新循环中尝试使用这个队列时,队列内部的Future对象属于旧循环,导致循环不匹配的RuntimeError。

修复后的代码

import asyncio
from asyncio.queues import LifoQueue


async def get(q):
    while True:
        n = await q.get()
        print(f"get {n}")
        yield n
        print(f"yield {n}")
        # 标记任务完成,避免队列内部计数异常
        q.task_done()


async def put(q, n):
    print(f"put {n}")
    await q.put(n)


async def listen(q):
    async for i in get(q):
        print(i)


async def write(q):
    await put(q, 0)
    await asyncio.sleep(1)
    await put(q, 1)
    await asyncio.sleep(1)
    await put(q, 2)
    await asyncio.sleep(1)
    await put(q, 3)
    # 等待队列所有任务处理完成
    await q.join()


async def main():
    # 在事件循环内部创建队列,确保绑定到当前循环
    q = LifoQueue()
    # 创建任务并等待完成
    await asyncio.gather(listen(q), write(q))


asyncio.run(main())

关键修复点

  • 调整队列创建位置:将LifoQueue的创建移到main协程内部,确保队列绑定到asyncio.run()创建的事件循环上。
  • 传递队列实例:将队列作为参数传递给各个协程,避免全局变量导致的循环不匹配问题。
  • 添加任务完成标记:每次从队列获取数据后调用q.task_done(),配合q.join()让写入协程等待所有数据处理完成,保证程序正常退出。
  • 替换协程等待方式:用asyncio.gather替代asyncio.wait,更简洁地等待多个协程完成,避免额外的任务管理操作。

运行输出

put 0
get 0
0
yield 0
put 1
get 1
1
yield 1
put 2
get 2
2
yield 2
put 3
get 3
3
yield 3

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 19:10:51