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

Python Socket与C客户端通信:接收行数统计不一致问题求助

Python服务端与C客户端通信程序的计数不一致问题

C客户端从命令行接收N个txt文件,为每个文件创建一个线程,每个线程与服务端建立连接并发送文件的所有行。Python服务端为每个接收到的连接创建一个线程,每个线程负责接收数据,同时服务端需要统计所有连接接收的总行数。问题在于,当服务端向客户端返回接收的行数时,每次执行结果都不一致。

客户端代码(C)

#include <stdio.h>
#include <stdlib.h>
#include <sys/types.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <errno.h>
#include <arpa/inet.h>
#define _GNU_SOURCE   // 启用GNU扩展
#include <unistd.h>
#include <string.h>   // 字符串函数
#include <assert.h>
#include <pthread.h>
#include <stdatomic.h>

// 要连接的主机和端口
#define HOST "127.0.0.1"
#define PORT 5050
#define Max_sequence_length 2048

pthread_mutex_t lock;

// 读写socket的readn和writen函数
size_t readn(int fd, void *ptr, size_t n) {  
   size_t   nleft;
   ssize_t  nread;
 
   nleft = n;
   while (nleft > 0) {
     if((nread = read(fd, ptr, nleft)) < 0) {
        if (nleft == n) return -1; /* 错误,返回-1 */
        else break; /* 错误,返回已读取的字节数 */
     } else if (nread == 0) break; /* EOF */
     nleft -= nread;
     ptr   += nread;
   }
   return(n - nleft); /* 返回>=0 */
}
size_t writen(int fd, void *ptr, size_t n) {  
   size_t   nleft;
   ssize_t  nwritten;
 
   nleft = n;
   while (nleft > 0) {
     if((nwritten = write(fd, ptr, nleft)) < 0) {
        if (nleft == n) return -1; /* 错误,返回-1 */
        else break; /* 错误,返回已写入的字节数 */
     } else if (nwritten == 0) break; 
     nleft -= nwritten;
     ptr   += nwritten;
   }
   return(n - nleft); /* 返回>=0 */
}

// 线程管理中会增减的原子变量
atomic_int totale_seq = 0;
atomic_int nthreads = 0;

void *gestionefile(char* nomefile){

    FILE *f = fopen(nomefile , "r");
     if(f == NULL){
        printf("打开文件%s失败 \n, 错误代码:%d", nomefile, errno);
        exit(1);
    }

    // 创建socket
    int inet_socket;
    inet_socket = socket(AF_INET, SOCK_STREAM, 0);
    if(inet_socket < 0){
        perror("创建socket失败");
        exit(1);
    };

    // 指定socket的地址和端口
    struct sockaddr_in server_address;
    server_address.sin_family = AF_INET;
    server_address.sin_port = htons(PORT);
    server_address.sin_addr.s_addr = inet_addr(HOST);

    // 连接socket
    int stato_connessione = connect(inet_socket, (struct sockaddr *) &server_address, sizeof(server_address));
    if(stato_connessione == -1){
        perror("连接socket时出错");
        exit(1);
    }
    printf("[已连接] 成功连接到端口为%d的socket \n" , server_address.sin_port);

    // 发送字符"B"指定连接类型
    char buffer[2];
    char* word = "B";
    strcpy(buffer, word);
    writen(inet_socket, &buffer, 1);
   
    // 读取文件行
    char* linea = NULL;
    size_t size = 0;
        while(getline(&linea , &size, f) != -1){

        // 检查行长度是否超过最大值
        assert(strlen(linea) <= Max_sequence_length);

        // 发送下一行的长度
        short dim = strlen(linea);
        short tmp = htons(dim);
        writen(inet_socket, &tmp, 2);

        // 发送行到服务端
        char buffer2[dim];
        strcpy(buffer2 , linea);
        writen(inet_socket, &buffer2, dim);
        
        totale_seq++;
        }

    fclose(f);
    free(linea);

    // 当前线程已发送完文件所有行
    printf("Client2已发送完文件%s的所有行 \n", nomefile);

    // 发送长度为0的序列作为结束标记
    short len_0 = 0;
    short zero = htons(len_0);
    writen(inet_socket, &zero , 2);
    
    // 如果是最后一个线程,等待服务端返回本次会话接收的行数
    pthread_mutex_lock(&lock);
    int is_last = (nthreads == 1);
    if(is_last){  
        int tmp1;   
        readn(inet_socket, &tmp1, 4);
        int s_ricevute = ntohl(tmp1);
        printf("服务端返回接收了%d行\n" , s_ricevute);
        printf("我总共发送了%d行",totale_seq);
        assert(s_ricevute == totale_seq);
    }
    nthreads--;
    pthread_mutex_unlock(&lock);
    
    close(inet_socket);
    pthread_exit(NULL);
}



