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

基于poll与线程池的C多客户端服务器程序设计疑问

TCP消息拆分下的线程池任务去重与并发处理方案

问题背景

我用C语言开发单服务器多客户端程序,客户端提交的操作请求可能耗时,因此采用线程池并发处理以避免阻塞主线程。由于要求POSIX兼容,无法使用epoll,故选择poll优化性能,规避每连接建线程的C10K问题。

当前服务器核心逻辑伪代码如下:

int main()
{
  // 假设这些变量已完成初始化
  ThreadPool thread_pool;
  Socket server_listening_socket;
  PolledFileDescriptors list_of_polled_fds;

  // 第一个pollfd对应监听套接字,监听读事件
  list_of_polled_fds[0].fd = server_listening_socket;
  list_of_polled_fds[0].events = POLLIN;

  while (true)
  {
    // 调用poll,无限等待事件
    poll(&list_of_polled_fds, number_of_fds, -1);
    for (int i = 0; i < number_of_fds; i++)
    {
      // 当前文件描述符触发读事件
      if (list_of_polled_fds[i].revents & POLLIN)
      {
 
        // 监听套接字触发事件:新客户端连接
        if (i == 0)
        {
          Socket client_socket = accept();
          AddClientConnectionToListOfPollFds(&list_of_polled_fds, client_socket);
        }

        // 已连接客户端触发事件:有数据发送过来
        else
        {  
          ThreadPoolTask task = {
            .argument = list_of_polled_fds[i].fd, // 客户端套接字
            .function = SomeFunctionToReadDataFromSocketAndProcessIt
          };
          AddTaskToThreadPool(&thread_pool, &task);
        }
      }
    }
  }

  return 0;
}

遇到的核心问题:如果客户端的10字节消息被TCP拆分成两个包,每个包都会触发poll的POLLIN事件,导致线程池出现两个针对同一套接字的同一条消息处理任务。但如果直接禁止同一套接字重复添加任务,又会影响同一客户端发送的独立请求的并行处理。需要找到方法识别同属单条消息的多个poll事件,既能避免冗余任务,又能支持同一客户端独立请求的并行处理。


解决方案

1. 为每个客户端连接维护独立状态机与消息缓冲区

给每个客户端套接字绑定一个专属状态结构,用来跟踪消息接收进度与处理状态:

typedef struct ClientConnState {
  int fd;                          // 客户端套接字
  char buffer[4096];               // 缓存未完成的消息片段
  size_t buffer_len;               // 当前缓冲区已存储的数据长度
  bool is_assembling_msg;          // 标记是否正在拼接未完成的消息
  pthread_mutex_t state_mutex;     // 保护状态与缓冲区的互斥锁
} ClientConnState;

当poll触发POLLIN事件时,按以下逻辑处理:

  • 先锁定该连接的状态锁,避免主线程与线程池任务的状态冲突
  • 如果未在拼接消息:读取数据到缓冲区,尝试解析完整消息。若解析成功,将完整消息封装为任务提交线程池;若数据仍不完整,标记is_assembling_msg为true,等待后续数据补充
  • 如果正在拼接消息:直接读取数据追加到缓冲区,再次尝试解析完整消息,成功后提交任务并重置状态

2. 定义应用层消息边界

必须在自定义协议中明确消息的分隔规则,这是区分单条消息与独立请求的核心:

  • 固定长度消息头:消息开头用固定字节数存储整个消息的总长度,读取时先解析长度,再累计读取对应字节数的数据,直到凑齐完整消息
  • 特殊分隔符:用特定字节序列(如\r\n、自定义标识)作为消息结束标志,读取时扫描缓冲区找到分隔符,截取完整消息

3. 线程池任务完成后重置连接状态

线程池中的任务处理完完整消息后,需要主动重置对应连接的is_assembling_msg状态(如果缓冲区已无剩余数据),确保该客户端后续发送的独立请求能正常触发新的线程池任务,实现并行处理。

4. 优化锁粒度避免性能损耗

给每个客户端的状态结构单独加互斥锁,而非使用全局锁。这样只有同一连接的状态操作会触发锁竞争,不会影响其他连接的并发处理,平衡线程安全与性能。


调整后的核心处理逻辑示例

// 已连接客户端触发事件:有数据发送过来
else
{  
  ClientConnState *state = GetClientStateByFd(list_of_polled_fds[i].fd);
  pthread_mutex_lock(&state->state_mutex);
  
  ssize_t read_len = read(state->fd, state->buffer + state->buffer_len, sizeof(state->buffer) - state->buffer_len);
  if (read_len <= 0) {
    // 处理连接关闭或错误
    pthread_mutex_unlock(&state->state_mutex);
    RemoveClientFromPollFds(state->fd);
    pthread_mutex_destroy(&state->state_mutex);
    free(state);
    continue;
  }
  state->buffer_len += read_len;

  // 循环解析缓冲区中的所有完整消息
  while (HasCompleteMessage(state->buffer, state->buffer_len)) {
    // 提取完整消息并复制,避免缓冲区被后续数据覆盖
    Message *complete_msg = ExtractCompleteMessage(state->buffer, state->buffer_len);
    // 封装任务,携带消息与连接状态
    ThreadPoolTask task = {
      .argument = complete_msg,
      .function = ProcessCompleteClientRequest,
      .conn_state = state
    };
    AddTaskToThreadPool(&thread_pool, &task);
    // 移除已处理的消息数据,更新缓冲区
    RemoveProcessedData(state->buffer, &state->buffer_len, complete_msg->total_len);
  }

  // 根据缓冲区剩余数据更新状态
  state->is_assembling_msg = (state->buffer_len > 0);
  
  pthread_mutex_unlock(&state->state_mutex);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 05:33:17