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

多线程场景下高效拆分与访问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),处理逻辑:

  1. 复制File3的首列索引;
  2. 每个位置对比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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 00:48:14