信息检索系统:大文档集下部分索引合并算法提速咨询
提升信息检索系统部分索引合并速度的优化方案
问题背景
为课程项目开发信息检索系统,需对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
相关产品推荐
相关产品推荐