int main(int argc, char *argv[]){

    pthread_mutex_init(&lock, NULL);
    
    if(argc == 1){
        perror("错误: client1需要从命令行接收至少一个文件");
        exit(1);
    }

    // 创建与命令行传入文件数量相等的线程
    // 每个线程创建一个B类型连接并将整个文件发送到服务端
    pthread_t id[argc - 1];
    for(int i = 0 ; i < argc - 1 ; i++){
        pthread_create(&id[i], NULL, gestionefile, argv[i+1] );
        nthreads++;
    }

    for(int i= 0; i < argc -1 ; i++){
    pthread_join(id[i],NULL);
    }

    printf("我总共发送了%d行",totale_seq);
    
    pthread_mutex_destroy(&lock);
    return 0;
}

服务端代码(Python)

import socket, sys, struct, threading, math

Max_sequence_length = 2048
PORT = 5050
HOST = "127.0.0.1"
ADDR = (HOST , PORT)
seq_tot = 0
connection_count = 0
lock = threading.Lock()

def recv_all(conn,n):
  chunks = b''
  bytes_recd = 0
  while bytes_recd < n:
    chunk = conn.recv(min(n - bytes_recd, 1024))
    if len(chunk) == 0:
      raise RuntimeError("socket连接中断")
    chunks += chunk
    bytes_recd = bytes_recd + len(chunk)
  return chunks
 
def gestisci_client(conn,addr):
  global connection_count
  with lock:
    connection_count +=1
  try:
    with conn:
      print(f"来自{addr}的新连接")
      # 接收通信类型
      tipo = recv_all(conn,1).decode()
      if(tipo == "A"):
        print("A类型连接")
      elif(tipo == "B"):
          gestione_b(conn)
  finally:
    with lock:
      connection_count -=1
 
def gestione_b(conn):
  n_sequenze_b = 0
  global seq_tot, connection_count
  
  while(1):
      # 接收行的长度
      data = recv_all(conn,2)
      lunghezza = struct.unpack('!h',data)[0]
  
      if(lunghezza == 0 ):
      
        with lock:
          seq_tot += n_sequenze_b
          # 检查是否是最后一个连接
          is_last = (connection_count == 1)
          current_total = seq_tot
          if is_last:
            seq_tot = 0
        
        if is_last:
            print(f"\n 最后一个连接,返回总行数{current_total}\n")
            conn.sendall(struct.pack("!i",current_total))
        
        break
      
      # 接收行并增加接收的总行数
      data = recv_all(conn,lunghezza).decode()
      n_sequenze_b = n_sequenze_b + 1
      


# 创建服务端socket
with socket.socket(socket.AF_INET , socket.SOCK_STREAM) as server : 
  # 将socket绑定到地址
  try:
      print("[等待中] : 服务端正在等待连接")
      server.bind(ADDR)
      server.listen()
      while True:
          conn, addr = server.accept()    
          thread = threading.Thread(target = gestisci_client , args = (conn, addr)) 
          thread.start()
              
  except KeyboardInterrupt: # 收到中断信号时关闭服务端
    pass

  print("[服务端关闭]")
  server.shutdown(socket.SHUT_RDWR)

问题分析与修复说明

核心问题

  1. 服务端线程判断逻辑不可靠:原代码用threading.active_count() == 2判断最后一个线程,这个值包含主线程,实际所有客户端线程结束后活跃数应为1,且该判断无锁保护,存在竞态条件,导致计数错误。
  2. 客户端原子变量竞态:原客户端判断nthreads == 1后未加锁,可能多个线程同时进入接收逻辑,或判断后nthreads被其他线程修改,导致最后一个线程无法正确接收计数。

修复措施

  1. 服务端维护连接计数器:新增connection_count变量,用锁保护其增减和判断,确保只有最后一个结束的连接会发送总行数。
  2. 客户端加锁保护线程判断:在判断是否为最后一个线程及修改nthreads时加互斥锁,避免竞态条件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 16:24:57