如何在异步向量化函数中处理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
相关产品推荐
相关产品推荐

