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

Python websockets实现客户端发消息后等待服务端响应的方案

问题根因

之前asyncio.Lock未生效的核心原因是锁的覆盖范围错误:仅将第一次send操作包裹在锁上下文内,锁释放后未等待服务端返回响应就直接执行了第二次发送,完全没有把「等待对应响应」的逻辑纳入锁保护范围,锁没有起到流程串行化的作用。

严格的「发一条消息等一条对应响应」逻辑本质是串行化请求-响应流程,websocket基于TCP实现,本身保证消息传输顺序,只要服务端按接收顺序回包,客户端按发送顺序等待接收即可保证请求和响应一一对应。


正确实现

服务端逻辑要求

服务端必须遵循「收到一条消息→完成业务处理→立刻返回对应响应」的逻辑,不要攒批回包、不要异步延迟回包打乱响应顺序,示例代码如下:

import asyncio
import websockets

class WebSocketServer:
    async def runServer(self):
        server = await websockets.serve(
            self.onConnect,
            "localhost",
            port=8765
        )
        print("Server started listening to new connections...")
        await server.wait_closed()

    async def onConnect(self, ws):
        try:
            while True:
                # 接收单条消息
                message = await ws.recv()
                print(f"Server received message: {message}")
                # 执行业务处理逻辑
                response = f"Response for: {message}"
                # 处理完成立刻返回对应响应
                await ws.send(response)
        except websockets.exceptions.ConnectionClosed:
            print("Client disconnected")

if __name__ == "__main__":
    server = WebSocketServer()
    asyncio.run(server.runServer())

客户端实现

场景1:无并发发送需求(最简单实现)

如果客户端不存在多协程同时触发发送的场景,完全不需要锁,严格按照「发送消息→阻塞等待对应响应→响应返回后再执行下一次发送」的顺序编写代码即可:

import asyncio
import websockets

class WebSocketClient:
    async def connect(self):
        try:
            async with websockets.connect("ws://localhost:8765") as ws:
                print("Connected to the server.")
                
                # 发送第一条消息
                await ws.send("First message")
                # 阻塞等待第一条消息的响应,未收到前不会执行后续逻辑
                first_resp = await ws.recv()
                print(f"First message response: {first_resp}")

                # 确认收到第一条响应后,再发送第二条消息
                await ws.send("Second message")
                second_resp = await ws.recv()
                print(f"Second message response: {second_resp}")

                # 后续可在此处循环接收服务端主动推送的消息
                while True:
                    push_msg = await ws.recv()
                    print(f"Server proactive push: {push_msg}")
        except Exception as e:
            print(f"Connection error: {e}")

if __name__ == "__main__":
    client = WebSocketClient()
    asyncio.run(client.connect())

场景2:存在多协程并发发送需求

如果业务中有多个协程可能独立触发消息发送,再使用asyncio.Lock做流程控制,注意锁必须覆盖「发送消息+等待对应响应」的完整流程,之前的实现错误就是仅锁了send操作,未将等待响应的逻辑纳入锁范围:

import asyncio
import websockets

class WebSocketClient:
    def __init__(self):
        self.ws = None
        self.send_lock = asyncio.Lock()

    async def send_and_wait_resp(self, message):
        # 获取锁,保证同一时间只有一个请求-响应流程在执行
        async with self.send_lock:
            await self.ws.send(message)
            # 等待响应的过程必须在锁上下文内,锁未释放前其他协程无法发送消息
            response = await self.ws.recv()
            return response

    async def connect(self):
        try:
            async with websockets.connect("ws://localhost:8765") as ws:
                self.ws = ws
                print("Connected to the server.")
                
                # 即使多个协程同时调用该方法,也不会出现连续发消息未等响应的情况
                resp1 = await self.send_and_wait_resp("First message")
                print(f"First message response: {resp1}")
                resp2 = await self.send_and_wait_resp("Second message")
                print(f"Second message response: {resp2}")

                while True:
                    push_msg = await ws.recv()
                    print(f"Server proactive push: {push_msg}")
        except Exception as e:
            print(f"Connection error: {e}")

if __name__ == "__main__":
    client = WebSocketClient()
    asyncio.run(client.connect())

补充说明

如果后续服务端引入异步处理逻辑、无法保证按请求接收顺序返回响应,需要给每个请求添加唯一ID,服务端返回响应时携带对应请求ID,客户端收到响应后先匹配ID再分发给对应的等待逻辑,避免请求响应错位。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 16:34:12