Node.js读取CSV流二次调用报错Failed to pipe. The response has been emitted already如何处理
问题根因
你遇到的报错本质是Node.js 可读流(Readable)是一次性消费资源:
- 第一次调用
parseStream传入this.source时,fast-csv 会自动 pipe 这个由 got 创建的 http 流,第一次读取结束后,流就进入ended/closed状态,所有数据已经被消费完毕,第二次再传入同一个流实例执行解析,自然无法读取内容,就抛出了The response has been emitted already的错误。 - 之前回调版本没有报错,大概率是测试时没有触发流完全消费,或者每次调用隐式创建了新的流实例,改为 Promise 版本后复用了存在
this.source的全局流实例,才暴露了问题。
无需每次重新打开流的解决方案
根据你的场景不同,可以选择两种方案:
方案1:资源预缓存(适合CSV体积不大的场景)
提前把整个CSV文件内容缓存到内存,每次调用GetPage时从缓存创建新的可读流,无需重复发起http请求:
class CSVStreamParser { private buffer: Buffer; private currentLocation = 0; // 其他属性省略 // 构造函数提前拉取全量资源缓存 constructor(url: string, private maxArrayLength: number, private delimiter: string) { this.init(url); } private async init(url: string) { // 一次性拉取全量CSV内容缓存到Buffer this.buffer = await got(url).buffer(); } async GetPage() : Promise<{OutputArray:any[], StartingIndex:number}>{ return new Promise((resolve, reject) => { const output:any[] = []; const startingIndex = this.currentLocation; // 每次调用从缓存Buffer创建全新的可读流 const source = Readable.from(this.buffer); try{ parseStream(source, {headers:true, maxRows:this.maxArrayLength, skipRows:this.currentLocation, ignoreEmpty:true, delimiter:this.delimiter}) .on('error', error => { console.log(`parseStream: ${error}`); // 异步错误要传给reject,不要直接throw reject(error); }) .on('data', row => { const obj = this.unflatten(row); output.push(obj); this.currentLocation++; }) .on('end', (rowCount: number) => { console.log(`Parsed ${this.currentLocation} rows`); resolve({OutputArray:output, StartingIndex:startingIndex}); }); } catch(ex){ console.log(`parseStream: ${ex}`); reject(ex); } }) } // unflatten等其他方法省略 }
方案2:流式分页队列(适合大体积CSV、不想全量缓存的场景)
初始化时就启动解析流,把解析后的行存入内存队列,调用GetPage时直接从队列取指定长度的行返回,源流全程只消费一次:
class CSVStreamParser { private parseQueue: any[] = []; private waitingResolve: null | ((value: any) => void) = null; private streamEnded = false; private currentLocation = 0; // 其他属性省略 constructor(url: string, private maxArrayLength: number, private delimiter: string) { // 初始化时启动流解析,数据持续存入队列 const source = got.stream(url); parseStream(source, {headers:true, ignoreEmpty:true, delimiter:this.delimiter}) .on('error', error => { console.log(`parseStream: ${error}`); if(this.waitingResolve) { this.waitingResolve(Promise.reject(error)); this.waitingResolve = null; } }) .on('data', row => { const obj = this.unflatten(row); this.parseQueue.push(obj); // 队列数据量满足分页要求、且有等待的GetPage请求,直接返回 if(this.waitingResolve && this.parseQueue.length >= this.maxArrayLength) { this.returnPage(); } }) .on('end', () => { this.streamEnded = true; // 流结束后把剩余不足一页的数据返回给等待的请求 if(this.waitingResolve) { this.returnPage(); } }) } private returnPage() { const output = this.parseQueue.splice(0, this.maxArrayLength); const startingIndex = this.currentLocation; this.currentLocation += output.length; this.waitingResolve!({OutputArray: output, StartingIndex: startingIndex}); this.waitingResolve = null; } async GetPage() : Promise<{OutputArray:any[], StartingIndex:number}> { // 队列有足够数据或流已结束,直接返回当前页 if(this.parseQueue.length >= this.maxArrayLength || this.streamEnded) { const output = this.parseQueue.splice(0, this.maxArrayLength); const startingIndex = this.currentLocation; this.currentLocation += output.length; return {OutputArray: output, StartingIndex: startingIndex}; } // 队列数据不足,等待新数据推送后返回 return new Promise(resolve => { this.waitingResolve = resolve; }) } // unflatten等其他方法省略 }
原有代码的优化点
- 不要给Promise的executor函数加
async关键字,你这里没有用到await,加了反而会导致内部抛出的错误无法被正常捕获 - 流的异步错误不能通过外层try/catch捕获,要在error事件回调里调用
reject(error),否则会出现未处理的Promise rejection。
内容的提问来源于stack exchange,提问作者Oded Answer
相关产品推荐
相关产品推荐

