多线程场景下高效拆分与访问mmap映射数据的最优方案
优化多线程mmap文件处理的CPU耗时方案
问题背景
用POSIX编写了多线程C程序处理两个大型TSV文件,生成对应输出时存在性能瓶颈。三个核心文件的结构如下:
- File 1:存储国家-键值对(例如:
African 0、Amerindian 1)。 - File 2:大型n×m矩阵,每个数值对应File1的键值。示例:
0\t3\t4\t1\t2... 0\t3\t4\t1\t2... - File 3:大型n×(m+1)矩阵,列值为0、1、9,首列是长度可变的整数索引。示例:
111111\t0\t1\t0\t0\t1... 22222233\t0\t1\t1\t1\t1...
处理目标
为每个国家生成与File3结构一致的输出文件(如Amerindian.txt),处理逻辑:
- 复制File3的首列索引;
- 每个位置对比File2的键值与当前国家键:不匹配输出
9,匹配则复制File3对应值。
示例输出(键为1的Amerindian.txt):
111111\t9\t9\t9\t0\t9... 22222233\t9\t9\t9\t1\t9...
当前实现
已采用mmap方案重写main函数与核心处理函数,核心代码如下:
main函数片段
// 库加载、函数/变量/输入参数声明... // 处理origins文件 FILE *originsFile = fopen(argv[optind], "r"); originCount = readOrigins(originsFile, &origins); printf("Number of origins parsed: %d\n", originCount); // 创建mmap映射 int tsv_fd = open(argv[optind + 1], O_RDONLY); int txt_fd = open(argv[optind + 2], O_RDONLY); struct stat tsv_stat, txt_stat; fstat(tsv_fd, &tsv_stat); fstat(txt_fd, &txt_stat); char *tsv_mapped = mmap(NULL, tsv_stat.st_size, PROT_READ, MAP_SHARED, tsv_fd, 0); char *txt_mapped = mmap(NULL, txt_stat.st_size, PROT_READ, MAP_SHARED, txt_fd, 0); close(txt_fd); close(tsv_fd); // 检查mmap错误 if (tsv_mapped == MAP_FAILED || txt_mapped == MAP_FAILED) { perror("Error mapping files"); // 错误处理:关闭文件描述符,可能退出 } // 统计文件行数 txtLineCount = countLinesInFile(txt_mapped, txt_stat.st_size); // 计算线程工作负载 int linesPerThread = txtLineCount / numThreads; int remainingLines = txtLineCount % numThreads; // 合理限制线程数(不超过总行数) if (txtLineCount < numThreads) { numThreads = txtLineCount; linesPerThread = txtLineCount / numThreads; remainingLines = txtLineCount % numThreads; } pthread_t threads[numThreads]; ThreadArg *threadArgs[numThreads]; for (int originIndex = 0; originIndex < originCount; ++originIndex) { printf("Processing origin: %s\n", origins[originIndex].name); int currentLine = 0; for (int i = 0; i < numThreads; ++i) { printf("Debug: Allocating ThreadArg for thread %d\n", i); threadArgs[i] = (ThreadArg *)malloc(sizeof(ThreadArg)); threadArgs[i]->startLine = currentLine; threadArgs[i]->endLine = currentLine + linesPerThread; if (i < remainingLines) threadArgs[i]->endLine++; currentLine = threadArgs[i]->endLine; threadArgs[i]->tsvMapped = tsv_mapped; threadArgs[i]->txtMapped = txt_mapped; threadArgs[i]->specificOriginIndex = origins[originIndex].code; threadArgs[i]->txt_stat = txt_stat.st_size; threadArgs[i]->tsv_stat = tsv_stat.st_size; threadArgs[i]->threadnum = i; printf("Debug: Thread %d - startLine: %d, endLine: %d, specificOriginIndex: %d\n", i, threadArgs[i]->startLine, threadArgs[i]->endLine, threadArgs[i]->specificOriginIndex); if (pthread_create(&threads[i], NULL, processFile, threadArgs[i])) { fprintf(stderr, "Failed to allocate memory for thread arguments\n"); // 清理资源并返回 } } }
processFile与getLinesFromMappedFile函数
void *processFile(void *arg) { ThreadArg *threadArg = (ThreadArg *)arg; char* getLinesFromMappedFile(const char* mappedFile, int startLine, int endLine, size_t fileSize, size_t* byteRange); char tempFileName[1024]; sprintf(tempFileName, "temp_output_%d_%d.tmp", threadArg->specificOriginIndex, threadArg->threadnum); FILE *outputFile = fopen(tempFileName, "a"); // 追加模式 if (!outputFile) { fprintf(stderr, "Failed to open output file.\n"); return NULL; } size_t byteRangeTsv; const char* segmentTsv = getLinesFromMappedFile(threadArg->tsvMapped, threadArg->startLine, threadArg->endLine, threadArg->tsv_stat, &byteRangeTsv); size_t byteRangeTxt; const char* segmentTxt = getLinesFromMappedFile(threadArg->txtMapped, threadArg->startLine, threadArg->endLine, threadArg->txt_stat, &byteRangeTxt); if (!segmentTsv || !segmentTxt) { fprintf(stderr, "Lines not found in one of the files.\n"); if (outputFile) fclose(outputFile); return NULL; } printf("%ld,%ld", byteRangeTsv, byteRangeTxt); // 按逻辑处理文件片段 for (size_t tsvIndex = 0, txtIndex = 0; tsvIndex < byteRangeTsv && txtIndex < byteRangeTxt; ) { if(txtIndex == 0){ // 复制File3的首列索引到输出 while (segmentTxt[txtIndex] != '\t' && segmentTxt[txtIndex] != '\0') { fprintf(outputFile, "%c", segmentTxt[txtIndex]); txtIndex++; tsvIndex++; } // 跳过File3的制表符 if (segmentTxt[txtIndex] == '\t') { fprintf(outputFile, "\t"); txtIndex++; } } // 处理行内数据 while (segmentTsv[tsvIndex] != '\n' && segmentTxt[txtIndex] != '\n') { if (segmentTsv[tsvIndex] == '\t') { fprintf(outputFile, "\t"); txtIndex++; tsvIndex++; } if (segmentTsv[tsvIndex] == threadArg->specificOriginIndex) { fprintf(outputFile, "%c\t", segmentTxt[txtIndex]); } else { fprintf(outputFile, "9\t"); } txtIndex++; tsvIndex++; } fprintf(outputFile, "\n"); // 跳转至下一行开头 while (segmentTsv[tsvIndex] != '\n' && tsvIndex < byteRangeTsv) tsvIndex++; while (segmentTxt[txtIndex] != '\n' && txtIndex < byteRangeTxt) txtIndex++; if (segmentTsv[tsvIndex] == '\n') tsvIndex++; if (segmentTxt[txtIndex] == '\n') txtIndex++; } fclose(outputFile); return NULL; } char* getLinesFromMappedFile(const char* mappedFile, int startLine, int endLine, size_t fileSize, size_t* byteRange) { if (mappedFile == NULL || startLine < 0 || endLine < startLine) { return NULL; // 参数无效 } const char* startPtr = NULL; const char* endPtr = NULL; int currentLine = 0; for (size_t i = 0; i < fileSize; ++i) { if (currentLine == startLine && startPtr == NULL) { startPtr = mappedFile + i; // 标记起始行位置 } if (mappedFile[i] == '\n') { if (currentLine == endLine) { endPtr = mappedFile + i; // 标记结束行位置 break; } currentLine++; } } if (startPtr == NULL) { return NULL; // 未找到起始行 } // 若未找到结束行,以文件末尾为结束位置 if (endPtr == NULL) { endPtr = mappedFile + fileSize; } // 计算字节范围 *byteRange = endPtr - startPtr; return (char*)startPtr; // 转换为非const类型匹配函数签名 }
核心疑问
以最低CPU耗时为目标,拆分并让每个线程访问mmap映射数据的最优方案是什么?
内容的提问来源于stack exchange,提问作者mugdi
相关产品推荐
相关产品推荐

