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
相关产品推荐
相关产品推荐

