【问题标题】:Save a very big CSV to mongoDB using mongoose使用 mongoose 将一个非常大的 CSV 保存到 mongoDB
【发布时间】:2014-07-31 09:08:01
【问题描述】:

我有一个包含超过 200'000 行的 CSV 文件。我需要将它保存到 MongoDB。

如果我尝试 for 循环,Node 将耗尽内存。

fs.readFile('data.txt', function(err, data) {
  if (err) throw err;

  data.split('\n');

  for (var i = 0; i < data.length, i += 1) {
    var row = data[i].split(',');

    var obj = { /* The object to save */ }

    var entry = new Entry(obj);
    entry.save(function(err) {
      if (err) throw err;
    }
  } 
}

如何避免内存不足?

【问题讨论】:

标签: javascript node.js mongodb mongoose


【解决方案1】:

欢迎使用流媒体。您真正想要的是一个“事件流”,它“一次一个块”处理您的输入,当然理想情况下是通过一个常见的分隔符,例如您当前使用的“换行符”。

对于真正高效的东西,您可以添加对 MongoDB "Bulk API" 插入的使用,以使您的加载尽可能快,而不会耗尽所有机器内存或 CPU 周期。

不提倡,因为有各种可用的解决方案,但这里有一个列表,它利用 line-input-stream package 使“行终止符”部分变得简单。

仅通过“示例”定义架构:

var LineInputStream = require("line-input-stream"),
    fs = require("fs"),
    async = require("async"),
    mongoose = require("mongoose"),
    Schema = mongoose.Schema;

var entrySchema = new Schema({},{ strict: false })

var Entry = mongoose.model( "Schema", entrySchema );

var stream = LineInputStream(fs.createReadStream("data.txt",{ flags: "r" }));

stream.setDelimiter("\n");

mongoose.connection.on("open",function(err,conn) { 

    // lower level method, needs connection
    var bulk = Entry.collection.initializeOrderedBulkOp();
    var counter = 0;

    stream.on("error",function(err) {
        console.log(err); // or otherwise deal with it
    });

    stream.on("line",function(line) {

        async.series(
            [
                function(callback) {
                    var row = line.split(",");     // split the lines on delimiter
                    var obj = {};             
                    // other manipulation

                    bulk.insert(obj);  // Bulk is okay if you don't need schema
                                       // defaults. Or can just set them.

                    counter++;

                    if ( counter % 1000 == 0 ) {
                        stream.pause();
                        bulk.execute(function(err,result) {
                            if (err) callback(err);
                            // possibly do something with result
                            bulk = Entry.collection.initializeOrderedBulkOp();
                            stream.resume();
                            callback();
                        });
                    } else {
                        callback();
                    }
               }
           ],
           function (err) {
               // each iteration is done
           }
       );

    });

    stream.on("end",function() {

        if ( counter % 1000 != 0 )
            bulk.execute(function(err,result) {
                if (err) throw err;   // or something
                // maybe look at result
            });
    });

});

因此,通常那里的“流”接口会“分解输入”以处理“一次一行”。这会阻止您一次加载所有内容。

主要部分是来自 MongoDB 的 "Bulk Operations API"。这允许您在实际发送到服务器之前一次“排队”许多操作。因此,在这种使用“模”的情况下,仅每处理 1000 个条目发送写入。你真的可以做任何事情,直到 16MB BSON 限制,但要让它易于管理。

除了批量处理的操作之外,async 库中还有一个额外的“限制器”。这并不是真正需要的,但这可确保在任何时候处理的文档基本上不超过“模数限制”。一般的批处理“插入”除了内存之外没有 IO 成本,但“执行”调用意味着 IO 正在处理。所以我们等待而不是排队。

对于“流处理”CSV 类型的数据,您肯定可以找到更好的解决方案。但总的来说,这为您提供了如何以内存高效的方式执行此操作而不占用 CPU 周期的概念。

【讨论】:

  • async.series 如何确保任何时候处理的文档不超过“模数限制”?我在这里看到的是对象在 1000 秒内排队并写入 db。如果接下来的 1000 个对象在前一个 bulk.execute() 完成之前准备好,它将触发另一个 bulk.execute()。我在这里错过了什么吗?
  • @JayyVis 显然是的,你错过了一些东西。正如我所说的“系列”有点做作,但实际执行是有道理的。使用模数,您一次不能“排队”超过 1000 个操作。 “系列”回调确保了这一点,因为“内部”执行需要完成。这要么是“批量插入”,要么是有效地排空队列的实际“执行”。这通常是您想要避免内存消耗的模式。加上事件滴答声的所有过程。
  • @JayyVis 不,你没有。 “异步”点是“暂停”进一步处理,直到完成。所以这可以在不同的层次上处理。也许您应该在发布“建设性批评”之前先尝试代码
  • 确实@JayKumar 的回答是正确的。执行命令是异步的,在执行操作完成之前可能会出现更多行事件,最终使用正在执行的同一个批量对象。当执行完成时,如果在开始执行之前对先前的批量对象有任何插入操作,那么在启动新的批量对象时,所有这些文档都会丢失。这假设将文档插入到正在执行的批量对象中不会引发任何异常。
【解决方案2】:

公认的答案很好,并试图涵盖这个问题的所有重要方面。

  1. 以行流的形式读取 CSV 文件
  2. 将文档批量写入MongoDB
  3. 读写同步

虽然它在前两个方面做得很好,但使用 async.series() 解决同步问题的方法无法按预期工作。

stream.on("line",function(line) {
    async.series(
        [
            function(callback) {
                var row = line.split(",");     // split the lines on delimiter
                var obj = {};             
                // other manipulation

                bulk.insert(obj);  // Bulk is okay if you don't need schema
                                   // defaults. Or can just set them.

                counter++;

                if ( counter % 1000 == 0 ) {
                    bulk.execute(function(err,result) {
                        if (err) throw err;   // or do something
                        // possibly do something with result
                        bulk = Entry.collection.initializeOrderedBulkOp();
                        callback();
                    });
                } else {
                    callback();
                }
           }
       ],
       function (err) {
           // each iteration is done
       }
   );
});

这里的 bulk.execute() 是一个 mongodb 写操作,它是一个异步 IO 调用。这允许 node.js 在 bulk.execute() 完成其数据库写入和回调之前继续执行事件循环。

因此它可能会继续从流中接收更多的“行”事件并将更多文档排队bulk.insert(obj),并且可以点击下一个模以再次触发 bulk.execute()。

让我们看看这个例子。

var async = require('async');

var bulk = {
    execute: function(callback) {
        setTimeout(callback, 1000);
    }
};

async.series(
    [
       function (callback) {
           bulk.execute(function() {
              console.log('completed bulk.execute');
              callback(); 
           });
       },
    ], 
    function(err) {

    }
);

console.log("!!! proceeding to read more from stream");

输出

!!! proceeding to read more from stream
completed bulk.execute

要真正确保我们在任何给定时间处理一批 N 个文档,我们需要使用 stream.pause() 和 stream.resume() 对文件流实施流控制

var LineInputStream = require("line-input-stream"),
    fs = require("fs"),
    mongoose = require("mongoose"),
    Schema = mongoose.Schema;

var entrySchema = new Schema({},{ strict: false });
var Entry = mongoose.model( "Entry", entrySchema );

var stream = LineInputStream(fs.createReadStream("data.txt",{ flags: "r" }));

stream.setDelimiter("\n");

mongoose.connection.on("open",function(err,conn) { 

    // lower level method, needs connection
    var bulk = Entry.collection.initializeOrderedBulkOp();
    var counter = 0;

    stream.on("error",function(err) {
        console.log(err); // or otherwise deal with it
    });

    stream.on("line",function(line) {
        var row = line.split(",");     // split the lines on delimiter
        var obj = {};             
        // other manipulation

        bulk.insert(obj);  // Bulk is okay if you don't need schema
                           // defaults. Or can just set them.

        counter++;

        if ( counter % 1000 === 0 ) {
            stream.pause(); //lets stop reading from file until we finish writing this batch to db

            bulk.execute(function(err,result) {
                if (err) throw err;   // or do something
                // possibly do something with result
                bulk = Entry.collection.initializeOrderedBulkOp();

                stream.resume(); //continue to read from file
            });
        }
    });

    stream.on("end",function() {
        if ( counter % 1000 != 0 ) {
            bulk.execute(function(err,result) {
                if (err) throw err;   // or something
                // maybe look at result
            });
        }
    });

});

【讨论】:

    猜你喜欢
    • 2018-01-20
    • 2022-01-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-27
    • 2019-01-18
    • 2016-02-15
    • 2016-12-11
    相关资源
    最近更新 更多