欢迎光临
我们一直在努力

13.实现一个 处理粘包和拆包,并引入心跳机制的 通信

文章目录

      • TCP 粘包与拆包现象的本质:字节流边界缺失
        • 1. 粘包(Sticky Packet)产生的底层诱因
        • 2. 拆包(Packet Fragmentation)产生的必然因素
        • 3. 深层逻辑汇总
  • 基于 C++ 重构的 TCP 长连接服务器:粘包/拆包解决方案与心跳机制设计
    • 一、前言
    • 二、TCP 粘包与拆包问题详解
      • 2.1 什么是粘包和拆包?
      • 2.2 产生原因
      • 2.3 常见的解决方案
    • 三、本项目协议设计:如何解决粘包问题
      • 3.1 自定义数据包格式
      • 3.2 为什么能解决粘包?
      • 3.3 核心实现:readn 精确读取
    • 四、心跳机制设计
      • 4.1 为什么需要心跳?
      • 4.2 本项目的双向心跳设计
      • 4.3 心跳交互时序图
    • 五、C++ 重构亮点
      • 5.1 线程安全的客户端管理器
      • 5.2 RAII 锁管理
    • 六、完整代码
      • 6.1 项目结构
      • 6.2 `socket.h`
      • 6.3 `socket.c`
      • 6.4 `clientlist.hpp`
      • 6.5 `server.cpp`
      • 6.6 `client.cpp`
      • 6.7 `Makefile`
    • 七、编译与运行
    • 八、数据流全景图
    • 九、总结

TCP 粘包与拆包现象的本质:字节流边界缺失

在 TCP/IP 协议栈中,TCP 是一种面向字节流(Stream-Oriented)的传输层协议。它维护着一个独立的发送缓冲区(Send Buffer)和接收缓冲区(Receive Buffer),数据在缓冲区中以二进制字节流的形式存在。应用层交付给 TCP 的数据会被拷贝进该缓冲区,而 TCP 协议对该缓冲区内的数据不感知应用层消息边界(Message Boundary)。正是这种“流式”特性,使得接收端的应用程序在调用 read() 或 recv() 系统调用时,无法天然区分哪几个字节属于同一条业务消息,从而衍生出粘包(Sticky Packet)与拆包(Packet Fragmentation)问题。

1. 粘包(Sticky Packet)产生的底层诱因

粘包特指接收端的一次读操作返回了包含两条及以上应用层消息的数据段。其成因可从发送端和接收端两个维度剖析:

  • 发送端侧:Nagle 算法导致的“小包合并” TCP 默认开启 Nagle 算法,旨在优化广域网中的小数据包传输,避免网络中充斥大量仅有几字节有效载荷的报文(Silly Window Syndrome)。该算法的核心逻辑是:只要发送缓冲区中还有未被确认(Unacknowledged)的数据,新生成的小于 MSS(最大分段大小)的数据包就不会立即发送,而是暂存于缓冲区中等待合并。当业务层高频发送多条短消息时,Nagle 算法将这些消息“攒”成一个足够大的 TCP 段发往对端,导致接收端一次性收到多消息合并的大段数据,即产生粘包。

  • 接收端侧:应用层读取滞后与缓冲区积压 接收端的 TCP 接收缓冲区会持续积累从网络到来的数据。若应用层处理数据的速度(即系统调用的调用频率)低于数据到达网卡的速度,缓冲区内的数据将发生积压。当应用程序下一次调用 recv 且指定的用户态缓冲区足够大时,内核协议栈会将当前接收缓冲区中已排队的连续字节流一次性拷贝至用户态,从而造成多条业务消息在用户态被“粘合”读取。

2. 拆包(Packet Fragmentation)产生的必然因素

