使用Node.js Stream导入CSV到MongoDB出现重复数据如何解决?
问题根源
- 你错误地暂停了原始文件读取流
fs.createReadStream,而非经过@fast-csv/parse处理后的解析流。CSV解析器内部存在缓冲区,已经读取到缓冲区的内容即便暂停上游文件流,仍会继续触发data事件,导致暂停逻辑完全失效。 data事件的回调使用了异步async函数,Node.js流的事件机制本身不会等待异步回调执行完成。你调用pause()之前,事件循环队列中已经堆积了多个待触发的data事件,这些事件不会因pause调用停止执行,会继续往documents数组写入数据,甚至部分行被多次处理导致重复插入。- 当判断
documents.length >= 200时你先把当前行push进数组再插入,相当于每批实际插入201条数据,不符合你最初设计的200条批量插入的预期。
修复代码
module.exports.seedProductCollection = (filepath) => new Promise((resolve, reject) => { let documents = []; let counter = 0; // 直接赋值CSV解析流,后续暂停、恢复都操作该流 const csvStream = fs.createReadStream(filepath).pipe(csv.parse({ headers: true })); csvStream .on('error', (error) => reject(error)) .on('data', (row) => { documents.push(row); if (documents.length >= 200) { csvStream.pause(); // 用setImmediate确保暂停生效后再执行异步插入操作 setImmediate(async () => { try { await db.Product.insertMany(documents); counter += documents.length; console.log('counter: ', counter); documents = []; csvStream.resume(); } catch (error) { reject(error); } }); } }) .on('end', async () => { if (documents.length > 0) { try { await db.Product.insertMany(documents); counter += documents.length; } catch (error) { return reject(error); } } resolve(counter); }); });
可选优化
如果CSV文件体积极大,可在insertMany方法中传入第二个参数{ ordered: false },提升批量插入效率,单条数据错误不会中断整批插入,可根据业务需求选择是否开启。
内容的提问来源于stack exchange,提问作者h_a
相关产品推荐
相关产品推荐

