Node.js克隆可读流多任务处理时MD5校验值为空如何解决
问题原因
- 校验值传参时机错误:调用
flushStream时passthroughStreamcheckSum的流还没读完,checkSum还是初始空字符串,字符串是基础类型传值,后续校验和计算完成后更新checkSum变量,不会同步到已经传入flushStream的参数里 - 现有
flushStream逻辑不符合需求:仅写入了流的原始内容,没有在文件末尾追加校验和 - 代码存在笔误与结构问题:调用的
processEacchRecord和实际定义的processRecords方法名不匹配;createTXTFile中定义的对象有重复的ckey,会导致校验和字段被覆盖 - 等待时序不合理:行数统计的结束事件没有和校验和计算的结束事件并行等待,可能出现统计结果还没生成就提前写TXT的问题
修复后的实现
import { createHash } from 'crypto' import { PassThrough } from 'stream' import * as readline from 'readline' import { once } from 'events' import * as fs from 'fs' // 原有Streamable等类型定义保持不变即可 async processStream(streamable: Streamable): Promise<number> { // 拆分三个分流分别处理校验计算、行数统计、文件写入 const passthroughForLineCount = new PassThrough() const passthroughForChecksum = new PassThrough() const passthroughForFileWrite = new PassThrough() let count = 0 let countError = 0 const md5Checksum = createHash("md5") // 1. 处理行数统计和JSON格式校验 const lineStream: readline.Interface = readline.createInterface({ input: passthroughForLineCount }) lineStream.on("line", line => { if (!line) return count++ try { JSON.parse(line) } catch (error) { countError++ } }) const lineCountPromise = once(lineStream, "close") // 2. 处理MD5校验和计算 const checksumPromise = new Promise<string>(resolve => { passthroughForChecksum.on("data", data => { md5Checksum.update(data) }) passthroughForChecksum.on("end", () => { resolve(md5Checksum.digest("hex")) }) }) // 3. 先把原始流内容写入JSON文件 const jsonFilePath = '你需要保存的JSON文件路径.json' const fileWritePromise = new Promise<void>((resolve, reject) => { const fileOut = fs.createWriteStream(jsonFilePath) passthroughForFileWrite.pipe(fileOut) fileOut.on('finish', resolve) fileOut.on('error', reject) }) // 原始流分流 streamable.stream.pipe(passthroughForChecksum) streamable.stream.pipe(passthroughForLineCount) streamable.stream.pipe(passthroughForFileWrite) // 并行等待所有流处理完成 const [finalChecksum] = await Promise.all([ checksumPromise, lineCountPromise, fileWritePromise ]) // 4. 往JSON文件末尾追加校验和(按需调整格式,这里默认单独占一行) await fs.promises.appendFile(jsonFilePath, `\n${finalChecksum}`) // 5. 写入统计结果TXT文件 await this.createTXTFile(finalChecksum, countError, count) return count } private async createTXTFile( checksum: string, errorCount: number, totalCount: number ) { // 修复重复key的问题 const data = { checksum, errorCount, totalCount } await this.flushString('你需要保存的统计文件路径.txt', JSON.stringify(data)) }
内容的提问来源于stack exchange,提问作者user5740953
相关产品推荐
相关产品推荐

