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

C语言多线程管道实现故障:重复执行同一函数问题排查

多线程管道实现问题:线程重复执行同一函数的排查与修复

问题描述

我正在实现一个多线程管道,将不同函数分配到多个线程,通过管道实现线程间通信。但运行程序时,虽然创建了多个线程,却出现线程重复执行同一函数的问题,怀疑是管道实现有误,但找不到具体错误点。

实现代码

#include <stdio.h>
#include <stdlib.h>
#include <unistd.h>
#include <pthread.h>
#include <stdbool.h>
int num_ints;
typedef void (*Function)(void* input, void* output);

typedef struct Pipeline {
    Function* functions;
    int numStages;
} Pipeline;
 
Pipeline* new_Pipeline() {
    Pipeline *this = malloc(sizeof(Pipeline));
    this->functions = NULL;
    this->numStages = 0;
    printf("created the pipeline\n");
    return this;   
}

bool Pipeline_add(Pipeline* this, Function f) {
    //reallocating memory to add a stage to the functions array
    this->functions = realloc(this->functions, (this->numStages +1) * sizeof(Function));
    if (this->functions == NULL) {
        return false;
    } 
    else {
        this->functions[this->numStages] = f;
        this->numStages++;
        printf("added a stage\n");
        return true;
    }
    
}

typedef struct {
    Function function;
    int inputPipe;
    int outputPipe;
} thread_args;

void* thread_func(void* arg) {
    thread_args *data = (thread_args*) arg;

    //get the input and output pipes from the args parameter
    int inPipe = data->inputPipe;
    int outPipe = data->outputPipe;
    data->function((void*)&inPipe,(void*)&outPipe);
    return NULL;
}

void Pipeline_execute(Pipeline* this) {
    //create threads 
    pthread_t threads[this->numStages];
    thread_args args[this->numStages];
    for (int i = 0; i < this->numStages; i++) {
         //creating the pipes
        int fd[2]; 
        pipe(fd);
        //creating input and output pipes for each stage
        args[i].function = this->functions[i];
        args[i].inputPipe = fd[0];
        args[i].outputPipe = fd[1];
        if (pthread_create(&threads[i], NULL, thread_func, &args[i]) != 0) {
            printf("created a thread\n");
            perror("pthread_create\n");
        }
        if (i == this->numStages -1) {
            close(fd[0]);
        }
    }
    //waiting for threads to finish
    for (int i = 0; i < this->numStages; i++) {
        if (pthread_join(threads[i], NULL) != 0) {
            perror("pthread_join\n");
            exit(1);
        }
    }
    //closing pipes
    for (int i = 0; i < this->numStages -1; i++) {
        
    }
    
}


void Pipeline_free(Pipeline* this) {
    free(this->functions);
    free(this);
}


bool Pipeline_send(void* channel, void* buffer, size_t size) {
    if ((write(*(int*)channel, buffer, size)) != -1) {
        return true;
    } else {
        return false;
    }
    
}


bool Pipeline_receive(void* channel, void* buffer, size_t size) {
    if ((read(*(int*)channel, buffer, size)) != -1) {
        return true; 
    } else {
        return false;
    }
}
//an application created to help test the implementation of pipes.

static void generateInts(void* input, void* output) {
    printf("generateInts: thread %p\n", (void*) pthread_self());
    for (int i = 1; i <= num_ints; i++) {
        if (!Pipeline_send(output, (void*) &i, sizeof(int))) exit(EXIT_FAILURE);
    }
}


static void squareInts(void* input, void* output) {
    printf("squareInts: thread %p\n", (void*) pthread_self());
    for (int i = 1; i <= num_ints; i++) {
        int number;
        if (!Pipeline_receive(input, (void*) &number, sizeof(int))) exit(EXIT_FAILURE);
        int result = number * number;
        if (!Pipeline_send(output, (void*) &result, sizeof(int))) exit(EXIT_FAILURE);
    }
}


static void sumIntsAndPrint(void* input, void* output) {
    printf("sumIntsAndPrint: thread %p\n", (void*) pthread_self());
    int number = 0;
    int result = 0;
    for (int i = 1; i <= num_ints; i++) {
        if (!Pipeline_receive(input, (void*) &number, sizeof(int))) exit (EXIT_FAILURE);
        result += number;
    }
    printf("sumIntsAndPrint: result = %i\n", result);
}

static void cleanupExit(Pipeline *p) {
    if (p) {
        Pipeline_free(p);
    }
    exit(EXIT_FAILURE);
}


int main() {
    scanf("%d", &num_ints);
    printf("Setting up pipeline to calculate the sum of squares of integers 1 to %i.\n", num_ints);

    Pipeline *p = new_Pipeline();

    if (p == NULL) cleanupExit(p);
    if (!Pipeline_add(p, generateInts)) cleanupExit(p);
    if (!Pipeline_add(p, squareInts)) cleanupExit(p);
    if (!Pipeline_add(p, sumIntsAndPrint)) cleanupExit(p);
    Pipeline_execute(p);

    Pipeline_free(p);
    return 0;
}

运行输出

