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

导入大型CSV至MongoDB遇MongoExpiredSessionError,如何完成全量上传?

问题描述

我使用csv-parser包处理大量CSV数据并尝试存储到MongoDB数据库中,部分数据已成功保存,但处理数百条记录后遇到error_1.MongoExpiredSessionError()错误,导致数据插入流程中断。我通过终端执行tsc .\importdata.ts -> node .\importdata.js运行函数,请问有没有办法让函数等待所有数据上传完成?

原代码:

const fs = require('fs');
const csv = require('csv-parser');
const { MongoClient } = require('mongodb');

interface JourneyData {
  Departure: string;
  Return: string;
  'Departure station id': string;
  'Departure station name': string;
  'Return station id': string;
  'Return station name': string;
  'Covered distance (m)': string;
  'Duration (sec.)': string;
}

async function importData() {
  const client = new MongoClient(
    'mongodb+srv://hasan....:.......@clustersolita.2ztdulk.mongodb.net/',
    { useNewUrlParser: true, useUnifiedTopology: true }
  );
  await client.connect();
  const db = client.db('CityBike');
  const journeysCollection = db.collection('journeys');

  try {
    await journeysCollection.deleteMany({});

    fs.createReadStream('2021-05.csv')
      .pipe(csv({ batchSize: 100 }))
      .on('data', async (row: JourneyData) => {
        const departureTime = row.Departure;
        const returnTime = row.Return;
        const departureStationId = parseInt(row['Departure station id'], 10);
        const departureStationName = row['Departure station name'];
        const returnStationId = parseInt(row['Return station id'], 10);
        const returnStationName = row['Return station name'];
        const coveredDistance = parseInt(row['Covered distance (m)'], 10);
        const duration = parseInt(row['Duration (sec.)'], 10);

        if (duration >= 10 && coveredDistance >= 10) {
          const journey = {
            Departure: departureTime,
            Return: returnTime,
            'Departure station id': departureStationId,
            'Departure station name': departureStationName,
            'Return station id': returnStationId,
            'Return station name': returnStationName,
            'Covered distance': coveredDistance,
            Duration: duration,
          };

          await journeysCollection.insertOne(journey);
        }
      })
      .on('end', async () => {
        console.log('Data import completed.');
        client.close();
      });
  } catch (error) {
    console.error('An error occurred:', error);
    client.close();
  }
}

importData().catch((err) => console.log(err));
解决方案

出现MongoExpiredSessionError的核心原因是:流的data事件会异步触发,而你在回调里用async/await执行insertOne时,流不会等待当前插入完成就继续触发下一个data事件,导致大量并发请求耗尽MongoDB会话,最终会话过期后插入操作失败。

以下是两种可行的修复方案:

方案一:批量插入+Promise包装流

通过批量攒数据再插入,减少数据库请求次数,同时把流操作包装成Promise,让函数等待所有处理完成:

const fs = require('fs');
const csv = require('csv-parser');
const { MongoClient } = require('mongodb');

interface JourneyData {
  Departure: string;
  Return: string;
  'Departure station id': string;
  'Departure station name': string;
  'Return station id': string;
  'Return station name': string;
  'Covered distance (m)': string;
  'Duration (sec.)': string;
}

