异步重试装饰器适配数据库耗时查询的DataFrame await错误排查
问题根因
你遇到的TypeError核心原因是异步装饰器错误地尝试await非可等待对象(DataFrame)。要么是装饰器逻辑对返回值执行了不必要的await,要么是被装饰的耗时代码本身是同步函数,却被异步装饰器直接处理了。
修复后的异步重试+超时装饰器
以下是修正后的装饰器代码,同时支持超时控制和重试逻辑:
import asyncio from functools import wraps def async_retry_with_timeout(max_retries: int, timeout: float): def decorator(func): @wraps(func) async def wrapper(*args, **kwargs): retry_count = 0 while retry_count <= max_retries: try: # 用asyncio.wait_for统一管控超时,仅await真正的可等待对象 result = await asyncio.wait_for(func(*args, **kwargs), timeout=timeout) return result except asyncio.TimeoutError: retry_count += 1 if retry_count > max_retries: raise TimeoutError(f"已达最大重试次数{max_retries},请求超时") except Exception as e: # 捕获数据库连接、查询异常等,触发重试 retry_count += 1 if retry_count > max_retries: raise e return wrapper return decorator
数据库操作类适配方案
确保数据库查询逻辑符合异步装饰器的要求:如果是异步驱动直接调用则直接使用,如果是同步查询(比如pandas的read_sql),需要用asyncio.to_thread包装为异步操作:
import pandas as pd import asyncio from your_db_module import get_async_db_conn, get_sync_db_conn class DBHandler: def __init__(self, db_config): self.db_config = db_config # 异步简单查询(直接用异步数据库驱动) @async_retry_with_timeout(max_retries=3, timeout=5) async def simple_query(self, sql): async with await get_async_db_conn(self.db_config) as conn: result = await conn.fetch(sql) return pd.DataFrame(result) # 耗时查询(同步逻辑转异步) @async_retry_with_timeout(max_retries=2, timeout=30) async def long_running_query(self, sql): # 用asyncio.to_thread将同步IO操作转为可await的异步任务 def sync_query(): with get_sync_db_conn(self.db_config) as conn: return pd.read_sql(sql, conn) return await asyncio.to_thread(sync_query)
关键修复要点
- 用
asyncio.wait_for替代手动超时判断,确保只对可等待对象执行await操作。 - 同步耗时逻辑必须通过
asyncio.to_thread包装,避免直接awaitDataFrame这类非可等待对象。 - 保留重试逻辑:超时或业务异常时自动重试,达到最大次数后抛出原异常,不吞报错。
装饰器扩展到其他耗时代码
这个装饰器可以直接复用在各类异步任务,或同步转异步的任务上:
- 异步HTTP请求:
import aiohttp @async_retry_with_timeout(max_retries=3, timeout=10) async def fetch_api_data(url): async with aiohttp.ClientSession() as session: async with session.get(url) as resp: return await resp.json()
- 同步CPU/IO密集任务:
import time def sync_heavy_calculation(): time.sleep(12) return "计算完成" @async_retry_with_timeout(max_retries=2, timeout=15) async def async_heavy_calculation(): return await asyncio.to_thread(sync_heavy_calculation)
内容的提问来源于stack exchange,提问作者Charlie
相关产品推荐
相关产品推荐

