当前位置: 云端笔记 » 编程 » Python » 【Python】socket编程

【Python】socket编程

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()
未经允许不得转载:云端笔记 » 【Python】socket编程

相关文章

评论 (0)

5 + 5 =

contact