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

使用asyncio.gather重试时遇‘cannot reuse already awaited coroutine’错误的解决咨询

解决asyncio.gather重试时"cannot reuse already awaited coroutine"错误

错误原因

协程对象是一次性的,一旦被await过就不能再次使用。你的代码里tasks列表存储的是已经创建好的协程对象,第一次await asyncio.gather(*chunk)后,这些协程就处于已完成状态,重试时再调用同一个对象就会触发错误。

修复方案

不要直接存储已创建的协程对象,而是存储能生成协程的逻辑(比如函数+参数),每次重试时重新创建新的协程对象。

修改后的代码示例

假设你原来的tasks是通过类似fetch_data(id)这样的函数生成的,先把生成协程的逻辑存起来:

# 替换原来的tasks生成逻辑,比如原来的tasks = [fetch_data(id) for id in ids]
# 改成存储函数、参数、关键字参数的元组
task_factories = [(fetch_data, (id,), {}) for id in ids]  # 根据你的实际函数调整参数

# 然后修改重试循环
for i in range((len(task_factories) // TASK_CHUNK_SIZE) + 1):
    while True:
        try:
            # 每次重试都重新生成当前chunk的协程对象
            chunk_factories = task_factories[TASK_CHUNK_SIZE*i : TASK_CHUNK_SIZE*(i+1)]
            current_tasks = [func(*args, **kwargs) for func, args, kwargs in chunk_factories]
            
            res = await asyncio.gather(*current_tasks)
            for rows in res:
                f.writelines([f'{row}\n' for row in rows])
            await asyncio.sleep(1)
            break
        except Exception as e:
            print(f'{(i + 1) * TASK_CHUNK_SIZE} -- RETRY {e}')
            await asyncio.sleep(10)
            continue

简化版(如果协程生成逻辑简单)

如果你的协程是基于某个固定参数列表生成的,也可以在重试时直接重新生成协程,不用额外存储工厂元组:

# 假设你有原始任务参数列表,比如ids = [1,2,3,...]
for i in range((len(ids) // TASK_CHUNK_SIZE) + 1):
    while True:
        try:
            # 每次重试都重新生成协程
            chunk_ids = ids[TASK_CHUNK_SIZE*i : TASK_CHUNK_SIZE*(i+1)]
            current_tasks = [fetch_data(id) for id in chunk_ids]
            
            res = await asyncio.gather(*current_tasks)
            for rows in res:
                f.writelines([f'{row}\n' for row in rows])
            await asyncio.sleep(1)
            break
        except Exception as e:
            print(f'{(i + 1) * TASK_CHUNK_SIZE} -- RETRY {e}')
            await asyncio.sleep(10)
            continue

补充说明

如果需要让asyncio.gather在单个任务出错时不直接抛出异常,而是返回异常对象以便单独处理,可以添加return_exceptions=True参数,这样就不用整个chunk重试:

res = await asyncio.gather(*current_tasks, return_exceptions=True)
for result in res:
    if isinstance(result, Exception):
        print(f"单个任务失败: {result}")
    else:
        f.writelines([f'{row}\n' for row in result])

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 17:00:11