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

信息检索系统:大文档集下部分索引合并算法提速咨询

提升信息检索系统部分索引合并速度的优化方案

问题背景

为课程项目开发信息检索系统,需对107000个.nxml格式文档构建索引。采用partial indexing技术避免内存占用过高,已生成多个partial posting和vocabulary文件,可合并为单个索引文件。小集合(54个文档)测试流程正常,但大集合测试时合并耗时远超预期,希望提升合并速度,考虑过线程并发但不熟悉其用法。

文件结构说明

部分文件按以下方式存储:

  • posting-0.txt, vocabulary-0.txt
  • posting-1.txt, vocabulary-1.txt
  • ...
  • posting-(n-1).txt, vocabulary-(n-1).txt
    通过文件名中的数字关联对应的partial posting和vocabulary文件,并存入内存队列等待合并。

各文件结构及排序规则:

  • 词汇表文件:term, doc_freq (df), pointer_to_posting_file,按term升序排列
  • 倒排列表文件:doc_id, term_freq (tf), position(s)_of_term_in_doc, pointer_to_doc_file,先按term、再按doc_id升序排列
  • 文档文件:doc_id, path, vector_length,按doc_id升序排列
    最终合并后的文件需保持上述排序规则。

当前合并代码

public static void copyPostingBlock(RandomAccessFile inputRaf, long pointer, RandomAccessFile outputRaf, int df) throws IOException {

    inputRaf.seek(pointer);
    for (int i = 0; i < df; i++){

        String postingLine = inputRaf.readLine();
        outputRaf.writeBytes(postingLine + "\n");

    }
}


public static void mergePartialFiles() throws IOException {

    int length = (int) (partialVocabularyFiles.size() / 2);
    System.out.println("length: " + length);

    File collectionIndex = new File(System.getProperty("user.dir") + "\\src\\resources\\CollectionIndex");

    for(int i = 0; i < length; i++){  // for each pair of partial files

        File vocab1 = partialVocabularyFiles.remove();
        File vocab2 = partialVocabularyFiles.remove();
        File posting1 = partialPostingFiles.remove();
        File posting2 = partialPostingFiles.remove();

        BufferedReader reader1 = new BufferedReader(new FileReader(vocab1));
        BufferedReader reader2 = new BufferedReader(new FileReader(vocab2));
        RandomAccessFile raf1 = new RandomAccessFile(posting1, "r");
        RandomAccessFile raf2 = new RandomAccessFile(posting2, "r");

        File mergedVocab = new File(collectionIndex, "merged-vocabulary-" + mergesNum + ".txt");
        File mergedPosting = new File(collectionIndex, "merged-posting-" + mergesNum++ + ".txt");
        RandomAccessFile rafMergedVocab = new RandomAccessFile(mergedVocab, "rw");
        RandomAccessFile rafMergedPosting = new RandomAccessFile(mergedPosting, "rw");

        String currentLine1 = reader1.readLine();
        String currentLine2 = reader2.readLine();

        while(currentLine1 != null && currentLine2 != null){   // O(N)

            String[] tokens1 = currentLine1.split(", ");
            String[] tokens2 = currentLine2.split(", ");

            String word1 = tokens1[0];
            String word2 = tokens2[0];
            int df1 = Integer.parseInt(tokens1[1]);
            int df2 = Integer.parseInt(tokens2[1]);
            long pointer1 = Long.parseLong(tokens1[2]);
            long pointer2 = Long.parseLong(tokens2[2]);

            long mergedPointer = rafMergedPosting.getFilePointer();

            if(word1.compareTo(word2) < 0){     // if w_i < w_j

                copyPostingBlock(raf1, pointer1, rafMergedPosting, df1);
                rafMergedVocab.writeBytes(word1 + ", " + df1 + ", " + mergedPointer + "\n");
                currentLine1 = reader1.readLine();

            }else if(word1.compareTo(word2) > 0) {    // if w_j < w_i

                copyPostingBlock(raf2, pointer2, rafMergedPosting, df2);
                rafMergedVocab.writeBytes(word2 + ", " + df2 + ", " + mergedPointer + "\n");
                currentLine2 = reader2.readLine();

            }else{      // if w_i == w_j

                // TODO: keep the word in memory until another word is discovered
                copyPostingBlock(raf1, pointer1, rafMergedPosting, df1);
                copyPostingBlock(raf2, pointer2, rafMergedPosting, df2);
                int newDf = df1 + df2;
                rafMergedVocab.writeBytes(word1 + ", " + newDf + ", " + mergedPointer + "\n");
                currentLine1 = reader1.readLine();
                currentLine2 = reader2.readLine();

            }

        }

        // handle leftover lines
        while(currentLine1 != null){

            String[] tokens1 = currentLine1.split(", ");
            String word1 = tokens1[0];
            int df1 = Integer.parseInt(tokens1[1]);
            long pointer1 = rafMergedPosting.getFilePointer();

            long mergedPointer = rafMergedPosting.getFilePointer();
            copyPostingBlock(raf1, pointer1, rafMergedPosting, df1);
            rafMergedVocab.writeBytes(word1 + ", " + df1 + ", " + mergedPointer + "\n");
            currentLine1 = reader1.readLine();

        }

        while(currentLine2 != null){

            String[] tokens2 = currentLine2.split(", ");
            String word2 = tokens2[0];
            int df2 = Integer.parseInt(tokens2[1]);
            long pointer2 = rafMergedPosting.getFilePointer();

            long pointer1 = rafMergedPosting.getFilePointer();
            copyPostingBlock(raf2, pointer2, rafMergedPosting, df2);
            rafMergedVocab.writeBytes(word2 + ", " + df2 + ", " + pointer2 + "\n");
            currentLine2 = reader2.readLine();

        }

        // after writing the new merged files, close and delete the partial files used for merging
        // add the new merged file to the queue

        reader1.close();
        reader2.close();
        raf1.close();
        raf2.close();
        rafMergedVocab.close();
        rafMergedPosting.close();

        vocab1.delete();
        vocab2.delete();
        posting1.delete();
        posting2.delete();

        partialVocabularyFiles.add(mergedVocab);
        partialPostingFiles.add(mergedPosting);

    }

    if(partialPostingFiles.size() != 1)
        mergePartialFiles();

}

