导入大型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
相关产品推荐
相关产品推荐

