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

FastAPI WebSocket消息排队,仅会话结束发送的问题求助

WebSocket消息延迟发送问题排查与解决建议

问题概述

部分WebSocket消息仅在会话结束或连接即将断开时才发送。已确认WebSocket对象一致,主类connect方法中的send_text能立即生效,但AutogenChat类的new_group_chat_received_message方法中的WebSocket调用被排队,无法即时推送到客户端。

代码片段

main.py

class ConnectionManager:
    def __init__(self):
        self.active_connections: Dict[str, AutogenChat] = {}

    async def connect(self, autogen_chat: AutogenChat, client_id: str):
        await autogen_chat.websocket.accept()
        logger.info(f"attempting to connect: {client_id}")  # Log when a connection is opened
        if client_id in self.active_connections:
            # Handle existing connection, e.g., close it or notify the client
            existing_chat = self.active_connections[client_id]
            await existing_chat.websocket.close(code=1001, reason="New connection made")
        self.active_connections[client_id] = autogen_chat
        await autogen_chat.websocket.send_text("connected") #this websocket send appears instantly in my client

    async def disconnect(self, autogen_chat: AutogenChat, client_id: str):
        # await autogen_chat.client_receive_queue.put_nowait("DO_FINISH")
        print(f"autogen_chat {autogen_chat.client_id} disconnected")
        self.active_connections[client_id].websocket = None
        self.active_connections.pop(client_id)


manager = ConnectionManager()



@app.get("/")
async def home():
    return "Welcome Home"


@app.get("/get-chat-id")
async def get_chat_id():
    client_id = str(uuid.uuid1())
    while client_id in client_ids:
        client_id = str(uuid.uuid1())
    client_ids.append(client_id)
    return JSONResponse({"client_id": client_id})

@app.websocket("/autogen/{client_id}")
async def websocket_endpoint(websocket: WebSocket, client_id: str):
    # clients[client_id] = websocket
    autogen_chat = AutogenChat(websocket=websocket, client_id=client_id, logger=logger)
    await manager.connect(autogen_chat, client_id)
    try:
        await manager.active_connections[client_id].connect(client_id)

    except Exception as e:
        logger.error(f"Error in WebSocket connection with client_id {client_id}: {e}")  # Log exceptions
    finally:
        try:
            logger.info(
                f"Closing WebSocket connection with client_id {client_id}")  # Log when a connection is closed
            await manager.disconnect(autogen_chat, client_id)
        except Exception as e:
            logger.error(f"Error while closing WebSocket connection with client_id {client_id}: {e}")


if __name__ == "__main__":
    uvicorn.run(app, host="0.0.0.0", port=8000)

AutogenChat类

