JavaScript读取流时数组无法重置的问题排查
问题:批量插入CSV数据到MongoDB时数组无法清空,持续增长
我尝试将远程CSV文件中的数据插入MongoDB(使用Mongoose),希望每次批量插入100条数据。但运行时发现productsToUpsert数组在达到100条后仍持续增长,无法按预期清空并重复批量插入流程。
原代码如下:
import csv from 'csv-parser' import fetch from 'node-fetch' import { Product } from '../models/Product' export async function handleCSVProcessing(targetCsv: string) { const batchSize = 100 try { const response = await fetch(targetCsv, { method: 'get', headers: { 'content-type': 'text/csv;charset=UTF-8', } }) if (!response || !response.ok || !response.body) { throw new Error(`Failed to fetch CSV file: ${response.statusText}`) } const productsToUpsert: any[] = [] response.body .pipe(csv()) .on('data', async (row: any) => { productsToUpsert.push({ // 转换逻辑此处省略 }) if (productsToUpsert.length >= batchSize) { console.log('productsToUpsert.length: ' + productsToUpsert.length) // 此处长度持续增长到101、102... await Product.bulkWrite(productsToUpsert); productsToUpsert.length = 0 // 理论上应该清空数组,但实际无效 } }) .on('end', async () => { if (productsToUpsert.length > 0) { await Product.bulkWrite(productsToUpsert); } }) .on('error', (error: any) => { console.error('Error processing CSV:', error); }); } catch (error) { console.error('Error fetching CSV:', error); } }
问题原因
核心问题是流的data事件不会等待你的异步回调完成。当你在data回调中执行await Product.bulkWrite()时,csv-parser的流仍在持续读取CSV内容并触发新的data事件,新的行数据会不断被push到productsToUpsert数组中,导致数组在批量插入过程中继续增长,无法按预期清空。
解决方案
需要在触发批量插入时暂停流,等待数据库操作完成后再恢复流,避免新数据在插入过程中被加入数组。修改后的代码如下:
import csv from 'csv-parser' import fetch from 'node-fetch' import { Product } from '../models/Product' export async function handleCSVProcessing(targetCsv: string) { const batchSize = 100 try { const response = await fetch(targetCsv, { method: 'get', headers: { 'content-type': 'text/csv;charset=UTF-8', } }) if (!response || !response.ok || !response.body) { throw new Error(`Failed to fetch CSV file: ${response.statusText}`) } const productsToUpsert: any[] = [] const csvStream = response.body.pipe(csv()) csvStream .on('data', async (row: any) => { productsToUpsert.push({ // 转换逻辑此处省略 }) if (productsToUpsert.length >= batchSize) { // 暂停流,避免继续读取新数据 csvStream.pause() try { console.log('执行批量插入,当前数组长度:', productsToUpsert.length) await Product.bulkWrite(productsToUpsert) // 清空数组 productsToUpsert.length = 0 } catch (dbError) { console.error('批量插入失败:', dbError) } finally { // 恢复流,继续处理下一批数据 csvStream.resume() } } }) .on('end', async () => { // 处理剩余不足batchSize的数据 if (productsToUpsert.length > 0) { try { await Product.bulkWrite(productsToUpsert) } catch (dbError) { console.error('剩余数据插入失败:', dbError) } } console.log('CSV数据处理完成') }) .on('error', (error: any) => { console.error('CSV处理出错:', error) }); } catch (error) { console.error('获取CSV文件失败:', error) } }
关键修改点
- 保存csv流的引用
const csvStream = response.body.pipe(csv()),方便后续控制流的暂停/恢复 - 在达到批量阈值时,先调用
csvStream.pause()暂停流读取 - 批量插入操作放在
try/catch中,确保无论成功或失败都能恢复流 - 插入完成后清空数组,再调用
csvStream.resume()恢复流处理
内容的提问来源于stack exchange,提问作者mrodo
相关产品推荐
相关产品推荐

