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
相关产品推荐
相关产品推荐