优化方案

1. 批量IO替代逐行读写(最直接的性能提升)

当前copyPostingBlock逐行读取写入,IO调用频繁开销极大。改为批量读取字节块,减少IO交互次数:

public static void copyPostingBlock(RandomAccessFile inputRaf, long pointer, RandomAccessFile outputRaf, int df) throws IOException {
    inputRaf.seek(pointer);
    byte[] buffer = new byte[8192]; // 8KB缓冲区,可根据磁盘性能调整为16KB/32KB
    int bytesRead;
    int linesCopied = 0;
    
    while (linesCopied < df && (bytesRead = inputRaf.read(buffer)) != -1) {
        int start = 0;
        for (int i = 0; i < bytesRead && linesCopied < df; i++) {
            if (buffer[i] == '\n') {
                outputRaf.write(buffer, start, i - start + 1);
                start = i + 1;
                linesCopied++;
            }
        }
        // 处理未读完的行
        if (start < bytesRead && linesCopied < df) {
            outputRaf.write(buffer, start, bytesRead - start);
        }
    }
}

进阶优化:生成partial索引时,额外记录每个posting块的字节长度,合并时直接按字节数批量复制,彻底避免逐行计数的开销:

  • 词汇表新增posting_block_byte_length字段
  • 合并时直接调用inputRaf.read(buffer, 0, byteLength)复制整个块

2. 多线程并发合并配对文件

当前串行处理每一对文件,可利用线程池并行处理不同文件对,充分利用CPU和磁盘IO带宽。

修改后的核心代码:

// 线程池大小设为CPU核心数,避免过多线程导致IO竞争
private static final ExecutorService executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());
// 用原子类保证线程安全的合并序号
private static final AtomicInteger mergesNum = new AtomicInteger(0);
// 使用线程安全队列存储文件
private static final ConcurrentLinkedQueue<File> partialVocabularyFiles = new ConcurrentLinkedQueue<>();
private static final ConcurrentLinkedQueue<File> partialPostingFiles = new ConcurrentLinkedQueue<>();

public static void mergePartialFiles() throws InterruptedException, IOException {
    int length = (int) (partialVocabularyFiles.size() / 2);
    System.out.println("length: " + length);

    File collectionIndex = new File(System.getProperty("user.dir") + "\\src\\resources\\CollectionIndex");
    List<Future<?>> futures = new ArrayList<>();

    for(int i = 0; i < length; i++){
        File vocab1 = partialVocabularyFiles.poll();
        File vocab2 = partialVocabularyFiles.poll();
        File posting1 = partialPostingFiles.poll();
        File posting2 = partialPostingFiles.poll();
        
        // 每对文件交给一个线程处理
        futures.add(executor.submit(() -> {
            try {
                mergePair(vocab1, vocab2, posting1, posting2, collectionIndex);
            } catch (IOException e) {
                e.printStackTrace();
            }
        }));
    }

    // 等待本轮所有合并任务完成
    for (Future<?> future : futures) {
        future.get();
    }

    if(partialPostingFiles.size() != 1)
        mergePartialFiles();
    else {
        // 合并完成后关闭线程池
        executor.shutdown();
    }
}

