解密解压CSV流仅返回首行后挂起,寻求技术解决方案
问题:解密解压CSV文件仅读取第一行后程序挂起
背景
原有代码可正常读取加密文件,解密解压后作为API响应返回:
let decipher = crypto.createDecipheriv(type, file.metadata?.key, file.metadata?.iv) let decompress = new fflate.Decompress() let decipherStream = fs.createReadStream(location).pipe(decipher) decipherStream.on('data', (data) => decompress.push(data)) decipherStream.on('finish', () => decompress.push(new Uint8Array(), true)) decompress.ondata = (data: any, final: any) => { if (!final) res.write(data) if (final) res.send() } res.attachment(file.originalName)
修改为非API处理流程后,尝试读取解密解压后的CSV内容,但仅能获取第一行数据,之后程序直接挂起,尝试的代码如下:
let decipher = crypto.createDecipheriv(type, file.metadata?.key, file.metadata?.iv) let decompress = new fflate.Decompress() let decipherStream = fs.createReadStream(location).pipe(decipher) decipherStream.on('data', (data) => decompress.push(data)) decipherStream.on('finish', () => decompress.push(new Uint8Array(), true)) let stream = new Readable() stream._read = function (){} decompress.ondata = (data: any, final: any) => { console.log(data) if (!final) stream.push(data) if (final) stream.push(null) } let headernames:any = [] let jsonobjs:any = [] let jsonout = '' await csv().fromStream(stream) .on('headers', (headers: any) => { headernames = headers; console.log("Headers are", headernames); }) .on('data', (data: any) => { jsonout = data.toString('utf8'); jsonobjs.push(jsonout); }) .on('done', () => { console.log("jsonout:",jsonout); return jsonobjs; })
问题根源
- 自定义Readable流未处理背压:空实现的
_read方法无法触发流的后续数据推送,CSV解析器读取第一行后,流无法继续提供数据,导致阻塞。 - CSV解析器的
await误用:csv().fromStream()返回的是事件发射器而非Promise,直接await无法等待解析完成,后续逻辑执行异常。 - 解压数据推送未考虑流缓冲区:直接推送数据到自定义流,当解析速度跟不上推送速度时,缓冲区溢出导致程序挂起。
修复代码
const { Readable } = require('stream'); const crypto = require('crypto'); const fflate = require('fflate'); const fs = require('fs'); const csv = require('csv-parser'); async function processEncryptedCSV(location, type, file) { return new Promise((resolve, reject) => { const decipher = crypto.createDecipheriv(type, file.metadata?.key, file.metadata?.iv); const decompress = new fflate.Decompress(); // 定义支持背压的可读流 const stream = new Readable({ read() {} }); const headernames = []; const jsonobjs = []; // 处理CSV解析逻辑 stream.pipe(csv()) .on('headers', headers => { headernames.push(...headers); console.log("Headers are", headernames); }) .on('data', data => { jsonobjs.push(data); }) .on('end', () => { console.log("CSV解析完成"); resolve({ headernames, jsonobjs }); }) .on('error', err => reject(err)); // 处理解密流数据 const decipherStream = fs.createReadStream(location).pipe(decipher); decipherStream.on('data', data => { // 处理背压:推送失败时暂停解密流,缓冲区空闲后恢复 if (!stream.push(data)) { decipherStream.pause(); stream.once('drain', () => decipherStream.resume()); } decompress.push(data); }); decipherStream.on('finish', () => decompress.push(new Uint8Array(), true)); decipherStream.on('error', err => reject(err)); // 处理解压后的数据 decompress.ondata = (data, final) => { if (!final) { if (!stream.push(data)) { decipherStream.pause(); stream.once('drain', () => decipherStream.resume()); } } else { stream.push(null); // 结束可读流 } }; decompress.onerror = err => reject(err); }); } // 使用示例 // processEncryptedCSV('/path/to/encrypted/file', 'aes-256-cbc', fileMeta) // .then(result => console.log(result)) // .catch(err => console.error(err));
关键修改说明
- 用Promise包裹流程:确保能正确等待整个解析流程完成,避免异步逻辑混乱。
- 处理流背压:推送数据时检查返回值,缓冲区满时暂停解密流,空闲后恢复,防止数据溢出。
- 修正CSV解析逻辑:使用
csv-parser的end事件监听解析完成,data事件直接接收解析后的JSON对象,无需手动转字符串。 - 完善错误监听:给解密、解压、CSV解析各环节添加错误捕获,避免程序静默崩溃。
内容的提问来源于stack exchange,提问作者waffletree
相关产品推荐
相关产品推荐

