使用管道传递数组:子进程向父进程传值后数组全零的问题排查
问题描述
需要实现并行程序,让父子进程分别计算大型二维数组的一半数据,子进程将计算结果通过管道传给父进程,最终合并得到完整数组。数组定义在Test结构体中:
typedef struct{ int *iterations; int height; int width; int start; int end; } Test;
因数组规模较大,采用分块读写函数chwrite和chread完成管道传输。当前代码流程为:父进程创建Test结构体并分配内存,调用processes函数;在该函数中创建管道、fork子进程,子进程计算数组前半部分并写入管道,父进程计算后半部分后读取管道数据并合并数组。但运行后发现最终的iterations数组全为零,需排查问题根源并修复。
核心代码及分块读写函数如下:
#include <stdio.h> #include <stdlib.h> #include <unistd.h> #include <sys/wait.h> #include <string.h> typedef struct{ int *iterations; int height; int width; int start; int end; } Test; // 示例计算函数 void Compute(Test *t) { for (int i = t->start; i < t->end; i++) { for (int j = 0; j < t->width; j++) { t->iterations[i * t->width + j] = i * t->width + j; } } } void chwrite(int fd, char *buf, int count, int chunksize){ int numChunks = (int)(count/chunksize); int remsize = count-numChunks*chunksize; int i, j; for(i=0,j=0; i<numChunks; j+=chunksize,i++){ write(fd, &(buf[j]), chunksize); } write(fd, &(buf[j]),remsize); } void chread(int fd, char *buf, int count, int chunksize){ int numChunks = (int)(count/chunksize); int remsize=count-numChunks*chunksize; int i,j; for(i=0,j=0; i<numChunks; j+=numChunks, i++){ read(fd, &(buf[j]),chunksize); } read(fd, &(buf[j]) ,remsize); } int main(){ Test t; t.height=1000; t.width=1000; if ((t.iterations = malloc(t.width * t.height * sizeof(int))) == NULL) { perror("Cannot allocate memory (iterations)"); exit(EXIT_FAILURE); } processes(&t); printf("First element: %d, Last element: %d\n", t.iterations[0], t.iterations[t.height*t.width-1]); free(t.iterations); return 0; } void processes(Test *Test){ int p2c[2], c2p[2]; int i,j,k; int half; int *dummy = malloc(Test->width*Test->height*sizeof(int)); if(dummy==NULL){ perror("Malloc Error"); exit(EXIT_FAILURE); } for (i = 0; i < Test->height; i++) { for (j = 0; j < Test->width; j++) { dummy[i * Test->width + j] = 0; } } half=Test->height >> 1; pipe(p2c); pipe(c2p); if(fork()==0){ Test->start=0; Test->end=half; Compute(Test); for(i=Test->start;i<Test->end;i++){ for(j=0;j<Test->width;j++){ dummy[i*Test->width+j]=Test->iterations[i*Test->width+j]; } } chwrite(c2p[1], (char*)dummy, Test->width * half * sizeof(int), 2048); printf("Sending Iterations\n"); close(p2c[0]); close(c2p[1]); exit(EXIT_SUCCESS); } else { Test->start=half; Test->end=Test->height; Compute(Test); sleep(1); chread(c2p[0], (char*)dummy, Test->width * half * sizeof(int), 2048); printf("Reading Iterations\n"); wait(NULL); for(i=0;i<Test->height;i++){ for(j=0;j<Test->width;j++){ Test->iterations[i*Test->width+j]=dummy[i*Test->width+j]; } } } }
问题根源分析
chread函数索引致命错误:循环中j += numChunks是错误逻辑,正确步长应为j += chunksize。这会导致读取的数据写入到错误内存位置,dummy数组大部分区域仍为初始0,最终覆盖iterations后全为0。- 父子进程内存空间完全独立:fork后父子进程拥有各自独立的内存副本,子进程对
Test结构体和iterations数组的修改不会同步到父进程,必须通过管道传输计算结果。 - 父进程错误覆盖自身计算结果:原代码将整个
dummy数组(后半部分为初始0)赋值给iterations,把父进程自己计算的有效后半部分数据覆盖成0。 - 管道未正确关闭所有未使用端:父子进程未关闭不需要的管道描述符,可能导致
read阻塞或数据传输异常。 sleep(1)是不可靠的同步方式:依赖固定睡眠等待子进程完成计算,实际运行中可能因计算时间过长导致父进程提前读取,或浪费不必要的时间。
修复方案
- 修正
chread的索引逻辑:将j += numChunks改为j += chunksize,确保数据写入到正确的内存位置。 - 仅合并子进程的计算结果:父进程只将
dummy中的前半部分(子进程计算结果)复制到iterations对应位置,保留自身计算的后半部分。 - 关闭所有未使用的管道端:父子进程分别关闭不需要的管道读写端,避免资源泄漏和阻塞。
- 移除
sleep(1):子进程写完数据后退出,父进程通过wait确保子进程完成,同时管道的EOF机制会让read正确返回。 - 子进程释放自身内存:子进程退出前释放
dummy数组,避免内存泄漏。
修复后的关键代码片段
修正后的chread函数
void chread(int fd, char *buf, int count, int chunksize){ int numChunks = (int)(count/chunksize); int remsize=count-numChunks*chunksize; int i,j; for(i=0,j=0; i<numChunks; j+=chunksize, i++){ // 修正j的步长 read(fd, &(buf[j]), chunksize); } read(fd, &(buf[j]) , remsize); }
父进程合并数据的逻辑
// 仅合并子进程传来的前半部分,保留父进程计算的后半部分 for(i=0;i<half;i++){ for(j=0;j<Test->width;j++){ Test->iterations[i*Test->width+j]=dummy[i*Test->width+j]; } }
完善后的processes函数管道处理
void processes(Test *Test){ int p2c[2], c2p[2]; int i,j,k; int half; int *dummy = malloc(Test->width*Test->height*sizeof(int)); if(dummy==NULL){ perror("Malloc Error"); exit(EXIT_FAILURE); } for (i = 0; i < Test->height; i++) { for (j = 0; j < Test->width; j++) { dummy[i * Test->width + j] = 0; } } half=Test->height >> 1; pipe(p2c); pipe(c2p); if(fork()==0){ // 关闭不需要的管道端 close(p2c[1]); close(c2p[0]); Test->start=0; Test->end=half; Compute(Test); for(i=Test->start;i<Test->end;i++){ for(j=0;j<Test->width;j++){ dummy[i*Test->width+j]=Test->iterations[i*Test->width+j]; } } chwrite(c2p[1], (char*)dummy, Test->width * half * sizeof(int), 2048); printf("Sending Iterations\n"); close(c2p[1]); close(p2c[0]); free(dummy); exit(EXIT_SUCCESS); } else { // 关闭不需要的管道端 close(p2c[0]); close(c2p[1]); Test->start=half; Test->end=Test->height; Compute(Test); chread(c2p[0], (char*)dummy, Test->width * half * sizeof(int), 2048); printf("Reading Iterations\n"); wait(NULL); // 仅合并子进程传来的前半部分 for(i=0;i<half;i++){ for(j=0;j<Test->width;j++){ Test->iterations[i*Test->width+j]=dummy[i*Test->width+j]; } } close(c2p[0]); close(p2c[1]); free(dummy); } }
内容的提问来源于stack exchange,提问作者Lach
相关产品推荐
相关产品推荐