// 抽取出合并单对文件的独立方法
private static void mergePair(File vocab1, File vocab2, File posting1, File posting2, File collectionIndex) throws IOException {
    BufferedReader reader1 = new BufferedReader(new FileReader(vocab1));
    BufferedReader reader2 = new BufferedReader(new FileReader(vocab2));
    RandomAccessFile raf1 = new RandomAccessFile(posting1, "r");
    RandomAccessFile raf2 = new RandomAccessFile(posting2, "r");

    int currentMergeNum = mergesNum.getAndIncrement();
    File mergedVocab = new File(collectionIndex, "merged-vocabulary-" + currentMergeNum + ".txt");
    File mergedPosting = new File(collectionIndex, "merged-posting-" + currentMergeNum + ".txt");
    RandomAccessFile rafMergedVocab = new RandomAccessFile(mergedVocab, "rw");
    RandomAccessFile rafMergedPosting = new RandomAccessFile(mergedPosting, "rw");

    String currentLine1 = reader1.readLine();
    String currentLine2 = reader2.readLine();

    while(currentLine1 != null && currentLine2 != null){
        String[] tokens1 = currentLine1.split(", ");
        String[] tokens2 = currentLine2.split(", ");

        String word1 = tokens1[0];
        String word2 = tokens2[0];
        int df1 = Integer.parseInt(tokens1[1]);
        int df2 = Integer.parseInt(tokens2[1]);
        long pointer1 = Long.parseLong(tokens1[2]);
        long pointer2 = Long.parseLong(tokens2[2]);

        long mergedPointer = rafMergedPosting.getFilePointer();

        if(word1.compareTo(word2) < 0){
            copyPostingBlock(raf1, pointer1, rafMergedPosting, df1);
            rafMergedVocab.writeBytes(word1 + ", " + df1 + ", " + mergedPointer + "\n");
            currentLine1 = reader1.readLine();
        }else if(word1.compareTo(word2) > 0) {
            copyPostingBlock(raf2, pointer2, rafMergedPosting, df2);
            rafMergedVocab.writeBytes(word2 + ", " + df2 + ", " + mergedPointer + "\n");
            currentLine2 = reader2.readLine();
        }else{
            copyPostingBlock(raf1, pointer1, rafMergedPosting, df1);
            copyPostingBlock(raf2, pointer2, rafMergedPosting, df2);
            int newDf = df1 + df2;
            rafMergedVocab.writeBytes(word1 + ", " + newDf + ", " + mergedPointer + "\n");
            currentLine1 = reader1.readLine();
            currentLine2 = reader2.readLine();
        }
    }

    // 处理剩余行
    while(currentLine1 != null){
        String[] tokens1 = currentLine1.split(", ");
        String word1 = tokens1[0];
        int df1 = Integer.parseInt(tokens1[1]);
        long mergedPointer = rafMergedPosting.getFilePointer();

        copyPostingBlock(raf1, pointer1, rafMergedPosting, df1);
        rafMergedVocab.writeBytes(word1 + ", " + df1 + ", " + mergedPointer + "\n");
        currentLine1 = reader1.readLine();
    }

    while(currentLine2 != null){
        String[] tokens2 = currentLine2.split(", ");
        String word2 = tokens2[0];
        int df2 = Integer.parseInt(tokens2[1]);
        long mergedPointer = rafMergedPosting.getFilePointer();

        copyPostingBlock(raf2, pointer2, rafMergedPosting, df2);
        rafMergedVocab.writeBytes(word2 + ", " + df2 + ", " + mergedPointer + "\n");
        currentLine2 = reader2.readLine();
    }

    // 关闭资源
    reader1.close();
    reader2.close();
    raf1.close();
    raf2.close();
    rafMergedVocab.close();
    rafMergedPosting.close();

    // 删除原文件
    vocab1.delete();
    vocab2.delete();
    posting1.delete();
    posting2.delete();

    // 添加合并后的文件到队列
    partialVocabularyFiles.add(mergedVocab);
    partialPostingFiles.add(mergedPosting);
}

3. 优化文件存储格式

当前文本格式解析和写入开销大,改为二进制格式可大幅提升效率:

  • 数值类型(df、pointer、doc_id等)用固定长度字节存储(如int占4字节,long占8字节)
  • 字符串用UTF-8编码加长度前缀存储
    合并时无需频繁调用split()和parseInt(),直接按字节读取解析。

4. 磁盘IO硬件/配置优化

  • 将partial文件和合并输出文件存储在不同物理磁盘,避免读写冲突
  • 使用SSD替代HDD,提升随机读写性能
  • 调整JVM IO缓冲区大小:启动参数添加-Dsun.nio.bufsize=1048576(设置为1MB)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 06:32:31