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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 16:55:17