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)
问题分析与修复说明
核心问题
- 服务端线程判断逻辑不可靠:原代码用
threading.active_count() == 2判断最后一个线程,这个值包含主线程,实际所有客户端线程结束后活跃数应为1,且该判断无锁保护,存在竞态条件,导致计数错误。 - 客户端原子变量竞态:原客户端判断
nthreads == 1后未加锁,可能多个线程同时进入接收逻辑,或判断后nthreads被其他线程修改,导致最后一个线程无法正确接收计数。
修复措施
- 服务端维护连接计数器:新增
connection_count变量,用锁保护其增减和判断,确保只有最后一个结束的连接会发送总行数。 - 客户端加锁保护线程判断:在判断是否为最后一个线程及修改
nthreads时加互斥锁,避免竞态条件。
内容的提问来源于stack exchange,提问作者jaeger22
相关产品推荐
相关产品推荐

