Linux下C语言线程与匿名管道通信异常问题求助
问题描述
我编写的程序计划创建5个Worker进程,每个进程绑定一个匿名管道并监听读取管道信息;同时创建一个Dispatcher线程,向所有匿名管道写入信息,供Worker进程读取并打印,管道存储在dispatcher_pipes数组中。此前在main函数中直接向管道写入数据时,Worker进程能正常接收并打印,但将写入逻辑移至Dispatcher线程后功能失效。编译环境为Linux终端的gcc工具。
当前失效代码
#include <stdio.h> #include <stdlib.h> #include <string.h> #include <unistd.h> #include <pthread.h> typedef struct msgbuf { long mtype; // if 0 it did not pass in no process, 1 - worker, 2 - alert watcher char mtext[100]; }msgbuf; int n_workers = 5; pid_t *pid_workers; void* dispatcher(void *args) { int (*dispatcher_pipe)[2] = (int (*)[2])args; pthread_t id=pthread_self(); printf("Dispatcher thread started.\n"); msgbuf send_msg; send_msg.mtype = 0; strcpy(send_msg.mtext, "Message from Dispatcher"); for (int i = 0; i < n_workers; i++) { close(dispatcher_pipe[i][0]); // Close the read end of the pipe write(dispatcher_pipe[i][1], &send_msg, sizeof(msgbuf)); close(dispatcher_pipe[i][1]); // Close the write end of the pipe } usleep(10); // simulates the processing of the order printf("Thread %ld: Worker finishing\n",id); pthread_exit(NULL); } void Worker(int *worker_pipe){ printf("Worker Process (PID=%d): starting!\n",getpid()); //Getting message from Dispatcher close(worker_pipe[1]); // Close the write end of the pipe msgbuf rec_msg; read(worker_pipe[0], rec_msg.mtext, sizeof(msgbuf)); printf("Worker message received from Dispatcher: %s\n", rec_msg.mtext); close(worker_pipe[0]); // Close the read end of the pipe sleep(1); printf("\nWorker process %d has died.\n", getpid()); exit(0); } void main(){ pthread_t dispatcher_thread; int dispatcher_pipe[n_workers][2]; pid_workers= malloc(n_workers * sizeof(pid_t)); //Creates the unnamed pipes for dispatcher for(int i=0;i<n_workers;i++){ if(pipe(dispatcher_pipe[i]) == -1){ printf("Error creating dispatcher pipe n1\n"); exit(0); } } //Creating Worker Processes for(int i=0;i<n_workers;i++){ pid_workers[i] = fork(); if (pid_workers[i] == 0) { Worker(dispatcher_pipe[i]); } else if (pid_workers[i] > 0) { } else { perror("fork"); exit(EXIT_FAILURE); } } //Creating Threads pthread_create(&dispatcher_thread, NULL, dispatcher, dispatcher_pipe); //Wait for the threads to finish pthread_join(dispatcher_thread, NULL); }
之前正常工作的代码
#include <stdio.h> #include <stdlib.h> #include <string.h> #include <unistd.h> #include <pthread.h> typedef struct msgbuf { long mtype; char mtext[100]; }msgbuf; int n_workers = 5; pid_t *pid_workers; void* dispatcher(void *args) { int (*dispatcher_pipe)[2] = (int (*)[2])args; pthread_t id=pthread_self(); printf("Dispatcher thread started.\n"); for(int i=0;i<n_workers;i++){ close(dispatcher_pipe[i][0]); // Close the read end of the pipe msgbuf send_msg; send_msg.mtype = 1; snprintf(send_msg.mtext, sizeof(send_msg.mtext), "Message from SystemManager to Worker %d", i); write(dispatcher_pipe[i][1], send_msg.mtext, sizeof(send_msg.mtext)); close(dispatcher_pipe[i][1]); // Close the write end of the pipe } usleep(10); // simulates the processing of the order printf("Thread %ld: Worker finishing\n",id); pthread_exit(NULL); } void Worker(int *worker_pipe){ printf("Worker Process (PID=%d): starting!\n",getpid()); //Getting message from Dispatcher close(worker_pipe[1]); // Close the write end of the pipe msgbuf rec_msg; read(worker_pipe[0], rec_msg.mtext, sizeof(msgbuf)); printf("Worker message received from Dispatcher: %s\n", rec_msg.mtext); close(worker_pipe[0]); // Close the read end of the pipe sleep(1); printf("\nWorker process %d has died.\n", getpid()); exit(0); } void main(){ pthread_t dispatcher_thread; int dispatcher_pipe[n_workers][2]; pid_workers= malloc(n_workers * sizeof(pid_t)); //Creates the unnamed pipes for dispatcher for(int i=0;i<n_workers;i++){ if(pipe(dispatcher_pipe[i]) == -1){ printf("Error creating dispatcher pipe n1\n"); exit(0); } } //Creating Worker Processes for(int i=0;i<n_workers;i++){ pid_workers[i] = fork(); if (pid_workers[i] == 0) { Worker(dispatcher_pipe[i]); } else if (pid_workers[i] > 0) { } else { perror("fork"); exit(EXIT_FAILURE); } } }
问题原因与修复方案
核心问题点
父进程未清理管道多余文件描述符
fork出Worker进程后,父进程仍持有所有管道的读写端。Dispatcher线程关闭某管道读端并写入后关闭写端,但父进程中该管道的写端未关闭,导致Worker进程的read不会触发EOF,会一直阻塞等待数据。数据读写结构不匹配
失效代码中,Dispatcher写入整个msgbuf结构体,但Worker只读取mtext字段,导致数据错位,mtext会混入mtype的内容,甚至读取不完整。main函数过早退出
原失效代码中,main在等待Dispatcher线程结束后直接退出,父进程终止可能导致Worker进程还未完成输出就被强制终止。
修改后的完整代码
#include <stdio.h> #include <stdlib.h> #include <string.h> #include <unistd.h> #include <pthread.h> #include <sys/wait.h> typedef struct msgbuf { long mtype; // if 0 it did not pass in no process, 1 - worker, 2 - alert watcher char mtext[100]; }msgbuf; int n_workers = 5; pid_t *pid_workers; void* dispatcher(void *args) { int (*dispatcher_pipe)[2] = (int (*)[2])args; pthread_t id = pthread_self(); printf("Dispatcher thread started.\n"); msgbuf send_msg; send_msg.mtype = 0; strcpy(send_msg.mtext, "Message from Dispatcher"); for (int i = 0; i < n_workers; i++) { // 父进程已关闭读端,无需重复操作 write(dispatcher_pipe[i][1], &send_msg, sizeof(msgbuf)); close(dispatcher_pipe[i][1]); // 写完后关闭写端 } usleep(10); printf("Thread %ld: Dispatcher finishing\n", id); pthread_exit(NULL); } void Worker(int *worker_pipe){ printf("Worker Process (PID=%d): starting!\n", getpid()); close(worker_pipe[1]); // 关闭写端 msgbuf rec_msg; // 读取整个结构体,与写入逻辑保持一致 read(worker_pipe[0], &rec_msg, sizeof(msgbuf)); printf("Worker message received from Dispatcher: %s\n", rec_msg.mtext); close(worker_pipe[0]); sleep(1); printf("\nWorker process %d has died.\n", getpid()); exit(0); } int main(){ pthread_t dispatcher_thread; int dispatcher_pipe[n_workers][2]; pid_workers = malloc(n_workers * sizeof(pid_t)); // 创建管道 for(int i=0; i<n_workers; i++){ if(pipe(dispatcher_pipe[i]) == -1){ printf("Error creating dispatcher pipe n1\n"); exit(EXIT_FAILURE); } } // 创建Worker进程 for(int i=0; i<n_workers; i++){ pid_workers[i] = fork(); if (pid_workers[i] == 0) { Worker(dispatcher_pipe[i]); } else if (pid_workers[i] > 0) { // 父进程关闭该管道的读端(仅保留写端用于发送数据) close(dispatcher_pipe[i][0]); } else { perror("fork"); exit(EXIT_FAILURE); } } // 创建并等待Dispatcher线程完成 pthread_create(&dispatcher_thread, NULL, dispatcher, dispatcher_pipe); pthread_join(dispatcher_thread, NULL); // 等待所有Worker进程退出 for(int i=0; i<n_workers; i++){ waitpid(pid_workers[i], NULL, 0); } free(pid_workers); return 0; }
内容的提问来源于stack exchange,提问作者Duarte Neves
相关产品推荐
相关产品推荐

