Asyncio客户端多次发送消息仅首次被服务端接收的问题排查
Asyncio客户端多次发送消息仅首次被服务端接收的问题排查
我来帮你排查这个问题哈!先理清楚你的场景:用asyncio写的TCP客户端和服务端,第一次发消息服务端能正常接收,但后续客户端明明显示发送成功,服务端却收不到消息;但把循环改成两次硬编码的发送逻辑,服务端又能收到两次消息。
核心原因分析
这个问题的根源在于TCP是无边界的流协议,再加上你当前的代码没有处理消息边界,同时客户端的同步input()调用可能间接影响了消息的发送/接收时机。
具体来说:
- TCP流的消息边界问题:你现在的服务端用
reader.read(100)来读取数据,这个方法会返回当前TCP缓冲区中最多100字节的可用数据,但它不知道哪部分是一个完整的消息。第一次发送时,消息较短且是首次发送,服务端刚好能读到完整的消息;但后续发送时,TCP可能把多次发送的数据合并成一个流,或者服务端的read(100)会陷入等待(因为它无法判断当前收到的字节是否是一个完整的消息)。 - 同步
input()的潜在影响:客户端里的input()是同步阻塞函数,会卡住整个asyncio事件循环。虽然你看到客户端打印了“消息发送成功”,但事件循环被阻塞可能导致操作系统的TCP缓冲区没有及时把数据推送给服务端(不过这不是主要原因,核心还是消息边界)。
解决方案
1. 给消息添加明确的边界(必做)
因为TCP是流协议,必须让服务端知道一个完整的消息在哪里结束。最常用的方式是给每个消息加换行符,服务端按行读取:
客户端修改send_message方法:
在JSON字符串末尾添加换行符,作为消息结束的标记:
async def send_message(message): if not reader or not writer: print("[Error]: Not connected to the message server.") return try: message_data = { "sender": "name", "message": message } # 添加换行符作为消息边界 request_data = json.dumps(message_data) + '\n' writer.write(request_data.encode('utf-8')) await writer.drain() print(f"[Info]: Message sent: {message}") except Exception as e: print(f"[Error]: Failed to send message. Error: {e}")
服务端修改handle_client的读取逻辑:
用reader.readline()代替read(100),按行读取完整的消息:
async def handle_client(reader, writer): addr = writer.get_extra_info('peername') print(f"[Info]: New connection from {addr}") clients.append((reader, writer)) try: while True: # 按行读取,自动识别换行符作为消息边界 data = await reader.readline() if not data: break # 连接关闭 # 去掉换行符和空白字符,解析JSON message_str = data.decode('utf-8').strip() if not message_str: continue print(f"[Info]: Received message: {message_str}") # 可选:解析JSON做后续处理 try: message_data = json.loads(message_str) print(f"[Detail] Sender: {message_data['sender']}, Content: {message_data['message']}") except json.JSONDecodeError as e: print(f"[Error]: Failed to parse message JSON. Error: {e}") except asyncio.CancelledError: pass except Exception as e: print(f"[Error]: Error handling client {addr}. Error: {e}") finally: # 确保从客户端列表中移除 if (reader, writer) in clients: clients.remove((reader, writer)) print(f"[Info]: Connection closed by {addr}") writer.close() await writer.wait_closed()
2. 替换同步input()为异步输入(可选但推荐)
客户端的input()是同步阻塞函数,会卡住asyncio的事件循环,影响异步任务的执行。你可以用aioconsole库的ainput()来实现异步输入:
- 先安装库:
pip install aioconsole
- 修改客户端的
handle_input方法:
from aioconsole import ainput async def handle_input(): while True: # 用异步输入代替同步input,不阻塞事件循环 user_input = await ainput("[You]: ") await send_message(user_input)
为什么硬编码两次发送能成功?
因为两次发送之间有input()的同步阻塞,给了TCP足够的时间把第一次发送的数据推送到服务端,服务端的read(100)刚好能读到完整的消息;而循环发送时,消息发送的间隔可能很短,TCP把多次发送的数据合并成了一个流,服务端的read(100)无法判断这是多个消息,导致后续的消息被“卡住”在缓冲区里,直到凑够100字节或者连接关闭。
备注:内容来源于stack exchange,提问作者hhhhh
相关产品推荐
相关产品推荐

