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

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);
        }
    }
}
问题原因与修复方案

核心问题点

  1. 父进程未清理管道多余文件描述符
    fork出Worker进程后,父进程仍持有所有管道的读写端。Dispatcher线程关闭某管道读端并写入后关闭写端,但父进程中该管道的写端未关闭,导致Worker进程的read不会触发EOF,会一直阻塞等待数据。

  2. 数据读写结构不匹配
    失效代码中,Dispatcher写入整个msgbuf结构体,但Worker只读取mtext字段,导致数据错位,mtext会混入mtype的内容,甚至读取不完整。

  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 08:27:08