使用asyncio调用OpenAI API长时间运行后出现挂起故障
异步调用OpenAI API挂起问题分析与解决
问题重现
使用asyncio并发调用OpenAI API翻译文本,初期正常输出结果,运行一段时间后程序无响应,既没有新结果输出,也没有重试提示。代码如下:
import asyncio from aiohttp import ClientSession import openai import os openai.api_key = os.getenv("OPENAI_API_KEY") async def _atranslate(sem, messages, **model_kwargs): max_retry = 2 async with sem: while max_retry > 0: try: response = await openai.ChatCompletion.acreate( messages=messages, **model_kwargs ) answer = response.choices[0]['message']['content'] print(answer) return answer except Exception as e: print(e) await asyncio.sleep(5) max_retry -= 1 print('retrying...') raise ConnectionError('cannot reach openai!') async def atranslate(text_list: list, source=None, target='English', max_workers=3, **model_kwargs): aio_session = ClientSession() openai.aiosession.set(aio_session) model_kwargs.setdefault('model', 'gpt-3.5-turbo') model_kwargs.setdefault('temperature', 1) model_kwargs.setdefault('timeout', 10) template = 'Translate the following {source} text into {target}:{text}' semaphore = asyncio.Semaphore(max_workers) tasks = [] for text in text_list: messages = [{ 'role': 'user', 'content': template.format( source=source, target=target, text=text )} ] tasks.append(asyncio.create_task(_atranslate(semaphore, messages, **model_kwargs))) results = await asyncio.gather(*tasks) await aio_session.close() return results if __name__ == '__main__': textList = '... (some texts are omitted for brevity)' translations = asyncio.run(atranslate(textList*20, 'Korean', 'English',30))
可能的原因
- 并发超限触发静默限制:设置的
max_workers=30远高于OpenAI推荐的gpt-3.5-turbo并发数(通常建议≤10),OpenAI可能对超限请求采取静默限流(不返回错误,仅保持连接挂起),导致任务一直等待响应。 - 连接池耗尽:未配置aiohttp的连接池参数,默认连接池无法支撑高并发请求,新请求无法获取可用连接,陷入无限等待。
- 超时机制失效:仅设置OpenAI客户端的
timeout参数,未对整个异步调用设置全局超时,若服务器端保持长连接不返回,任务会一直挂起。 - 重试逻辑覆盖不全:部分隐性异常(如连接挂起类异常)未被捕获,导致任务卡在
await openai.ChatCompletion.acreate()步骤,无法进入重试流程。
解决办法与优化建议
降低并发数
将max_workers调整到5-10之间,符合OpenAI的速率限制要求,避免触发静默限流。配置aiohttp连接池
创建ClientSession时指定连接池参数,限制最大连接数和超时时间,避免连接耗尽:from aiohttp import TCPConnector connector = TCPConnector(limit=10, limit_per_host=10, timeout=10) aio_session = ClientSession(connector=connector)添加全局超时控制
用asyncio.wait_for包裹OpenAI API调用,确保单个任务不会无限等待:# 在_atranslate的try块中修改 response = await asyncio.wait_for( openai.ChatCompletion.acreate(messages=messages, **model_kwargs), timeout=15 # 全局超时,比OpenAI的timeout稍长 )优化重试逻辑
- 采用指数退避策略(重试间隔逐渐增加),避免短时间内重复请求触发更严格的限制;
- 捕获细分异常类型(如
openai.error.RateLimitError),针对性处理; - 在重试时打印请求标识,方便定位问题。
修改后的重试逻辑示例:
async def _atranslate(sem, messages, **model_kwargs): max_retry = 3 retry_delay = 2 # 初始延迟 async with sem: for attempt in range(max_retry): try: response = await asyncio.wait_for( openai.ChatCompletion.acreate(messages=messages, **model_kwargs), timeout=15 ) answer = response.choices[0]['message']['content'] print(answer) return answer except openai.error.RateLimitError as e: print(f"Rate limit hit: {e}, retrying in {retry_delay}s...") await asyncio.sleep(retry_delay) retry_delay *= 2 # 指数退避 except Exception as e: print(f"Error: {e}, retrying in {retry_delay}s...") await asyncio.sleep(retry_delay) retry_delay *= 2 raise ConnectionError('Failed to reach OpenAI after retries!')避免全局任务挂起
在asyncio.gather中设置return_exceptions=True,单个任务失败不会导致整个程序挂起,而是返回异常对象,便于后续排查:results = await asyncio.gather(*tasks, return_exceptions=True)升级OpenAI库
确保使用最新版本的OpenAI Python库,修复旧版本中异步客户端的潜在bug:pip install --upgrade openai
内容的提问来源于stack exchange,提问作者martin li
相关产品推荐
相关产品推荐

