如何复用已传递给类方法的Readable流,无需重新生成?
Node.js Readable流复用解决方案
Readable流的核心特性是单向消费:一旦数据被读取完毕(比如FTP上传过程中读取了整个流),流的内部指针会移动到末尾,后续read()调用会返回null,这就是你遇到数据丢失的根本原因。以下是几种无需重新生成原始流即可复用数据的方案:
方案一:缓存流数据到内存,生成多份可读流
先将原始流的所有数据收集到内存Buffer中,再用这个Buffer创建多个独立的Readable流,分别供给FTP上传和前端返回:
const readable: Readable = await createCsv(csv); // 收集所有流数据到内存 const chunks: Buffer[] = []; for await (const chunk of readable) { chunks.push(chunk); } const csvBuffer = Buffer.concat(chunks); // 用于FTP上传的流 const ftpStream = Readable.from(csvBuffer); await this.ftp.uploadCSV(ftpStream, pathway); // 返回给前端的流 const responseStream = Readable.from(csvBuffer); return { readable: responseStream };
适用场景:CSV数据量较小,不会造成内存占用过高的情况。
方案二:封装流工厂函数,按需生成新流
把创建CSV流的逻辑封装成工厂函数,需要时调用生成新的可读流,避免重复写生成逻辑:
// 定义工厂函数,封装CSV流创建逻辑 const createCsvStream = async () => await createCsv(csv); // 生成流用于FTP上传 const ftpStream = await createCsvStream(); await this.ftp.uploadCSV(ftpStream, pathway); // 生成新流返回给前端 const responseStream = await createCsvStream(); return { readable: responseStream };
适用场景:CSV数据量较大,不想缓存到内存,但能接受重复生成流的计算开销。
方案三:使用PassThrough分流,同时处理多个消费方
通过stream.PassThrough创建中转流,让原始流的数据同时流向FTP客户端和缓存(或另一个中转流),实现一次生成原始流,同时满足多个消费需求:
const { PassThrough } = require('stream'); const readable: Readable = await createCsv(csv); const passThrough = new PassThrough(); const chunks: Buffer[] = []; // 收集中转流的数据用于后续返回前端 passThrough.on('data', chunk => chunks.push(chunk)); // 将原始流导入中转流 readable.pipe(passThrough); // 用中转流完成FTP上传 await this.ftp.uploadCSV(passThrough, pathway); // 用缓存的数据创建返回前端的流 const responseStream = Readable.from(Buffer.concat(chunks)); return { readable: responseStream };
如果你的FTP客户端支持同时处理流,也可以用Promise.all让两个消费方同时处理:
const { PassThrough } = require('stream'); const readable: Readable = await createCsv(csv); const ftpStream = new PassThrough(); const responseStream = new PassThrough(); // 原始流同时分流到两个中转流 readable.pipe(ftpStream); readable.pipe(responseStream); // 等待FTP上传和流处理完成 await Promise.all([ this.ftp.uploadCSV(ftpStream, pathway), new Promise(resolve => responseStream.on('end', resolve)) ]); return { readable: responseStream };
适用场景:需要一次性处理原始流,且不想重复生成流或缓存大量数据的场景。
内容的提问来源于stack exchange,提问作者Ctfrancia
相关产品推荐
相关产品推荐

