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

Python Asyncio:如何实现重试操作的并行等待

解决异步调用重试串行问题的方案

嘿,这个问题我之前也碰到过!你的核心问题出在两个关键点上:

  • 重试装饰器里用了同步的time.sleep(),这会直接阻塞整个asyncio事件循环,让其他需要重试的协程根本没机会并行运行;
  • 当前的重试逻辑虽然写在异步函数里,但同步sleep会把事件循环卡死,导致所有失败任务的重试操作只能一个接一个来。

怎么改?看这两步:

1. 把同步睡眠换成异步睡眠

把装饰器里的time.sleep(sleep_ms / 1000)替换成await asyncio.sleep(sleep_ms / 1000)。asyncio的sleep是异步的,协程在等待的时候会主动让出事件循环的控制权,这样其他需要重试的任务就能同时跑起来了。

2. 优化重试计数的打印(可选但更清晰)

原来的打印逻辑有点混乱(比如第一次重试会显示成2/10),我调整了计数方式,让重试次数的展示更直观。

修改后的完整代码:

先看更新后的重试装饰器:

import functools
import random
import asyncio

def retry_with_backoff(retries=5, backoff_in_ms=100):
    def wrapper(f):
        @functools.wraps(f)
        async def wrapped(*args, **kwargs):
            attempt = 0
            while True:
                try:
                    return await f(*args, **kwargs)
                except Exception as e:
                    attempt += 1
                    print(f'Fetch error on attempt {attempt}: {e}')
                    if attempt >= retries:
                        print(f'Failed after {attempt} retries, giving up')
                        raise
                    # 计算退避时间,用异步sleep让出控制权
                    sleep_ms = (backoff_in_ms * (2 ** (attempt - 1)) + random.uniform(0, 1))
                    await asyncio.sleep(sleep_ms / 1000)
                    print(f'Retrying ({attempt}/{retries})...')
        return wrapped
    return wrapper

你的业务函数和调用逻辑不用变,保持原来的就行:

@retry_with_backoff(retries=10)
async def fetch_tx_details(sig):
    # 这里放你的业务逻辑,比如调用第三方API、查询数据库等
    pass

# 并行发起调用的逻辑完全不变
txs = await asyncio.gather(*[fetch_tx_details(s["signature"]) for s in sigs])

为什么这样就能实现并行重试了?

  • 用await asyncio.sleep()后,某个协程进入重试等待时,会把事件循环的控制权交还给asyncio,其他需要重试的协程就能被调度执行,所有重试任务就会并行进行;
  • 每个fetch_tx_details的重试过程都是独立的异步任务,asyncio会自动调度它们,不会因为某个任务在等重试而卡住整个流程。

额外给你两个优化建议:

  • 如果不想因为个别任务最终失败导致整个gather调用失败,可以给asyncio.gather加上return_exceptions=True参数,这样失败的任务会返回异常对象,你可以在结果里过滤处理;
  • 如果重试的任务数量特别多,建议加个asyncio.Semaphore来限制并发数,避免把目标服务打崩。比如在fetch_tx_details里加个信号量控制:
    semaphore = asyncio.Semaphore(50)  # 限制同时50个并发
    
    @retry_with_backoff(retries=10)
    async def fetch_tx_details(sig):
        async with semaphore:
            # 你的业务逻辑
            pass
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 14:07:35