aiogram机器人集成OpenAI API的异步阻塞问题求助
问题核心
你的机器人出现阻塞、延迟、其他用户无法使用的原因是同步调用OpenAI接口阻塞了asyncio事件循环。因为聊天处理函数是异步的,但gpt_talk是同步函数,调用它会占用整个事件循环的线程,直到OpenAI返回结果,期间所有其他用户的消息、命令都无法被处理。
解决方案
1. 改用OpenAI异步客户端
OpenAI官方提供了异步SDK,替换同步接口即可避免阻塞事件循环:
import openai from openai import AsyncOpenAI # 初始化异步客户端(可通过环境变量或直接传入API_KEY) client = AsyncOpenAI(api_key="你的OpenAI_API_KEY") # 改写为异步函数 async def gpt_talk(msg): response = await client.completions.create( model="text-davinci-003", prompt=msg, temperature=0.9, max_tokens=1000, top_p=1, frequency_penalty=0.0, presence_penalty=0.6, ) return response.choices[0].text.strip() # 同理,gpt_prog也要改成异步函数,调用对应的OpenAI异步接口 async def gpt_prog(msg): # 你的代码逻辑,使用await调用异步接口 pass
2. 修正消息处理函数的调用
在异步处理函数中,必须用await调用异步函数,否则会返回未执行的协程对象导致错误:
@dp.message_handler() # Catches all messages except commands async def chatting(message: types.Message): usr_model = db.get_model(message.from_user.id)[0] msg = message.text # 补全原代码中缺失的msg变量 if usr_model == 1: text = await gpt_talk(msg) elif usr_model == 2: text = await gpt_prog(msg) await message.answer(text)
3. 可选:控制并发请求数(避免API限流)
如果担心大量用户同时请求触发OpenAI的限流,可以用asyncio.Semaphore限制同时处理的请求数:
import asyncio # 限制最多同时处理5个请求,可根据OpenAI的API配额调整 semaphore = asyncio.Semaphore(5) async def gpt_talk(msg): async with semaphore: response = await client.completions.create( # 保持原参数 ) return response.choices[0].text.strip()
关于asyncio的学习
不需要深入学习asyncio的所有细节,掌握以下基础即可解决当前问题:
- 异步函数的定义(
async def) - 使用
await等待异步任务完成 - 理解同步代码会阻塞事件循环的核心概念
- 常用异步工具:
Semaphore(控制并发)、Task(创建后台任务)
后续遇到更复杂的异步场景(比如定时任务、多任务协作)再深入学习即可。
队列实现的正确姿势(可选)
如果需要对请求进行严格排队处理,可以使用asyncio.Queue,示例如下:
import asyncio request_queue = asyncio.Queue() # 后台worker任务,负责处理队列中的请求 async def request_worker(): while True: message = await request_queue.get() try: usr_model = db.get_model(message.from_user.id)[0] msg = message.text if usr_model == 1: text = await gpt_talk(msg) elif usr_model == 2: text = await gpt_prog(msg) await message.answer(text) finally: request_queue.task_done() # 机器人启动时启动worker async def on_startup(dp): asyncio.create_task(request_worker()) # 修改消息处理函数,将请求放入队列 @dp.message_handler() async def chatting(message: types.Message): await request_queue.put(message)
内容的提问来源于stack exchange,提问作者latteisuser
相关产品推荐
相关产品推荐

