1.核心概念
Socket(套接字)是操作系统提供的一套网络通信 API。它处于传输层(TCP/UDP)之上,应用层之下。
如果你觉得网络协议太抽象,可以用“打跨国电话”来理解 Socket 的核心概念:
- Socket (套接字):电话机。想通话,双方都必须有一台电话机。
- IP 地址:国家和城市区号(比如:192.168.1.100)。用来在网络中唯一确定一台主机。
- Port (端口号):分机号(比如:8080)。一台电脑可能同时运行微信、浏览器、游戏,端口号用来唯一确定某个具体的应用程序。范围是 0 ~ 65535。
- TCP 协议:双向确认电话。必须等对方接听、说“喂”,建立稳定连接后才能开始聊,不漏字,不丢包。
2.TCP Socket 通信模型(工作流程)
在 TCP 协议下,Server(服务端)和Client(客户端)的通信有着严格的先后顺序和握手步骤:
Server (服务端) Client (客户端)
============== ==============
1. 创建电话机 (socket)
2. 绑定手机号 (bind)
3. 开启待机铃声 (listen)
4. 苦苦等待来电 (accept) <--- 阻塞等待
5. 创建电话机 (socket)
6. 拨打对方号码 (connect)
================== 握手成功,建立连接 ==================
7. 听对方说话 (recv) <-------- 阻塞等待 -------- 8. 说话/发送 (send)
9. 回复对方 (send) --------------------------> 10. 听回复 (recv)
11. 挂断电话 (close) 12. 挂断电话 (close)
3.实现一个支持持续对话的 C/S 架构单线程程序。
服务端代码:server.py
服务端需要绑定本机的 IP和端口,并保持循环监听。
import socket
def start_server():
# 1. 创建 socket 对象
# AF_INET 表示使用 IPv4 协议,SOCK_STREAM 表示使用 TCP 协议(面向连接的流式套接字)
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
# 解决:重启服务器时提示 "Address already in use" 的问题
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
# 2. 绑定 IP 地址和端口号 ('0.0.0.0' 表示监听本机所有的网络接口)
server_address = ('0.0.0.0', 8888)
server_socket.bind(server_address)
print(f"[系统] 服务器已启动,正在监听端口 {server_address[1]}...")
# 3. 开始监听,5 表示允许排队等待连接的最大队列长度
server_socket.listen(5)
try:
while True:
print("[系统] 等待客户端连接...")
# 4. 阻塞等待连接。有客户端连接时,返回一个全新的连接套接字和客户端的 IP/Port
client_socket, client_address = server_socket.accept()
print(f"[系统] 成功连接到客户端: {client_address}")
# 与当前连接的客户端进行持续交互
while True:
# 5. 接收客户端发来的数据,最大接收 1024 字节
# 注意:网络传输的是 bytes,需要 decode() 解码成字符串
data = client_socket.recv(1024)
# 如果收到空数据,说明客户端已经主动断开连接
if not data:
print(f"[系统] 客户端 {client_address} 已断开连接。")
break
message = data.decode('utf-8')
print(f"[收到客户端消息]: {message}")
# 6. 回复客户端
reply_message = f"服务器已收到你的消息: [{message}]"
client_socket.send(reply_message.encode('utf-8'))
# 关闭当前客户端的连接,准备迎接下一个客户端
client_socket.close()
except KeyboardInterrupt:
print("\n[系统] 服务器正在关闭...")
finally:
# 关闭监听套接字
server_socket.close()
if __name__ == "__main__":
start_server()
客户端代码:client.py
客户端直接向服务端的 IP 和端口发起连接。
import socket
def start_client():
# 1. 创建客户端 socket 对象
client_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
# 2. 目标服务器的 IP 和端口
# 如果在同一台电脑上测试,IP 可以写 '127.0.0.1'
server_address = ('127.0.0.1', 8888)
try:
# 3. 向服务端发起 TCP 连接(三次握手在此处发生)
client_socket.connect(server_address)
print("[系统] 成功连接到服务器!输入 'exit' 可以退出。")
while True:
# 4. 获取用户输入的文本
user_input = input("请输入发送给服务器的消息: ")
if user_input.lower() == 'exit':
print("[系统] 正在断开连接...")
break
if not user_input.strip():
print("发送内容不能为空,请重新输入。")
continue
# 5. 发送数据 (必须 encode 转换为字节流)
client_socket.send(user_input.encode('utf-8'))
# 6. 接收服务器的返回数据
server_reply = client_socket.recv(1024)
print(f"[服务器回复]: {server_reply.decode('utf-8')}\n")
except ConnectionRefusedError:
print("[错误] 无法连接到服务器,请确保服务端已启动。")
finally:
# 7. 关闭套接字,释放资源
client_socket.close()
print("[系统] 已退出。")
if __name__ == "__main__":
start_client()
4. 核心方法及注意事项
- send(bytes) 与 recv(bufsize):
网络上跑的全部是字节流。Python 3 严格区分了 str 和 bytes。所以发送前必须 .encode(‘utf-8’),接收后必须 .decode(‘utf-8’)。 - 什么是”阻塞(Blocking)”:
上面的 accept() 和 recv() 默认都是阻塞的。也就是说,如果没有新客户端连接,或者连上的客户端一直不发消息,代码就会卡在这一行静静等待,不会往下走。 - 当前代码的局限性:
由于采用了单线程单循环结构,这个临时的 server.py 同一时间只能服务一个客户端。如果有第二个客户端连进来,它必须排队,直到第一个客户端输入 exit 断开后,第二个客户端才会被 accept() 接纳。
10 个请求同时进来会发生什么?
- 第 1 个客户端进来了,accept() 成功接待,代码进入了内部的 while 循环,开始执行 client_socket.recv()。
- 此时,第 2 到第 10 个客户端也发出连接请求。因为操作系统有缓存队列(我们在 listen(5) 里设置的长度),这 9 个客户端显示连接成功。
- 但是! 此时服务器的单线程正卡在第 1 个客户端的 recv() 上。如果第 1 个客户端不发消息、也不断开连接,服务器的代码就永远不会回到外层的 accept()。
- 结果就是:第 2 到第 10 个客户端虽然显示连上了,但他们发送的任何消息,服务器都完全不理会(因为没有执行到属于他们的 recv 代码)。他们只能死等,直到第 1 个客户端断开。
5. 改进版
方法 A:多线程(Threading)
接待完一个客户端后,立刻创建一个新线程,把这个客户端的 client_socket 交给线程去负责 recv()。主线程立刻回到 accept() 继续接待下一个人。
server_threaded.py
import socket
import threading
def handle_client(client_socket, client_address):
"""每个线程独立的函数,专门负责这一个客户端的 recv 和 send"""
print(f"[线程] 开始服务客户端: {client_address}")
try:
while True:
data = client_socket.recv(1024)
if not data:
break
print(f"[{client_address} 发来]: {data.decode('utf-8')}")
message = data.decode('utf-8')
# 6. 回复客户端
reply_message = f"服务器已收到你的消息: [{message}]"
client_socket.send(reply_message.encode('utf-8'))
finally:
client_socket.close()
print(f"[线程] 客户端 {client_address} 已退出,释放资源。")
def start_server():
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server.bind(('0.0.0.0', 8888))
server.listen(10) # 增大队列,允许更多人排队
while True:
# 主线程只负责 accept() 迎宾
client_socket, client_address = server.accept()
print(f"[系统] 迎来新客户端: {client_address}")
# 只要来一个人,就开一个新线程,把 client_socket 传进去
t = threading.Thread(target=handle_client, args=(client_socket, client_address))
t.daemon = True # 设置为守护线程,主线程退时一起退
t.start() # 启动线程,去执行 handle_client
if __name__ == "__main__":
start_server()
方法 B: I/O 多路复用
selectors 模块实现 I/O 多路复用(事件驱动),一些python web框架底层就是使用这种方法
server_selectors.py
import socket
import selectors
# 1. 初始化选择器(自动根据系统选择最高效的底层:Mac 用 kqueue, Linux 用 epoll)
sel = selectors.DefaultSelector()
def accept_handler(server_socket, mask):
"""【迎宾事件处理】当有新客户端发起连接请求时触发"""
client_socket, client_address = server_socket.accept()
print(f"[系统] 迎来新客户端: {client_address}")
# 将新创建的客户端套接字设置为【非阻塞】模式(多路复用的核心前提!)
client_socket.setblocking(False)
# 将这个客户端套接字也注册到选择器中,监听它的【读就绪(EVENT_READ)】事件
# 并把具体的业务处理函数 read_handler 绑定在 data 参数中
sel.register(client_socket, selectors.EVENT_READ, data=read_handler)
def read_handler(client_socket, mask):
"""【点单/通信事件处理】当连上的客户端有数据发过来、或者断开连接时触发"""
client_address = client_socket.getpeername() # 获取客户端的 IP/Port
try:
# 接收数据。因为是非阻塞的,此时调用必然能立刻拿到数据,绝不卡死
data = client_socket.recv(1024)
if data:
message = data.decode('utf-8')
print(f"[{client_address} 发来]: {message}")
# 回复客户端
reply = f"服务器(多路复用版)收到: {message}"
client_socket.sendall(reply.encode('utf-8'))
else:
# 收到空数据,说明客户端主动关闭了连接
print(f"[系统] 客户端 {client_address} 正常断开连接。")
# 必须先从选择器中注销监视
sel.unregister(client_socket)
# 再关闭套接字释放资源
client_socket.close()
except ConnectionResetError:
# 处理客户端异常断开(比如直接掐断程序)
print(f"[系统] 客户端 {client_address} 异常断开。")
sel.unregister(client_socket)
client_socket.close()
def start_server():
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server.bind(('0.0.0.0', 8888))
server.listen(100) # 增大队列,高并发标配
# 关键点 1:必须将监听套接字设为【非阻塞】
server.setblocking(False)
# 关键点 2:将监听套接字注册到选择器,监听【读就绪】(即有新请求进来)
# 并将 accept_handler 函数绑定到这个事件上
sel.register(server, selectors.EVENT_READ, data=accept_handler)
print("[系统] I/O 多路复用并发服务器启动,正在监听 8888 端口...")
try:
while True:
# 关键点 3:进入大循环,无限阻塞等待
# 当没有任何事件发生时,程序停在这里;一旦有任何 socket 动了,它立刻醒来并返回
events = sel.select()
# 遍历所有已经准备就绪的事件
for key, mask in events:
# key.fileobj 是当前触发事件的 socket 对象
# key.data 是我们注册时绑定的处理函数(accept_handler 或 read_handler)
callback = key.data
# 关键点 4:直接执行对应的回调函数,快速处理
callback(key.fileobj, mask)
except KeyboardInterrupt:
print("\n[系统] 服务器正在关闭...")
finally:
sel.close()
if __name__ == "__main__":
start_server()
mask 的本质是什么?
mask 是一个整数(在底层是位掩码 Bitmask),用来表示当前这个套接字(Socket)上具体发生了什么事件。在 selectors 模块中,主要有两种核心事件:
- selectors.EVENT_READ(值为 1):表示读就绪。即有新连接进来(accept)或者有新数据到了(recv)。
- selectors.EVENT_WRITE(值为 2):表示写就绪。即服务器的发送缓冲区有空闲,可以往外发数据了(send)。
问题:服务端不是要写数据吗 为何没有写事件监听?
核心原因:操作系统的缓冲区通常总是“可写”的
在 TCP 连接建立后,操作系统在底层为每个 Socket 分配了两个缓冲区:一个接收缓冲区(读),一个发送缓冲区(写)。
- 读事件(EVENT_READ)的触发条件:只有当客户端发来数据,服务器的接收缓冲区里有数据时,才会触发读事件。因为平时绝大多数时间缓冲区是空的,所以服务器平时不触发读,只有客户端说话时才触发。
- 写事件(EVENT_WRITE)的触发条件:只要服务器的发送缓冲区没有满,就会一直触发写事件。因为网速通常很快,发送缓冲区在 99.9% 的时间里都是空闲且未满的。这就导致了一个致命问题:如果你把一个 Socket 注册了“写事件”,由于发送缓冲区几乎总是空的,sel.select() 就会像发了疯一样疯狂提示你“可写了!可写了!”,导致你的 while True 循环瞬间飙到 CPU 100%,也就是所谓的“忙轮询(Busy Loop)”。
我们要回复的内容非常短(只有几十个字节),而操作系统的发送缓冲区通常有几百 KB。这就好比你要把一小杯水(几十字节)倒进一个巨大的空水桶(发送缓冲区)里,它绝对能瞬间倒进去,绝对不会发生卡顿(阻塞)。既然直接调用 sendall() 能够秒发成功、不会卡死程序,我们自然就偷了个懒,不需要大费周章地去注册和监听写事件了。
events = sel.select()解析
这个就是监听注册在sel上的事件,在start_server里注册了EVENT_READ读事件,因为socket其实在系统层面也是文件,所以接受新建客户端连接发来的数据也是读,这个注册在server上,在accept_handler里把read_handler注册到当前client_socket。
sel.register(client_socket, selectors.EVENT_READ, data=read_handler)
这两个注册都是系统级别的
关键代码解析
events = sel.select()
# 遍历所有已经准备就绪的事件
for key, mask in events:
# key.fileobj 是当前触发事件的 socket 对象
# key.data 是我们注册时绑定的处理函数(accept_handler 或 read_handler)
callback = key.data
# 关键点 4:直接执行对应的回调函数,快速处理
callback(key.fileobj, mask)
(1) events = sel.select()
如果有新建socket连接或者有已经创建socket客户端发来的数据,这里会唤醒,进入for 循环
- 如果是新建socket连接:
callback就是注册的accept_handler,key.fileobj就是server_socket,然后调用accept_handler函数;
在accept_handler里创建client_socket并将这个客户端套接字也注册到选择器中,监听它的【读就绪(EVENT_READ)】事件,并把具体的业务处理函数 read_handler 绑定在 data 参数中 - 如果是已经创建socket客户端发来的数据:
callback就是注册的read_handler,key.fileobj就是对应的client_socket,然后调用read_handler函数去读数据。
在底层操作系统内核里,确实有实实在在地维护着一张这样(甚至更高效)的表!无论是 Linux 的 epoll 还是 Mac 的 kqueue,它们能做到高性能的核心秘密,就是把原来需要由 Python 程序员自己用代码维护和遍历的表,直接做进了操作系统的内核内存中。为了让你对底层的这张“表”有一个具象的认识,我们以目前互联网服务器用得最多的 Linux epoll 为例,看看它在内核里到底维护了什么:
内核里真实存在的两张“表”
Linux 内核在处理多路复用时,并不是只用一张死板的表格,而是为了极致的性能,在内核内存中同时维护了两套数据结构:
第一张表:内核监视红黑树(The Interest List)
当你调用 sel.register(socket, …) 时,内核就会把这个 Socket 的文件描述符(FD)和你想监听的事件,放进一棵红黑树(Red-Black Tree)里。
- 为什么用红黑树? 因为红黑树增加、删除、查找一个 Socket 的速度极快(时间复杂度是 (O(\log n)))。即使你注册了 10 万个客户端 Socket,内核也能在微秒级内完成管理。
- 里面的内容:实实在在地记录着 [FD 3 -> 监听读], [FD 4 -> 监听读]。
第二张表:就绪链表(The Ready List)
能够秒回的关键!内核还维护了一个双向链表。
- 当没有任何人发消息、也没有新连接时,这个链表是完全空的。
- 一旦客户端 A 发来数据,网卡触发中断,内核发现是 FD 4 动了。内核就会顺手把 FD 4 丢进这个“就绪链表”里。
Python 层的 key.fileobj 和 key.data 又是怎么对上的?
操作系统内核只认数字(比如 FD 3、FD 4),它可不知道 Python 里的 accept_handler 或者是哪个具体的 socket 对象。那 Python 是怎么把内核的数字和我们的函数对上的呢?其实,在 Python 的 selectors.DefaultSelector 对象内部,在 Python 应用程序的内存里,也镜像维护了一张小表(一个普通的 Python 字典 _fd_to_key):
{4: SelectorKey(fileobj=<socket.socket fd=4, family=2, type=1, proto=0, laddr=('0.0.0.0', 8888)>, fd=4, events=1, data=<function accept_handler at 0x10117e200>),
5: SelectorKey(fileobj=<socket.socket fd=5, family=2, type=1, proto=0, laddr=('127.0.0.1', 8888), raddr=('127.0.0.1', 54637)>, fd=5, events=1, data=<function read_handler at 0x1012a9300>),
6: SelectorKey(fileobj=<socket.socket fd=6, family=2, type=1, proto=0, laddr=('127.0.0.1', 8888), raddr=('127.0.0.1', 54657)>, fd=6, events=1, data=<function read_handler at 0x1012a9300>)}
我们来修改server_selectors.py,在accept_handler下新增 pprint.pprint(dict(sel.get_map())),完整代码如下:
import socket
import selectors
import pprint
# 1. 初始化选择器(自动根据系统选择最高效的底层:Mac 用 kqueue, Linux 用 epoll)
sel = selectors.DefaultSelector()
def accept_handler(server_socket, mask):
"""【迎宾事件处理】当有新客户端发起连接请求时触发"""
client_socket, client_address = server_socket.accept()
print(f"[系统] 迎来新客户端: {client_address}")
# 将新创建的客户端套接字设置为【非阻塞】模式(多路复用的核心前提!)
client_socket.setblocking(False)
# 将这个客户端套接字也注册到选择器中,监听它的【读就绪(EVENT_READ)】事件
# 并把具体的业务处理函数 read_handler 绑定在 data 参数中
sel.register(client_socket, selectors.EVENT_READ, data=read_handler)
pprint.pprint(dict(sel.get_map()))
def read_handler(client_socket, mask):
"""【点单/通信事件处理】当连上的客户端有数据发过来、或者断开连接时触发"""
client_address = client_socket.getpeername() # 获取客户端的 IP/Port
try:
# 接收数据。因为是非阻塞的,此时调用必然能立刻拿到数据,绝不卡死
data = client_socket.recv(1024)
if data:
message = data.decode('utf-8')
print(f"[{client_address} 发来]: {message}")
# 回复客户端
reply = f"服务器(多路复用版)收到: {message}"
client_socket.sendall(reply.encode('utf-8'))
else:
# 收到空数据,说明客户端主动关闭了连接
print(f"[系统] 客户端 {client_address} 正常断开连接。")
# 必须先从选择器中注销监视
sel.unregister(client_socket)
# 再关闭套接字释放资源
client_socket.close()
except ConnectionResetError:
# 处理客户端异常断开(比如直接掐断程序)
print(f"[系统] 客户端 {client_address} 异常断开。")
sel.unregister(client_socket)
client_socket.close()
def start_server():
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server.bind(('0.0.0.0', 8888))
server.listen(100) # 增大队列,高并发标配
# 关键点 1:必须将监听套接字设为【非阻塞】
server.setblocking(False)
# 关键点 2:将监听套接字注册到选择器,监听【读就绪】(即有新请求进来)
# 并将 accept_handler 函数绑定到这个事件上
sel.register(server, selectors.EVENT_READ, data=accept_handler)
print("[系统] I/O 多路复用并发服务器启动,正在监听 8888 端口...")
try:
while True:
# 关键点 3:进入大循环,无限阻塞等待,类似于之前的 `ready = selector.select(0.5)`
# 当没有任何事件发生时,程序停在这里;一旦有任何 socket 动了,它立刻醒来并返回
events = sel.select()
# 遍历所有已经准备就绪的事件
for key, mask in events:
# key.fileobj 是当前触发事件的 socket 对象
# key.data 是我们注册时绑定的处理函数(accept_handler 或 read_handler)
callback = key.data
# 关键点 4:直接执行对应的回调函数,快速处理
callback(key.fileobj, mask)
except KeyboardInterrupt:
print("\n[系统] 服务器正在关闭...")
finally:
sel.close()
if __name__ == "__main__":
start_server()
当创建一个socket连接时,服务端控制台输出信息:
{4: SelectorKey(fileobj=<socket.socket fd=4, family=2, type=1, proto=0, laddr=('0.0.0.0', 8888)>, fd=4, events=1, data=<function accept_handler at 0x10117e200>),
5: SelectorKey(fileobj=<socket.socket fd=5, family=2, type=1, proto=0, laddr=('127.0.0.1', 8888), raddr=('127.0.0.1', 54637)>, fd=5, events=1, data=<function read_handler at 0x1012a9300>)}
当新增一个socket连接时,服务端控制台输出信息:
{4: SelectorKey(fileobj=<socket.socket fd=4, family=2, type=1, proto=0, laddr=('0.0.0.0', 8888)>, fd=4, events=1, data=<function accept_handler at 0x10117e200>),
5: SelectorKey(fileobj=<socket.socket fd=5, family=2, type=1, proto=0, laddr=('127.0.0.1', 8888), raddr=('127.0.0.1', 54637)>, fd=5, events=1, data=<function read_handler at 0x1012a9300>),
6: SelectorKey(fileobj=<socket.socket fd=6, family=2, type=1, proto=0, laddr=('127.0.0.1', 8888), raddr=('127.0.0.1', 54657)>, fd=6, events=1, data=<function read_handler at 0x1012a9300>)}
这说明这个表确实存在,数字是key及SelectorKey信息
6. 粘包现象
“粘包现象”(Packet Stickiness / Packet Coalescing)是网络编程新手在处理 TCP 套接字通信时,最容易崩溃、也最难以绕过去的一个“惊天巨坑”。用大白话解释,粘包就是:客户端明明分两次发送了 “hello” 和 “world”,结果服务端在一次 recv() 中却一口气收到了 “helloworld”,两条独立的消息像“年糕”一样粘死在一起了。为了让你彻底看清并解决这个魔鬼现象,我们从它的本质原因、发生场景和解决方案来分析:
(1)为什么会发生“粘包”?
发生粘包的根本原因只有一句话:TCP 协议是一个面向“流(Stream)”的协议,它的底层根本没有“消息”或“报文”的边界概念!人类的直觉:我调用一次 send(“hello”),这就是一封独立的信;再调用 send(“world”),这是第二封信。TCP 的现实:TCP 看待数据就像是一根自来水管里流出来的水(字节流)。它不关心你是一杯水倒进去的,还是一桶水倒进去的。对它来说,流进来的全都是水滴(字节)。当客户端连续快速地发送数据时,操作系统底层的发送缓冲区会把它们连续地排在一起;服务端的 recv(1024) 就像是一台大功率抽水机,它会在接收缓冲区里有多少就抽多少(最大 1024 字节)。结果,就把两段没有任何标记的水流一次性全抽了上来,这就导致了粘包。(注:如果一条大消息因为缓冲区不够大,被切成了两半接收,那种现象叫“半包”)
(2)粘包发生的两个经典时机
- 发送端粘包(Nagle 算法):
为了提高网络传输效率,操作系统的 TCP 协议默认开启了 Nagle 算法。如果你的消息非常短(比如只有几个字节),OS 不会傻傻地来一个发一个(因为网络包头太重了),而是会等一等,把几个小包攒成一个大包一起发出去。 - 接收端粘包(接收过快):
如果服务端在处理别的事情,导致 sel.select() 醒来得慢了一点,此时客户端已经连续发了 3 条消息,它们在服务端的接收缓冲区里整齐地排好队了。当服务器终于执行 recv(1024) 时,就会一网打尽,把 3 条消息当成一条读出来。
(3)解决方法
方法 A:
TLV 格式 / 固定包头长度法(全行业最标准、最硬核)在发送真正的聊天数据之前,先发送一个固定长度(比如 4 字节)的“整数”,告诉对方接下来这条消息到底有多少个字节。
【一个标准网络包的结构】
┌───────────────────────────┬──────────────────────────────────┐
│ 包头:消息长度 (固定4字节) │ 包体:真正的业务数据 (可变长度) │
│ 例如:数字 5 │ 例如:"hello" │
└───────────────────────────┴──────────────────────────────────┘
服务端解析流程:
- 1.不管三七二十一,先死板地接收且只接收 4 个字节(recv(4))。
- 2.把这 4 个字节还原成数字(比如读出数字是 5)。
- 3.接下来,极其精准地只接收 5 个字节(recv(5)),这就拿到了 “hello”。哪怕后面紧接着死死粘着 “world”,也绝对不会多读一个字节。这就从物理上彻底斩断了粘包。
方法 B:
特殊分隔符法(特定场景使用)规定每条消息的末尾必须带上特定的结束符(比如换行符 \n 或自定义的 [END])。服务端在读取字节流时,自己去扫描这个结束符,手动切分字符串。(HTTP 协议的请求头就是用 \r\n\r\n 来区分头和身体的)。
(4)实战:用 Python struct 模块消灭粘包
Python 的 struct 模块可以把普通的整数(如消息长度)完美打包成固定 4 个字节的二进制字节流,专门用来做网络包头。
客户端
import socket
import struct
def send_msg(sock, text):
"""防粘包发送:先将包体长度打包成固定 4 字节二进制头,拼接后一次性送出"""
body = text.encode('utf-8')
# '!I' 表示使用标准的网络字节序(大端)将长度转为 4 字节二进制
header = struct.pack('!I', len(body))
sock.sendall(header + body)
def recv_msg(sock):
"""防粘包接收:严格按照 4 字节包头精准解包,完美 Hold 住服务端的应答,防止半包"""
try:
# 1. 精准接收 4 字节包头
header = sock.recv(4)
if not header:
return None
# 2. 解出接下来的响应内容到底有多少个字节
body_len = struct.unpack('!I', header)[0]
# 3. 循环接收,直到收满 body_len 字节为止,彻底杜绝半包现象
data_buffer = b""
while len(data_buffer) < body_len:
packet = sock.recv(body_len - len(data_buffer))
if not packet:
break
data_buffer += packet
return data_buffer.decode('utf-8')
except OSError:
return None
def start_client():
client = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
try:
client.connect(('127.0.0.1', 8888))
print("[系统] 成功连接到高性能并发服务器!\n")
# 【测试 1】:连续快速交互,中间不加任何 time.sleep,验证绝对不会粘包
print("[客户端] 发送: hello")
send_msg(client, "hello")
# 原地等待服务端回复。此处的阻塞完美杜绝了客户端进程因“光速闪退”而震碎服务端 Socket
reply1 = recv_msg(client)
print(f"[服务器回复]: {reply1}\n")
# 【测试 2】
print("[客户端] 发送: world")
send_msg(client, "world")
reply2 = recv_msg(client)
print(f"[服务器回复]: {reply2}\n")
except ConnectionRefusedError:
print("[错误] 无法连接到服务器,请确保服务端已启动。")
finally:
# 通信结束,优雅关闭连接,展现良好网络市民的文明规范
client.close()
print("[系统] 连接已安全断开,程序退出。")
if __name__ == "__main__":
start_client()
服务端
import socket
import selectors
import struct
# 1. 初始化选择器(在 Mac 上会自动输出并使用 KqueueSelector)
sel = selectors.DefaultSelector()
class ClientState:
"""每个客户端专属的『内存蓄水池』,用于在非阻塞模式下完美解决粘包与半包"""
def __init__(self):
self.buffer = b"" # 存放该客户端独有的、未处理完的原始字节流
def accept_handler(server_socket, mask):
"""【迎宾事件处理】大门被敲响时触发。负责开门迎客并为新客安排专属坐席"""
client_socket, client_address = server_socket.accept()
print(f"[系统] 迎来新客户端: {client_address}")
# 核心前提:必须将客人的套接字设为非阻塞模式
client_socket.setblocking(False)
# 为这个新进来的客户端实例化一个专属的蓄水池
state = ClientState()
# 将(处理函数, 蓄水池对象)打包成元组作为 data 注册进内核,实现动态绑定
sel.register(client_socket, selectors.EVENT_READ, data=(read_handler, state))
def read_handler(client_socket, mask, state):
"""【点单/通信事件处理】融入了 TLV 解包循环以及防御 Mac 连接暴毙的强壮逻辑"""
# 【第一道防线】防御 Mac 平台特有的 Errno 22 (Invalid argument)
# 如果客户端发完数据光速退出了,此时查询名字在 Mac 底层会直接抛 OSError
try:
client_address = client_socket.getpeername()
except (OSError, ConnectionResetError) as e:
print(f"[系统] 客户端在建立视图前已骤断或失效(成功拦截 Errno 22): {e}")
try:
sel.unregister(client_socket)
client_socket.close()
except Exception:
pass
return
# 【第二道防线】正常的通信处理,包裹在 try-except 中以应对网络波动和异常掐断
try:
# 1. 顺着网线把当前网卡里积压的数据一网打尽(非阻塞,有多少收多少,绝不卡死)
recv_data = client_socket.recv(1024)
if not recv_data:
# 收到空数据,说明客户端调用了 close() 进行了优雅的正常断开
print(f"[系统] 客户端 {client_address} 正常断开连接。")
sel.unregister(client_socket)
client_socket.close()
return
# 2. 把刚抽上来的水,追加到这个客户端专属的蓄水池里
state.buffer += recv_data
# 3. 核心:通过循环不断切分蓄水池。即使几条消息粘在了一起,也能在这里像切香肠一样被一一切开
while True:
# 情况 A:如果池子里的数据连 4 字节的包头都不够,说明长度信息还没传输全,跳出循环继续等待下一次 recv
if len(state.buffer) < 4:
break
# 情况 B:够了 4 字节,先把头切出来,拆箱(unpack)算出接下来的业务内容到底有多长
header = state.buffer[:4]
# '!I' 表示网络字节序(大端无符号整数),unpack 返回元组,取第 0 个元素
body_len = struct.unpack('!I', header)[0]
# 情况 C:如果池子的总长度,小于(4字节头 + 算出来的身体长度),说明产生了“半包”
# 此时绝不能强行去读,必须跳出循环,静静等待下一次接收事件把剩下的身体拼进来
if len(state.buffer) < 4 + body_len:
break
# 情况 D:完美!头齐了,身体也全了。精准地把这一条独立的消息从池子里切出来
body_data = state.buffer[4:4 + body_len]
# 斩断已经处理完的数据,让池子里剩下的数据往前挪,准备迎接下一次循环检查
state.buffer = state.buffer[4 + body_len:]
# 业务逻辑:解码并打印干净、无粘包的独立消息
message = body_data.decode('utf-8')
print(f"[{client_address} 发来]: {message}")
# 回复客户端(注意:因为客户端也防粘包,服务器发回去的数据也要带上 4 字节包头!)
reply_body = f"服务器(多路复用版)收到: [{message}]".encode('utf-8')
reply_header = struct.pack('!I', len(reply_body))
client_socket.sendall(reply_header + reply_body)
except (ConnectionResetError, OSError) as e:
# 处理客户端异常强制断开(比如直接在终端掐断客户端程序)
print(f"[系统] 客户端 {client_address} 异常断开或传输受阻: {e}")
try:
sel.unregister(client_socket)
client_socket.close()
except Exception:
pass
def start_server():
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
# 允许端口快速重用,彻底解决 "Address already in use" 报错
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server.bind(('0.0.0.0', 8888))
server.listen(100)
# 必须将大门的监听套接字也设为非阻塞模式
server.setblocking(False)
# 把大门注册进内核。因为大门不需要蓄水池,data 直接传单一函数即可
sel.register(server, selectors.EVENT_READ, data=accept_handler)
# 打印多路复用后端驱动类型,亲眼见证内核大杀器的激活
print(f"[系统] 当前底层多路复用驱动为: {sel}")
print("[系统] 商用级防粘包并发服务器启动,正在监听 8888 端口...")
try:
while True:
# 关键点:进入大循环。没有任何事件发生时,线程在内核层深度休眠,CPU 占用率为 0%
# 一旦大门被敲响(新连接)或专属座位有动静(发数据),立刻被内核一巴掌拍醒并返回
events = sel.select()
for key, mask in events:
# 动态分流:通过判断触发事件的 fileobj 是大门还是客桌,来采取不同的解包策略
if key.fileobj is server:
callback = key.data
callback(key.fileobj, mask)
else:
# 如果是老客桌,解包出(处理函数, 属于它的专属蓄水池)
callback, state = key.data
callback(key.fileobj, mask, state)
except KeyboardInterrupt:
print("\n[系统] 收到关闭指令,服务器正在安全退出...")
finally:
sel.close()
if __name__ == "__main__":
start_server()







赣ICP备2025054460号-1