class AutogenChat:


    def __init__(self, client_id=None, websocket=None, logger=None):
        self.groupchat = None
        self.websocket = websocket
        self.logger = logger
        self.client_id = client_id


        self.product_manager_assistant = autogen.AssistantAgent(
            name="product_manager",
            llm_config=llm_config,
            system_message="Your job is...."
        )



        self.ux_designer_assistant = autogen.AssistantAgent(
            name="user_experience_designer",
            llm_config=llm_config,
            system_message="Your job is to..."
        )

        self.user_proxy = UserProxyWebAgent(
            name="human_admin",
            human_input_mode="ALWAYS",
            max_consecutive_auto_reply=5,
            default_auto_reply="APPROVED",
            is_termination_msg=lambda x: x.get("content", "") and x.get("content", "").rstrip().endswith("TERMINATE"),
            code_execution_config=False,
            llm_config=llm_config,
            system_message="""A human admin. Interact with the team to receive feedback on ideas. Plan execution needs to be approved by this admin. 
           Only say APPROVED in most cases, and say TERMINATE when nothing is to be done further. Do not say others."""
        )
        try:


            self.user_proxy.set_websocket(self.websocket)
            self.user_proxy.set_logger(self.logger)

        except Exception as e:
            logger.error(f"Error in processing: {e}")

    async def new_group_chat_received_message(self, recipient, messages, sender, config):
        if messages:
            content = messages[-1]['content']
            # print ("message is ")
            # print [messages[-1]]
            if 'name' in messages[-1]:
                name = messages[-1]['name']
                #if (name != "human_admin"):
                payload = {'message': content, 'sender': name}
                json_payload = json.dumps(payload)
                try:
                    print ("before trying")
                    await self.websocket.send_text(json_payload) #This websocket send is queued up and only appears in the client after the session ends
                    print ("after trying")

                except Exception as e:
                    self.logger.error(f"Error in sending websocket message: {e}")  # Log exceptions
                self.logger.info(f"Sent message from new group chat recieved: {json_payload}")
        return False, None

    async def start_group_chat(self, message):
        self.user_proxy.register_reply([autogen.Agent, None], reply_func=self.new_group_chat_received_message,
                                       config={"callback": None})
        self.ux_designer_assistant.register_reply([autogen.Agent, None],
                                                  reply_func=self.new_group_chat_received_message,
                                                  config={"callback": None})
        self.product_manager_assistant.register_reply([autogen.Agent, None],
                                                      reply_func=self.new_group_chat_received_message,
                                                      config={"callback": None})

        self.groupchat = autogen.GroupChat(
            agents=[self.user_proxy, self.ux_designer_assistant, self.product_manager_assistant],
            messages=[], max_round=12)
        self.manager = GroupChatManagerWeb(groupchat=self.groupchat,
                                           llm_config=llm_config,
                                           human_input_mode="ALWAYS")

        await self.user_proxy.a_initiate_chat(
            self.manager,
            clear_history=True,
            message=message
        )
        self.user_proxy.stop_reply_at_receive(self.manager)

    async def send_message(self, message):
        print ("sending message " + message)
        await self.user_proxy.a_send(message,
                    self.manager)
        return self.user_proxy.last_message()["content"]

    async def connect(self,client_id):
        await self.websocket.send_text("Connecting to websocket...")
        while True:
            data = await self.websocket.receive_text()
            #future_calls = asyncio.gather(receive_from_client(autogen_chat, client_id))
            self.logger.info(f"WebSocket connection established with client_id: {client_id}")
            if data:
                message = json.loads(data)
                content = message["message"]["content"]
                message_type = message["message"]["messageType"]
                if message_type == "start":
                    print("we're in start")
                    #await autogen_chat.start_single_chat(content)
                    await self.start_group_chat(content)
                elif message_type == "feedback":
                    print ("we're in feedback")
                    await self.send_message(content)
                elif content == "DO_FINISH":
                    break

核心问题原因

asyncio事件循环被阻塞:在AutogenChat.connect方法的while循环中,调用await self.start_group_chat(content)时,a_initiate_chat是一个长时间运行的协程,会持续占用事件循环,导致new_group_chat_received_message中的await self.websocket.send_text无法被及时调度。所有WebSocket发送操作会被排队,直到a_initiate_chat执行完毕(会话终止)才会批量执行。

解决思路

  1. 将聊天任务移至后台执行
    修改connect方法中调用start_group_chat的代码,用asyncio.create_task启动后台任务,避免阻塞主WebSocket接收循环:

    import asyncio
    
    # ... 其他代码 ...
    if message_type == "start":
        print("we're in start")
        # 用create_task将聊天任务放入后台,不阻塞主循环
        asyncio.create_task(self.start_group_chat(content))
    

    这样主循环可以继续监听客户端消息,事件循环能及时调度WebSocket发送操作。

  2. 确认Autogen回调的异步兼容性
    检查Autogen的register_reply方法是否支持异步回调执行。如果Autogen是同步调用回调,需将异步WebSocket发送操作包装为线程安全的协程调用(如asyncio.run_coroutine_threadsafe),但优先使用Autogen官方推荐的异步回调方式。

  3. 验证WebSocket对象的一致性
    确保self.websocket在整个生命周期中未被意外修改,所有发送操作均在正确的asyncio上下文里执行。

关于asyncio的学习建议

非常有必要深入学习asyncio,因为你的应用基于FastAPI/UVicorn的异步IO框架。核心需要掌握:

  • 事件循环的工作原理:协程调度、await的作用、阻塞与非阻塞操作的区别
  • 后台任务的创建与管理(asyncio.create_task)
  • 协程的并发与同步(asyncio.gather、asyncio.wait等)
  • 异步环境下的资源共享与状态管理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 02:14:54