async function importData() {
  const client = new MongoClient(
    'mongodb+srv://hasan....:.......@clustersolita.2ztdulk.mongodb.net/',
    { useNewUrlParser: true, useUnifiedTopology: true }
  );
  await client.connect();
  const db = client.db('CityBike');
  const journeysCollection = db.collection('journeys');
  const batchSize = 100;
  let batch: any[] = [];

  try {
    await journeysCollection.deleteMany({});

    // 用Promise包裹流操作,确保函数等待所有数据处理完成
    await new Promise((resolve, reject) => {
      fs.createReadStream('2021-05.csv')
        .pipe(csv())
        .on('data', async (row: JourneyData) => {
          const departureTime = row.Departure;
          const returnTime = row.Return;
          const departureStationId = parseInt(row['Departure station id'], 10);
          const departureStationName = row['Departure station name'];
          const returnStationId = parseInt(row['Return station id'], 10);
          const returnStationName = row['Return station name'];
          const coveredDistance = parseInt(row['Covered distance (m)'], 10);
          const duration = parseInt(row['Duration (sec.)'], 10);

          if (duration >= 10 && coveredDistance >= 10) {
            const journey = {
              Departure: departureTime,
              Return: returnTime,
              'Departure station id': departureStationId,
              'Departure station name': departureStationName,
              'Return station id': returnStationId,
              'Return station name': returnStationName,
              'Covered distance': coveredDistance,
              Duration: duration,
            };
            batch.push(journey);

            // 达到批量大小执行插入
            if (batch.length >= batchSize) {
              await journeysCollection.insertMany(batch);
              batch = [];
            }
          }
        })
        .on('end', async () => {
          // 插入剩余的不足批量的数据
          if (batch.length > 0) {
            await journeysCollection.insertMany(batch);
          }
          console.log('Data import completed.');
          await client.close();
          resolve(null);
        })
        .on('error', async (err) => {
          console.error('Stream error:', err);
          await client.close();
          reject(err);
        });
    });
  } catch (error) {
    console.error('An error occurred:', error);
    await client.close();
  }
}

importData().catch((err) => console.log(err));

方案二:使用stream.pipeline+异步迭代器

利用Node.js的stream/promises模块的pipeline方法,结合异步迭代器控制流的处理顺序,确保每条数据插入完成后再处理下一条:

const fs = require('fs');
const csv = require('csv-parser');
const { MongoClient } = require('mongodb');
const { pipeline } = require('stream/promises');

interface JourneyData {
  Departure: string;
  Return: string;
  'Departure station id': string;
  'Departure station name': string;
  'Return station id': string;
  'Return station name': string;
  'Covered distance (m)': string;
  'Duration (sec.)': string;
}

async function importData() {
  const client = new MongoClient(
    'mongodb+srv://hasan....:.......@clustersolita.2ztdulk.mongodb.net/',
    { useNewUrlParser: true, useUnifiedTopology: true }
  );
  await client.connect();
  const db = client.db('CityBike');
  const journeysCollection = db.collection('journeys');

  try {
    await journeysCollection.deleteMany({});

    await pipeline(
      fs.createReadStream('2021-05.csv'),
      csv(),
      async function* (source) {
        for await (const row of source) {
          const departureTime = row.Departure;
          const returnTime = row.Return;
          const departureStationId = parseInt(row['Departure station id'], 10);
          const departureStationName = row['Departure station name'];
          const returnStationId = parseInt(row['Return station id'], 10);
          const returnStationName = row['Return station name'];
          const coveredDistance = parseInt(row['Covered distance (m)'], 10);
          const duration = parseInt(row['Duration (sec.)'], 10);

          if (duration >= 10 && coveredDistance >= 10) {
            const journey = {
              Departure: departureTime,
              Return: returnTime,
              'Departure station id': departureStationId,
              'Departure station name': departureStationName,
              'Return station id': returnStationId,
              'Return station name': returnStationName,
              'Covered distance': coveredDistance,
              Duration: duration,
            };
            await journeysCollection.insertOne(journey);
          }
        }
      }
    );

    console.log('Data import completed.');
    await client.close();
  } catch (error) {
    console.error('An error occurred:', error);
    await client.close();
  }
}

importData().catch((err) => console.log(err));

关键改进点

  • 用Promise或stream.pipeline包装流操作,确保importData函数等待所有数据处理完成再结束
  • 用insertMany替代多次insertOne,减少数据库请求次数,降低会话资源消耗
  • 新增流的error事件处理,避免异常导致数据库连接泄漏
  • 确保client.close()在所有操作完成后执行,不会提前关闭MongoDB连接

内容的提问来源于stack exchange,提问作者Hasan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 19:17:53