You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 08:17:26