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

异步重试装饰器适配数据库耗时查询的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:52:42