拆包指一条完整的应用层消息被截断,分散在多次读操作中才被完整接收。其根本原因在于链路层与传输层的最大载荷长度限制。

  • 传输层硬限制:MSS(最大分段大小)的强制切片 在 TCP 三次握手期间,双方会交换 MSS 选项(通常取值 = MTU – 40 字节,标准以太网下为 1460 字节)。当应用层单条消息长度超过 MSS 时,TCP 发送端的分段逻辑(Segmentation)会强制将该数据块拆分为多个 TCP 段(Segment),并在 IP 层分别封装后发送。这在接收端表现为:接收方需要调用多次 read 才能拼凑出完整的应用层原始数据。

  • 误区澄清:IP 分片(IP Fragmentation)并非主因 业内常有误解将拆包归咎于 IP 层的分片重组。事实上,TCP 设计之初就极力避免 IP 分片。因为 IP 分片缺失超时重传机制(丢失一片需重传整个 IP 数据报),开销极高。通过协商 MSS 并使每个 TCP 段长度小于 MTU,TCP 将分片行为上移至传输层自身完成,从而规避了网络层的 IP 分片。因此,TCP 拆包的核心根源在于 MSS 对报文大小的强约束。

  • 接收窗口(Receive Window)动态滑动的影响 接收方通告的窗口大小(Window Size)动态变化。若当前可用窗口小于单条消息的长度,发送方必须根据窗口值将消息切割成多个符合窗口大小的报文段发送,这也会导致接收端逻辑上的拆包现象。

3. 深层逻辑汇总

将上述成因映射至 TCP 内核状态机,可总结如下:

  • 粘包 = 发送缓冲区的延迟合并(Nagle)与 接收缓冲区的读取时机(Delay Read)耦合作用的结果。
  • 拆包 = **路径 MTU 发现(PMTU)**与 MSS 协商值对上层数据的硬性切割。

技术结论:TCP 协议保证了字节流的有序性和可靠性,但不保证应用层消息边界的完整性。因此,粘包与拆包是 TCP 流式特性的必然副产物,而非协议异常。解决此问题的唯一途径是在应用层设计协议栈(Protocol Stack)时引入明确的定界符(Delimiter)、**固定长度(Fixed Length)或长度字段(Length Field)**机制,以此在无序的字节流中重建消息帧(Frame)边界。


基于 C++ 重构的 TCP 长连接服务器:粘包/拆包解决方案与心跳机制设计


一、前言

