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

如何在异步向量化函数中处理OpenAI速率限制错误并实现重试?

解决OpenAI速率限制错误的异步重试实现方案

问题需求

需要在异步向量化函数中处理OpenAI的速率限制错误(429),触发错误时自动重试,同时保留错误发生前已生成的部分结果,询问当前实现的问题及更优方案。

当前代码的问题分析

你当前的代码已经实现了基础的429错误重试逻辑,但存在几个可优化点:

  • 固定等待60秒不够灵活,未利用OpenAI返回的Retry-After头信息,也没有指数退避策略,可能导致不必要的等待或频繁重试
  • 异常捕获范围不够精准:如果使用的是OpenAI官方SDK,速率限制会抛出RateLimitError而非通用的HTTPException,当前捕获逻辑可能无法命中正确异常
  • 没有设置最大重试次数,极端情况下可能陷入无限循环

更优处理方案

推荐采用以下策略提升重试逻辑的健壮性:

  • 指数退避重试:每次失败后等待时间翻倍(如10s→20s→40s…),直到设定的最大等待时间,平衡重试效率与服务压力
  • 动态读取Retry-After头:优先使用响应中返回的推荐等待时间,比固定值更精准
  • 限制最大重试次数:避免无限循环,超过次数后抛出错误
  • 精准捕获异常:针对OpenAI的RateLimitError做专门处理,其他异常正常抛出

优化后的代码示例

import asyncio
import time
from openai import RateLimitError  # 导入OpenAI官方异常类
import logging

CHUNK_SIZE = 512
MAX_RETRIES = 5  # 最大重试次数
INITIAL_WAIT = 10  # 初始等待时间(秒)

@app.get("/vectorize-data", status_code=200)
async def vectorize_data():
    text_splitter = CharacterTextSplitter(chunk_size=CHUNK_SIZE, chunk_overlap=0)
    raw_data = db.load_data_documentation()
    raw_data2 = db.load_data2_documentation()

    docs_data = text_splitter.split_documents(vectorizer.create_multiple_data_documents(raw_data))
    docs_data2 = text_splitter.split_documents(vectorizer.create_multiple_data2_documents(raw_data2))

    documents = [*docs_data, *docs_data2]
    start_time = time.time()

    # 按批次处理文档
    for i in range(0, len(documents), CHUNK_SIZE):
        retry_count = 0
        current_wait = INITIAL_WAIT
        while retry_count < MAX_RETRIES:
            try:
                Milvus.from_documents(
                    documents[i:i + CHUNK_SIZE],
                    embeddings,
                    collection_name=VECTOR_DB_NAME,
                    connection_args=CONNECTION_ARGS,
                )
                break  # 成功则跳出重试循环
            except RateLimitError as e:
                retry_count += 1
                logger.warning(f"速率限制触发,第{retry_count}次重试,当前等待{current_wait}秒")
                
                # 尝试从错误信息中提取Retry-After时间
                retry_after = None
                if hasattr(e, 'response') and e.response.headers.get('Retry-After'):
                    retry_after = int(e.response.headers['Retry-After'])
                
                # 使用Retry-After或指数退避时间
                wait_time = retry_after if retry_after else current_wait
                await asyncio.sleep(wait_time)
                
                # 更新下一次等待时间(指数退避)
                current_wait = min(current_wait * 2, 60)  # 最大等待不超过60秒
            except Exception as e:
                logger.exception(f"处理文档批次{i//CHUNK_SIZE}时发生错误: {str(e)}")
                raise  # 非速率限制错误直接抛出

        if retry_count >= MAX_RETRIES:
            raise RuntimeError(f"文档批次{i//CHUNK_SIZE}重试{MAX_RETRIES}次后仍失败")

    elapsed_time = time.time() - start_time
    return {
        "total_time_s": elapsed_time,
    }

更低层级重试的实现思路

如果希望在**向量化生成(embeddings)**的更低层级实现重试,可以直接封装embeddings的调用逻辑,例如:

async def get_embeddings_with_retry(texts, embeddings_model, max_retries=5):
    retry_count = 0
    current_wait = 10
    while retry_count < max_retries:
        try:
            # 同步模型需转换为异步调用,或使用asyncio.to_thread
            return await asyncio.to_thread(embeddings_model.embed_documents, texts)
        except RateLimitError as e:
            retry_count += 1
            retry_after = int(e.response.headers.get('Retry-After', current_wait))
            await asyncio.sleep(retry_after)
            current_wait = min(current_wait * 2, 60)
    raise RuntimeError("生成embeddings重试多次失败")

之后在批量处理文档时,先调用这个函数生成embeddings,再插入Milvus,这样单个embedding生成失败时只重试该部分,能保留更多已完成的结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 00:27:06