如何用NodeJS基于归并排序排序10GB+大词文件(1GB内存限制)
10GB+大文件在1GB内存限制下的NodeJS原生归并排序实现
针对大文件内存不足的排序场景,外部归并排序是标准解决方案,分为分片局部排序和多路归并两个核心阶段,以下是基于NodeJS原生API的具体实现方案:
一、核心流程概述
- 分片排序:将10GB的大文件拆分为多个可完全载入内存的小文件(如每个700MB,预留内存余量),对每个小文件的词进行排序后写入磁盘临时文件。
- 多路归并:同时读取所有排序后的临时文件,通过最小堆(优先级队列)每次取出当前最小的词,写入最终排序文件,直到所有临时文件处理完毕。
二、分片排序实现
拆分时需确保每个词的完整性(拆分点必须在逗号后),避免截断词导致排序错误。使用NodeJS流API逐块读取文件,累计到指定内存阈值后排序写入临时文件。
const fs = require('fs'); const path = require('path'); // 配置项:每个临时文件的目标内存占用(700MB,预留30%内存给系统) const MAX_CHUNK_MEMORY = 700 * 1024 * 1024; const INPUT_FILE = './large_input.txt'; const TEMP_DIR = './temp_chunks'; // 创建临时目录 if (!fs.existsSync(TEMP_DIR)) { fs.mkdirSync(TEMP_DIR, { recursive: true }); } async function splitAndSortChunks() { const stream = fs.createReadStream(INPUT_FILE, { highWaterMark: 64 * 1024 }); // 64KB缓冲区 let buffer = ''; let chunkIndex = 0; let currentMemory = 0; for await (const chunk of stream) { buffer += chunk.toString(); const lastCommaIndex = buffer.lastIndexOf(','); if (lastCommaIndex === -1) continue; // 提取可处理的完整词,剩余内容留到下一轮 const processablePart = buffer.slice(0, lastCommaIndex); buffer = buffer.slice(lastCommaIndex + 1); const words = processablePart.split(','); currentMemory += Buffer.byteLength(processablePart, 'utf8'); // 达到内存阈值则排序写入临时文件 if (currentMemory >= MAX_CHUNK_MEMORY) { words.sort(); const tempFile = path.join(TEMP_DIR, `chunk_${chunkIndex}.txt`); fs.writeFileSync(tempFile, words.join(',')); chunkIndex++; currentMemory = 0; } } // 处理最后剩余的词 if (buffer.trim()) { const words = buffer.split(',').filter(w => w.trim()); words.sort(); const tempFile = path.join(TEMP_DIR, `chunk_${chunkIndex}.txt`); fs.writeFileSync(tempFile, words.join(',')); } return chunkIndex + 1; // 返回临时文件总数 }
三、多路归并实现
使用自定义最小堆管理多个临时文件的当前待选词,每次取出最小词写入输出文件,再从对应文件读取下一个词补充到堆中,直到所有文件处理完成。
// 最小堆实现:存储{ word: 待选词, stream: 文件流, index: 临时文件索引 } class MinHeap { constructor() { this.heap = []; } push(element) { this.heap.push(element); this.bubbleUp(this.heap.length - 1); } bubbleUp(index) { while (index > 0) { const parentIndex = Math.floor((index - 1) / 2); if (this.heap[parentIndex].word <= this.heap[index].word) break; [this.heap[parentIndex], this.heap[index]] = [this.heap[index], this.heap[parentIndex]]; index = parentIndex; } } pop() { const min = this.heap[0]; const end = this.heap.pop(); if (this.heap.length > 0) { this.heap[0] = end; this.sinkDown(0); } return min; } sinkDown(index) { const left = 2 * index + 1; const right = 2 * index + 2; let smallest = index; const length = this.heap.length; if (left < length && this.heap[left].word < this.heap[smallest].word) { smallest = left; } if (right < length && this.heap[right].word < this.heap[smallest].word) { smallest = right; } if (smallest !== index) { [this.heap[index], this.heap[smallest]] = [this.heap[smallest], this.heap[index]]; this.sinkDown(smallest); } } isEmpty() { return this.heap.length === 0; } } async function mergeChunks(chunkCount, outputFile) { const heap = new MinHeap(); const outputStream = fs.createWriteStream(outputFile); let completedStreams = 0; let isFirstWrite = true; // 控制输出文件的逗号分隔格式 // 初始化所有临时文件流,读取第一个词入堆 for (let i = 0; i < chunkCount; i++) { const filePath = path.join(TEMP_DIR, `chunk_${i}.txt`); const stream = fs.createReadStream(filePath, { highWaterMark: 16 * 1024 }); let buffer = ''; const processBuffer = () => { while (true) { const commaIndex = buffer.indexOf(','); if (commaIndex === -1) { if (stream.readableEnded) { if (buffer.trim()) heap.push({ word: buffer, stream, index: i }); break; } // 等待下一块数据 stream.once('data', (chunk) => { buffer += chunk.toString(); processBuffer(); }); break; } const word = buffer.slice(0, commaIndex); buffer = buffer.slice(commaIndex + 1); heap.push({ word, stream, index: i }); } }; stream.on('data', (chunk) => { buffer += chunk.toString(); processBuffer(); }); stream.on('end', () => { completedStreams++; if (completedStreams === chunkCount && heap.isEmpty()) { outputStream.end(); } }); } // 循环取出最小词写入输出文件 const writeNext = () => { if (heap.isEmpty()) { if (completedStreams === chunkCount) outputStream.end(); return; } const minElement = heap.pop(); // 处理逗号分隔:第一个词不加前置逗号,后续词加 const writeContent = isFirstWrite ? minElement.word : `,${minElement.word}`; isFirstWrite = false; outputStream.write(writeContent, (err) => { if (err) throw err; writeNext(); }); }; // 等待写入完成 await new Promise((resolve) => { outputStream.on('finish', resolve); writeNext(); }); // 清理临时文件与目录 for (let i = 0; i < chunkCount; i++) { fs.unlinkSync(path.join(TEMP_DIR, `chunk_${i}.txt`)); } fs.rmdirSync(TEMP_DIR); } // 主执行函数 async function main() { const outputFile = './sorted_output.txt'; const chunkCount = await splitAndSortChunks(); await mergeChunks(chunkCount, outputFile); console.log('排序完成,输出文件:', outputFile); } main().catch(err => console.error('排序失败:', err));
四、关键注意事项
- 内存管控:拆分阶段的内存阈值需预留足够系统开销(如示例中的700MB),避免内存溢出。
- 词完整性:必须基于逗号位置拆分,绝对不能截断单个词,否则会导致排序逻辑错误。
- 堆效率:最小堆的时间复杂度为O(log n)(n为临时文件数),是多路归并的核心高效保障。
- 流处理:全程使用NodeJS流API,避免一次性加载大文件,确保内存占用稳定在限制内。
- 临时文件清理:排序完成后及时删除临时文件,避免磁盘空间浪费。
内容的提问来源于stack exchange,提问作者Qwes
相关产品推荐
相关产品推荐