在实际的网络编程中,TCP 长连接服务面临两大核心挑战:

  • 粘包/拆包问题 —— TCP 是面向流的协议,不保证消息边界
  • 连接活性检测 —— 如何及时发现"半开连接"(对端异常断开但本端不感知)
  • 本项目通过自定义应用层协议(长度字段 + 类型字段)彻底解决粘包问题,并引入双向心跳机制实现连接的活性检测。本文将完整剖析其设计原理与 C++ 实现。


    二、TCP 粘包与拆包问题详解

    2.1 什么是粘包和拆包?

    TCP 是面向字节流的传输协议,它不像 UDP 那样有消息边界的概念。对于 TCP 而言,数据就像一条连续的水流,没有"一条消息"和"另一条消息"的区分。

    这就导致了两种经典问题:

    发送方连续发送两条消息:
    ┌──────────┐ ┌──────────────┐
    │ "Hello" │ │ "World!" │
    │ (5字节) │ │ (6字节) │
    └──────────┘ └──────────────┘

    接收方可能的读取结果:

    【粘包】一次性读到两条消息粘在一起:
    ┌────────────────────┐
    │ "HelloWorld!" │ → 接收方不知道这是两条消息
    │ (11字节一次到达) │
    └────────────────────┘

    【拆包】一条消息被拆成多次到达:
    ┌──────┐ ┌─────────────────┐
    │ "Hel"│ │ "loWorld!" │ → 接收方不知道数据在哪断开
    │(3字节)│ │ (8字节) │
    └──────┘ └─────────────────┘

    2.2 产生原因

    粘包原因拆包原因
    发送方 Nagle 算法将多个小包合并发送 数据大小超过 MSS,TCP 层自动分片
    接收方读取缓冲区不及时 接收方缓冲区不够,只读取了部分数据
    多条短消息间隔时间极短 网络 MTU 限制

    2.3 常见的解决方案

    方案原理优缺点
    固定长度 每条消息固定 N 字节 简单但浪费带宽,不够灵活
    特殊分隔符 如 HTTP 用 \\r\\n 分隔 简单,但数据本身可能包含分隔符
    长度字段(本项目采用) 在消息头部携带数据长度 ✅ 灵活高效,业界主流方案

    三、本项目协议设计:如何解决粘包问题

    3.1 自定义数据包格式

    本项目在应用层定义了如下数据包格式:

    ┌──────────────────┬──────────────────┬──────────────────┐
    │ 数据长度 │ 数据包类型 │ 数据块 │
    │ (4 字节) │ (1 字节) │ (N 字节) │
    │ 网络字节序 │ H=心跳 M=消息 │ 实际负载数据 │
    └──────────────────┴──────────────────┴──────────────────┘
    ↑ ↑ ↑
    告诉接收方 告诉接收方 真正的业务数据
    后面有多少字节 这是什么类型的数据

    • 数据长度字段(4字节):存储的是 类型字段(1字节) + 数据块(N字节) 的总长度,采用网络字节序(大端)
    • 类型字段(1字节):'H' 表示心跳包,'M' 表示业务消息
    • 数据块(N字节):实际传输的数据内容

    3.2 为什么能解决粘包?

    关键在于 recvMessage() 的两步读取策略:

    接收流程:
    Step 1: 精确读取 4 字节 → 得到 dataLen(本次消息的总长度)
    Step 2: 精确读取 1 字节 → 得到消息类型
    Step 3: 精确读取 (dataLen – 1) 字节 → 得到完整数据块

    TCP 字节流
    ┌──────────────────────────────────────────────────┐
    │ [len1][type1][data1…] [len2][type2][data2…] │
    └──────────────────────────────────────────────────┘
    ↑ ↑
    第一条消息 第二条消息
    读 len1 字节后停止 ← 通过长度精确划分边界

    无论底层 TCP 如何粘包或拆包,应用层始终按 先读长度 → 再读指定字节数 的方式工作:

    • 遇到粘包:读完第一条消息的指定长度后,剩余数据留在缓冲区,下次循环继续读 → ✅ 正确拆分
    • 遇到拆包:readn() 函数会循环读取直到凑够指定字节数才返回 → ✅ 正确拼接

    3.3 核心实现:readn 精确读取

    int readn(int fd, char* buffer, int size)
    {
    int left = size; // 还需要读取的字节数
    int readBytes = 0;
    char* ptr = buffer;

    while (left > 0) {
    readBytes = read(fd, ptr, left);
    if (readBytes == 1) {
    if (errno == EINTR) {
    continue; // 被信号中断,重试
    } else {
    perror("read");
    return 1;
    }
    } else if (readBytes == 0) {
    // 对端关闭连接
    printf("对方主动断开了连接…\\n");
    return 1;
    }
    left -= readBytes; // 更新剩余字节数
    ptr += readBytes; // 移动写入指针
    }
    return size left; // 返回实际读到的总字节数
    }

    readn() 是解决拆包的核心——它保证读到恰好 size 字节才返回。即使 TCP 一次只给了一部分数据,它也会继续读,直到凑满。


    四、心跳机制设计

    4.1 为什么需要心跳?

    TCP 连接存在一种隐蔽的故障:半开连接(Half-Open Connection)。

    当一端异常崩溃(如断电、网络中断)而没有正常发送 FIN 包时,另一端并不知道连接已经失效,会一直维护这个"僵尸连接",浪费资源。

    正常断开:
    Client ──── FIN ────→ Server (四次挥手,双方都知道连接结束了)

    异常断开(半开连接):
    Client ✕(突然断电)
    Server 还在傻等… 以为连接仍然有效

    心跳机制就是定期互相"打招呼"来确认对方还活着。

    4.2 本项目的双向心跳设计

    时间轴 ──────────────────────────────────────────────────────→

    客户端(每3秒):
    ├── 发送心跳包 "ping" ──→
    ├── count++
    ├── 收到回复 → count = 0
    ├── count > 5 → 判定服务器断开,退出

    服务端(每3秒巡检):
    ├── 遍历所有客户端
    ├── 每个客户端 count++
    ├── 收到心跳 → count = 0
    ├── count > 5 → 判定客户端断开,移除

    双向检测的意义:

    • 服务端巡检:及时发现并清理死连接,释放资源
    • 客户端心跳:及时发现服务器不可达,可以重连或提示用户

    4.3 心跳交互时序图

    Client Server
    │ │
    │──── Heart("ping") ────────────────→│ 客户端发送心跳
    │ │ count 重置为 0
    │←─── Heart("ping") 回复 ────────────│ 服务端回复心跳
    │ │
    │ count 重置为 0 │
    │ │
    │ … 3秒后 … │
    │ │
    │──── Heart("ping") ────────────────→│
    │ │
    │ (假设此时网络中断) │
    │ │
    │ count=1, count=2, … count=6 │ 服务端 count > 5
    │ → 判定服务器断开,退出 │ → 移除客户端,关闭连接


    五、C++ 重构亮点

    5.1 线程安全的客户端管理器

    原 C 版本使用链表 + 全局互斥锁管理客户端,重构后使用 C++ 的 ClientManager 类封装:

    class ClientManager {
    std::unordered_map<int, ClientInfo> m_clients; // O(1) 查找
    mutable std::mutex m_mutex; // 内置锁保护
    };

    优势:

    特性C 版本(链表)C++ 版本(unordered_map)
    查找复杂度 O(n) O(1)
    线程安全 需手动加锁/解锁 类内部封装,RAII 自动管理
    内存管理 手动 malloc/free 自动管理
    遍历操作 手动遍历链表 forEach() + lambda 回调

    5.2 RAII 锁管理

    使用 std::lock_guard 确保异常安全,不会出现忘记解锁的问题:

    void addClient(int fd) {
    std::lock_guard<std::mutex> lock(m_mutex); // 构造时加锁
    // … 操作 …
    // 函数返回时自动解锁(即使异常退出)
    }


    六、完整代码

    6.1 项目结构

    project/
    ├── socket.h (C 语言网络基础库头文件)
    ├── socket.c (C 语言网络基础库实现)
    ├── clientlist.hpp (C++ 客户端管理器)
    ├── server.cpp (服务端主程序)
    ├── client.cpp (客户端主程序)
    └── Makefile (编译脚本)


    6.2 socket.h

    #pragma once
    #include <stdbool.h>
    #include <arpa/inet.h>

    // 数据包类型
    enum Type { Heart, Message };
    // 数据包格式: 数据长度|数据包类型|数据块
    // int char char*
    // 4字节 1字节 N字节

    // 初始化一个套接字
    int initSocket(void);
    // 设置监听
    int setListen(int lfd, unsigned port);
    // 接收客户端连接
    int acceptConnect(int lfd, struct sockaddr* addr);
    // 连接服务器
    int connectToHost(int fd, unsigned port, const char* ip);
    // 读出指定的字节数
    int readn(int fd, char* buffer, int size);
    // 写入指定的字节数
    int writen(int fd, const char* buffer, int length);
    // 发送数据
    bool sendMessage(int fd, const char* buffer, int length, enum Type t);
    // 接收数据
    int recvMessage(int fd, char** buffer, enum Type* t);


    6.3 socket.c

    #include "socket.h"
    #include <stdio.h>
    #include <errno.h>
    #include <string.h>
    #include <stdlib.h>
    #include <unistd.h>

    int initSocket(void)
    {
    int lfd = socket(AF_INET, SOCK_STREAM, 0);
    if (lfd == 1) {
    perror("socket");
    return 1;
    }
    return lfd;
    }

    static void initSockaddr(struct sockaddr* addr, unsigned port, const char* ip)
    {
    struct sockaddr_in* addrin = (struct sockaddr_in*)addr;
    addrin->sin_family = AF_INET;
    addrin->sin_port = htons(port);
    addrin->sin_addr.s_addr = inet_addr(ip);
    }

    int setListen(int lfd, unsigned port)
    {
    struct sockaddr addr;
    initSockaddr(&addr, port, "0.0.0.0");

    // 端口复用
    int opt = 1;
    setsockopt(lfd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt));

    int ret = bind(lfd, &addr, sizeof(addr));
    if (ret == 1) {
    perror("bind");
    return 1;
    }

    ret = listen(lfd, 128);
    if (ret == 1) {
    perror("listen");
    return 1;
    }
    return 0;
    }

    int acceptConnect(int lfd, struct sockaddr* addr)
    {
    int connfd;
    if (addr == NULL) {
    connfd = accept(lfd, NULL, NULL);
    } else {
    socklen_t len = sizeof(struct sockaddr);
    connfd = accept(lfd, addr, &len);
    }
    if (connfd == 1) {
    perror("accept");
    return 1;
    }
    return connfd;
    }

    int connectToHost(int fd, unsigned port, const char* ip)
    {
    struct sockaddr addr;
    initSockaddr(&addr, port, ip);
    int ret = connect(fd, &addr, sizeof(addr));
    if (ret == 1) {
    perror("connect");
    return 1;
    }
    return 0;
    }

    int readn(int fd, char* buffer, int size)
    {
    int left = size;
    int readBytes = 0;
    char* ptr = buffer;

    while (left > 0) {
    readBytes = read(fd, ptr, left);
    if (readBytes == 1) {
    if (errno == EINTR) {
    continue;
    } else {
    perror("read");
    return 1;
    }
    } else if (readBytes == 0) {
    printf("对方主动断开了连接…\\n");
    return 1;
    }
    left -= readBytes;
    ptr += readBytes;
    }
    return size left;
    }

    int writen(int fd, const char* buffer, int length)
    {
    int left = length;
    int writeBytes = 0;
    const char* ptr = buffer;

    while (left > 0) {
    writeBytes = write(fd, ptr, left);
    if (writeBytes == 1) {
    if (errno == EINTR) {
    continue;
    } else {
    perror("write");
    return 1;
    }
    }
    ptr += writeBytes;
    left -= writeBytes;
    }
    return length;
    }

    bool sendMessage(int fd, const char* buffer, int length, enum Type t)
    {
    // 数据总长度 = 长度字段(4字节) + 类型(1字节) + 数据块(length字节)
    int dataLen = length + 1 + sizeof(int);

    char* data = (char*)malloc(dataLen);
    if (data == NULL) {
    return false;
    }

    // 1. 写入长度字段(网络字节序)
    int netlen = htonl(length + 1);
    memcpy(data, &netlen, sizeof(int));

    // 2. 写入类型字段
    char ch = (t == Heart) ? 'H' : 'M';
    memcpy(data + sizeof(int), &ch, sizeof(char));

    // 3. 写入数据块
    memcpy(data + sizeof(int) + 1, buffer, length);

    int ret = writen(fd, data, dataLen);
    free(data);

    return ret == dataLen;
    }

    int recvMessage(int fd, char** buffer, enum Type* t)
    {
    // 1. 读取长度字段(4字节)
    int dataLen = 0;
    int ret = readn(fd, (char*)&dataLen, sizeof(int));
    if (ret == 1) {
    *buffer = NULL;
    return 1;
    }
    dataLen = ntohl(dataLen);

    // 2. 读取类型字段(1字节)
    char ch;
    ret = readn(fd, &ch, 1);
    if (ret == 1) {
    *buffer = NULL;
    return 1;
    }
    *t = (ch == 'H') ? Heart : Message;

    // 3. 读取数据块(dataLen – 1 字节)
    char* tmpbuf = (char*)calloc(dataLen, sizeof(char));
    if (tmpbuf == NULL) {
    *buffer = NULL;
    return 1;
    }

    ret = readn(fd, tmpbuf, dataLen 1);
    if (ret != dataLen 1) {
    free(tmpbuf);
    *buffer = NULL;
    return 1;
    }

    *buffer = tmpbuf;
    return ret;
    }


    6.4 clientlist.hpp

    #pragma once
    #include <unordered_map>
    #include <mutex>
    #include <pthread.h>
    #include <cstdio>

    // 客户端信息结构体
    struct ClientInfo {
    int fd; // 通信套接字
    int count; // 心跳超时计数
    pthread_t pid; // 接收线程 ID
    };

    // 客户端管理器(封装所有操作,线程安全)
    class ClientManager {
    public:
    ClientManager() = default;
    ~ClientManager() = default;

    // 添加客户端
    void addClient(int fd) {
    std::lock_guard<std::mutex> lock(m_mutex);
    ClientInfo info;
    info.fd = fd;
    info.count = 0;
    info.pid = 0;
    m_clients[fd] = info;
    printf("[Manager] 添加客户端 fd=%d, 当前在线: %zu\\n", fd, m_clients.size());
    }

    // 删除客户端
    bool removeClient(int fd) {
    std::lock_guard<std::mutex> lock(m_mutex);
    auto it = m_clients.find(fd);
    if (it == m_clients.end()) {
    return false;
    }
    // 取消对应的接收线程
    if (it->second.pid != 0) {
    pthread_cancel(it->second.pid);
    }
    m_clients.erase(it);
    printf("[Manager] 删除客户端 fd=%d, 当前在线: %zu\\n", fd, m_clients.size());
    return true;
    }

    // 获取客户端信息(只读指针)
    ClientInfo* getClient(int fd) {
    std::lock_guard<std::mutex> lock(m_mutex);
    auto it = m_clients.find(fd);
    if (it == m_clients.end()) {
    return nullptr;
    }
    return &(it->second);
    }

    // 重置客户端的 count(收到心跳回复时调用)
    void resetCount(int fd) {
    std::lock_guard<std::mutex> lock(m_mutex);
    auto it = m_clients.find(fd);
    if (it != m_clients.end()) {
    it->second.count = 0;
    }
    }

    // 递增客户端的 count
    void incrementCount(int fd) {
    std::lock_guard<std::mutex> lock(m_mutex);
    auto it = m_clients.find(fd);
    if (it != m_clients.end()) {
    it->second.count++;
    }
    }

    // 遍历所有客户端(执行回调函数)
    template<typename Func>
    void forEach(Func func) {
    std::lock_guard<std::mutex> lock(m_mutex);
    for (auto& pair : m_clients) {
    func(pair.second);
    }
    }

    // 获取在线客户端数量
    size_t size() const {
    std::lock_guard<std::mutex> lock(m_mutex);
    return m_clients.size();
    }

    // 检查是否存在某个 fd
    bool exists(int fd) const {
    std::lock_guard<std::mutex> lock(m_mutex);
    return m_clients.find(fd) != m_clients.end();
    }

    private:
    std::unordered_map<int, ClientInfo> m_clients; // key: fd, value: ClientInfo
    mutable std::mutex m_mutex; // 保护 m_clients 的互斥锁
    };


    6.5 server.cpp

    #include <stdio.h>
    #include <string.h>
    #include <pthread.h>
    #include <stdlib.h>
    #include <unistd.h>
    #include "socket.h"
    #include "clientlist.hpp"

    // 全局客户端管理器
    ClientManager g_clientManager;

    // 接收消息的线程函数
    void* parseRecvMessage(void* arg) {
    int fd = *(int*)arg;
    free(arg); // 释放传入的 fd 内存

    printf("[Thread] 接收线程启动, fd=%d\\n", fd);

    while (1) {
    char* buffer = nullptr;
    enum Type t;
    int len = recvMessage(fd, &buffer, &t);

    if (buffer == NULL) {
    printf("[Thread] fd=%d 接收失败或连接关闭, 线程退出\\n", fd);
    g_clientManager.removeClient(fd);
    pthread_exit(NULL);
    } else {
    if (t == Heart) {
    printf("[Heart] fd=%d 收到心跳包: %s\\n", fd, buffer);
    g_clientManager.resetCount(fd); // 清零超时计数
    sendMessage(fd, buffer, len, Heart); // 回复心跳
    } else {
    printf("[Message] fd=%d 收到数据: %s\\n", fd, buffer);
    const char* reply = "愿世界和平…";
    sendMessage(fd, reply, strlen(reply), Message);
    }
    free(buffer);
    }
    }
    return NULL;
    }

    // 心跳巡检线程(每隔3秒检查所有客户端)
    void* heartBeat(void* arg) {
    printf("[HeartBeat] 心跳巡检线程启动\\n");

    while (1) {
    // 遍历所有客户端
    g_clientManager.forEach([&](ClientInfo& info) {
    info.count++;
    printf("[HeartBeat] fd=%d, count=%d\\n", info.fd, info.count);

    if (info.count > 5) {
    printf("[HeartBeat] 客户端 fd=%d 超时(count=%d),断开连接\\n",
    info.fd, info.count);
    close(info.fd);
    // 从管理器中移除(会自动取消线程)
    g_clientManager.removeClient(info.fd);
    }
    });

    sleep(3);
    }
    return NULL;
    }

    int main() {
    printf("=== 服务器启动 ===\\n");

    unsigned short port = 8888;
    int lfd = initSocket();
    if (lfd == 1) {
    printf("初始化套接字失败\\n");
    return 1;
    }

    if (setListen(lfd, port) == 1) {
    printf("设置监听失败\\n");
    close(lfd);
    return 1;
    }

    printf("服务器监听端口: %d\\n", port);

    // 创建心跳巡检线程
    pthread_t pid1;
    pthread_create(&pid1, NULL, heartBeat, NULL);
    pthread_detach(pid1); // 分离线程,自动回收资源

    // 主循环:接受客户端连接
    while (1) {
    int sockfd = acceptConnect(lfd, NULL);
    if (sockfd == 1) {
    continue;
    }

    printf("[Main] 新客户端连接, fd=%d\\n", sockfd);

    // 添加到客户端管理器
    g_clientManager.addClient(sockfd);

    // 传递 sockfd 给线程(在堆上分配,避免栈悬空)
    int* fdPtr = (int*)malloc(sizeof(int));
    *fdPtr = sockfd;

    // 创建接收线程
    ClientInfo* info = g_clientManager.getClient(sockfd);
    if (info) {
    pthread_create(&info->pid, NULL, parseRecvMessage, fdPtr);
    pthread_detach(info->pid);
    printf("[Main] 为 fd=%d 创建接收线程, tid=%lu\\n",
    sockfd, (unsigned long)info->pid);
    } else {
    free(fdPtr);
    close(sockfd);
    }
    }

    close(lfd);
    return 0;
    }


    6.6 client.cpp

    #include <stdio.h>
    #include <string.h>
    #include <pthread.h>
    #include <stdlib.h>
    #include <unistd.h>
    #include "socket.h"

    // 客户端信息结构体
    struct ClientInfo {
    int fd;
    int count;
    pthread_mutex_t mutex; // 每个客户端独立的锁
    };

    // 接收消息线程
    void* parseRecvMessage(void* arg) {
    ClientInfo* info = (ClientInfo*)arg;
    printf("[Client] 接收线程启动\\n");

    while (1) {
    char* buffer = nullptr;
    enum Type t;
    int ret = recvMessage(info->fd, &buffer, &t);

    if (buffer == NULL) {
    printf("[Client] 接收失败,线程退出\\n");
    pthread_exit(NULL);
    }

    if (t == Heart) {
    printf("[Heart] 收到心跳回复: %s\\n", buffer);
    pthread_mutex_lock(&info->mutex);
    info->count = 0; // 清零超时计数
    pthread_mutex_unlock(&info->mutex);
    } else {
    printf("[Message] 收到服务器数据: %s\\n", buffer);
    }
    free(buffer);
    }
    return NULL;
    }

    // 心跳发送和检测线程
    void* heartBeat(void* arg) {
    ClientInfo* info = (ClientInfo*)arg;
    printf("[Client] 心跳线程启动\\n");

    while (1) {
    pthread_mutex_lock(&info->mutex);
    info->count++;
    printf("[Heart] fd=%d, count=%d\\n", info->fd, info->count);

    if (info->count > 5) {
    printf("[Heart] 连续5次未收到心跳回复,断开连接\\n");
    close(info->fd);
    pthread_mutex_unlock(&info->mutex);
    exit(0); // 退出整个程序
    }
    pthread_mutex_unlock(&info->mutex);

    // 发送心跳包
    sendMessage(info->fd, "ping", 4, Heart);
    sleep(3);
    }
    return NULL;
    }

    int main() {
    printf("=== 客户端启动 ===\\n");

    // 在堆上分配 ClientInfo,避免悬空指针
    ClientInfo* info = new ClientInfo();
    info->fd = initSocket();
    info->count = 0;
    pthread_mutex_init(&info->mutex, NULL);

    if (info->fd == 1) {
    printf("初始化套接字失败\\n");
    delete info;
    return 1;
    }

    unsigned short port = 8888;
    const char* ip = "127.0.0.1";

    if (connectToHost(info->fd, port, ip) == 1) {
    printf("连接服务器失败\\n");
    close(info->fd);
    delete info;
    return 1;
    }

    printf("连接服务器 %s:%d 成功\\n", ip, port);

    // 创建接收线程
    pthread_t pid_recv;
    pthread_create(&pid_recv, NULL, parseRecvMessage, info);
    pthread_detach(pid_recv);

    // 创建心跳线程
    pthread_t pid_heart;
    pthread_create(&pid_heart, NULL, heartBeat, info);
    pthread_detach(pid_heart);

    // 主线程:定期发送业务数据
    int count = 0;
    while (1) {
    char data[64];
    snprintf(data, sizeof(data), "你好, 服务器! 第%d次", ++count);
    printf("[Main] 发送数据: %s\\n", data);
    sendMessage(info->fd, data, strlen(data), Message);
    sleep(5); // 每5秒发送一次
    }

    // 清理(实际上永远不会执行到这里)
    pthread_mutex_destroy(&info->mutex);
    close(info->fd);
    delete info;
    return 0;
    }


    6.7 Makefile

    CC = gcc
    CXX = g++
    CFLAGS = -Wall -O2
    CXXFLAGS = -Wall -O2 -std=c++11
    LDFLAGS = -lpthread

    # 可执行文件
    TARGETS = server client

    all: $(TARGETS)

    server: server.o socket.o
    $(CXX) -o $@ $^ $(LDFLAGS)

    client: client.o socket.o
    $(CXX) -o $@ $^ $(LDFLAGS)

    %.o: %.cpp
    $(CXX) $(CXXFLAGS) -c $< -o $@

    %.o: %.c
    $(CC) $(CFLAGS) -c $< -o $@

    clean:
    rm -f $(TARGETS) *.o

    .PHONY: all clean


    七、编译与运行

    # 编译
    make clean && make

    # 终端1:启动服务器
    ./server

    # 终端2:启动客户端
    ./client

    # 可同时启动多个客户端测试
    ./client
    ./client


    八、数据流全景图

    ┌─────────────────────────────────────────────────────────────────────┐
    │ 服务器 (Server) │
    │ │
    │ ┌─────────────┐ ┌─────────────────┐ ┌──────────────────────┐ │
    │ │ 主线程 │ │ 心跳巡检线程 │ │ 接收线程 (每客户端) │ │
    │ │ accept() │ │ 每3秒遍历 │ │ recvMessage() │ │
    │ │ 创建线程 │ │ count++ │ │ 解析类型 → 处理 │ │
    │ └──────┬──────┘ │ >5 → 断开 │ └──────────┬───────────┘ │
    │ │ └────────┬────────┘ │ │
    │ ▼ ▼ ▼ │
    │ ┌──────────────────────────────────────────────────────────────┐ │
    │ │ ClientManager (线程安全) │ │
    │ │ unordered_map<int, ClientInfo> + mutex │ │
    │ │ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐ │ │
    │ │ │fd=5 │ │fd=6 │ │fd=7 │ │ … │ │ │
    │ │ │count=0 │ │count=2 │ │count=0 │ │ │ │ │
    │ │ └────────┘ └────────┘ └────────┘ └────────┘ │ │
    │ └──────────────────────────────────────────────────────────────┘ │
    └─────────────────────────────────────────────────────────────────────┘
    ↕ TCP 连接
    ┌─────────────────────────────────────────────────────────────────────┐
    │ 客户端 (Client) │
    │ │
    │ ┌─────────────┐ ┌─────────────────┐ ┌──────────────────────┐ │
    │ │ 主线程 │ │ 心跳线程 │ │ 接收线程 │ │
    │ │ 发送业务数据 │ │ 每3秒: │ │ recvMessage() │ │
    │ │ 每5秒一次 │ │ count++ │ │ 处理回复 │ │
    │ │ │ │ 发送心跳 │ │ count 归零 │ │
    │ │ │ │ >5 → 退出 │ │ │ │
    │ └─────────────┘ └─────────────────┘ └──────────────────────┘ │
    └─────────────────────────────────────────────────────────────────────┘


    九、总结

    技术点解决方案
    粘包问题 消息头部携带长度字段,接收方按长度精确读取
    拆包问题 readn() 循环读取直到凑满指定字节数
    半开连接检测 双向心跳:客户端定期发送,服务端定期巡检
    多线程安全 ClientManager 类封装 + std::mutex + RAII 锁管理
    高效查找 unordered_map 哈希表 O(1) 查找
    资源管理 pthread_detach 分离线程,堆分配避免悬空指针

    本项目展示了一个生产级别 TCP 长连接服务的核心要素:可靠的消息边界划分和健壮的连接活性检测。这两者是构建即时通讯、IoT 设备管理、游戏服务器等实时系统的基础。

    赞(0)
    未经允许不得转载:171主机测评 » 13.实现一个 处理粘包和拆包,并引入心跳机制的 通信
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址