Setting up pipeline to calculate the sum of squares of integers 1 to 10.
created the pipeline
added a stage
added a stage
added a stage
generateInts: thread 0x5677640
squareInts: thread 0x5e78640
sumIntsAndPrint: thread 0x6679640
==162384== 
==162384== HEAP SUMMARY:
==162384==     in use at exit: 584 bytes in 4 blocks
==162384==   total heap usage: 9 allocs, 5 frees, 10,096 bytes allocated
==162384==
==162384== LEAK SUMMARY:
==162384==    definitely lost: 0 bytes in 0 blocks
==162384==    indirectly lost: 0 bytes in 0 blocks
==162384==      possibly lost: 544 bytes in 2 blocks
==162384==    still reachable: 40 bytes in 2 blocks

问题分析

从运行输出看,三个线程实际上已经执行了各自对应的函数(generateInts、squareInts、sumIntsAndPrint均有打印),但程序逻辑存在核心问题导致管道通信失效,可能让你误以为线程执行了同一函数:

  • 管道连接错误:当前为每个线程单独创建管道,但未将前一个线程的输出管道与后一个线程的输入管道连接,导致每个线程的管道独立,数据无法在流水线中传递。
  • 管道文件描述符管理混乱:未正确关闭不需要的管道端,导致read操作阻塞,线程无法正常结束;最后一个线程的输入管道被提前关闭,无法读取数据。
  • 固定循环次数隐患:所有阶段都使用固定的num_ints循环次数,而非根据管道实际读取的数据量判断,若数据传递异常会导致线程阻塞或错误。

修复方案

1. 修正Pipeline_execute函数,正确连接流水线管道

void Pipeline_execute(Pipeline* this) {
    pthread_t threads[this->numStages];
    thread_args args[this->numStages];
    int prev_out_fd = -1; // 保存前一个阶段的输出管道写端

    for (int i = 0; i < this->numStages; i++) {
        int fd[2]; 
        if (pipe(fd) == -1) {
            perror("pipe");
            exit(EXIT_FAILURE);
        }

        args[i].function = this->functions[i];
        // 第一个阶段无输入管道,设为-1
        args[i].inputPipe = (i == 0) ? -1 : prev_out_fd;
        args[i].outputPipe = fd[1];

        if (pthread_create(&threads[i], NULL, thread_func, &args[i]) != 0) {
            perror("pthread_create");
            exit(EXIT_FAILURE);
        }

        // 关闭当前管道的读端(由下一个阶段使用)
        close(fd[0]);
        // 非第一个阶段,关闭前一个阶段的输出写端(当前线程已保存)
        if (i != 0) {
            close(prev_out_fd);
        }
        // 更新前一个输出写端为当前管道的写端
        prev_out_fd = fd[1];
    }

    // 关闭最后一个阶段的输出写端(最后一个阶段无需输出)
    close(prev_out_fd);

    // 等待线程完成
    for (int i = 0; i < this->numStages; i++) {
        if (pthread_join(threads[i], NULL) != 0) {
            perror("pthread_join");
            exit(EXIT_FAILURE);
        }
    }
}

2. 修正阶段函数,处理管道关闭与动态读取逻辑

static void generateInts(void* input, void* output) {
    printf("generateInts: thread %p\n", (void*) pthread_self());
    for (int i = 1; i <= num_ints; i++) {
        if (!Pipeline_send(output, (void*) &i, sizeof(int))) exit(EXIT_FAILURE);
    }
    // 关闭输出管道,通知后续阶段数据发送完毕
    close(*(int*)output);
}

static void squareInts(void* input, void* output) {
    printf("squareInts: thread %p\n", (void*) pthread_self());
    int number;
    ssize_t read_bytes;
    // 循环读取直到管道关闭(read返回0)
    while ((read_bytes = read(*(int*)input, (void*) &number, sizeof(int))) > 0) {
        int result = number * number;
        if (!Pipeline_send(output, (void*) &result, sizeof(int))) exit(EXIT_FAILURE);
    }
    // 关闭输入、输出管道
    close(*(int*)input);
    close(*(int*)output);
}

static void sumIntsAndPrint(void* input, void* output) {
    printf("sumIntsAndPrint: thread %p\n", (void*) pthread_self());
    int number = 0;
    int result = 0;
    ssize_t read_bytes;
    while ((read_bytes = read(*(int*)input, (void*) &number, sizeof(int))) > 0) {
        result += number;
    }
    printf("sumIntsAndPrint: result = %i\n", result);
    // 关闭输入管道
    close(*(int*)input);
}

3. 修正Pipeline_send和Pipeline_receive函数,确保完整读写

bool Pipeline_send(void* channel, void* buffer, size_t size) {
    int fd = *(int*)channel;
    ssize_t written = write(fd, buffer, size);
    return written == size; // 确保数据完整写入
}

bool Pipeline_receive(void* channel, void* buffer, size_t size) {
    int fd = *(int*)channel;
    ssize_t read_bytes = read(fd, buffer, size);
    // 读取到完整数据或管道关闭(返回0)时返回true,读取失败返回false
    return (read_bytes == size) || (read_bytes == 0);
}

关键修复点说明

  • 流水线管道链式连接:每个阶段的输入管道承接前一个阶段的输出管道,形成数据传递的链路。
  • 及时关闭管道描述符:每个线程只保留需要的管道端,避免资源泄漏和read操作无限阻塞。
  • 动态读取逻辑:通过read的返回值判断数据是否读取完毕,适配流水线的动态数据流。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 03:37:02