【问题标题】:Insert streamed XML data database插入流式 XML 数据数据库
【发布时间】:2017-08-05 22:48:48
【问题描述】:

我正在尝试有效地插入大量数据(XML 文件大小超过 70GB),而不会使我的 MongoDB 服务器崩溃。目前这就是我在 NodeJS 中使用 xml-stream 所做的事情:

var fs = require('fs'),
    path = require('path'),
    XmlStream = require('xml-stream'),
    MongoClient = require('mongodb').MongoClient,
    assert = require('assert'),
    ObjectId = require('mongodb').ObjectID,
    url = 'mongodb://username:password@my.server:27017/mydatabase',
    amount = 0;

var stream = fs.createReadStream(path.join(__dirname, 'motor.xml'));
var xml = new XmlStream(stream);

xml.collect('ns:Statistik');
xml.on('endElement: ns:Statistik', function(item) {
    var insertDocument = function(db, callback) {
        db.collection('vehicles').insertOne(item, function(err, result) {
            amount++;
            if (amount % 1000 == 0) {
                console.log("Inserted", amount);
            }
            callback();
        });
    };

    MongoClient.connect(url, function(err, db) {
        insertDocument(db, function() {
            db.close();
        });
    });
});

当我调用xml.on() 时,它基本上会返回我当前所在的树/元素。由于这是直接的 JSON,我可以将它作为参数提供给我的 db.collection().insertOne() 函数,它会按照我想要的方式将它插入到数据库中。

所有代码实际上都像现在一样工作,但是在大约 3000 次插入后它停止了(大约需要 10 秒)。我怀疑这是因为我每次在 XML 文件中看到一棵树时打开一个数据库连接,插入数据,然后关闭连接,在这种情况下大约 3000 次。

我可以以某种方式合并 insertMany() 函数并以 100 秒(或更多)为单位执行此操作,但我不太确定这将如何处理所有流式传输和异步操作。

所以我的问题是:如何将大量 XML(到 JSON)插入到我的 MongoDB 数据库中而不会崩溃?

【问题讨论】:

    标签: javascript node.js xml mongodb


    【解决方案1】:

    您认为.insertMany() 比每次都写要好,所以这只是收集"stream" 上的数据的问题。

    由于执行是“异步”的,您通常希望避免在堆栈上进行过多的活动调用,因此通常您在调用.insertMany() 之前.pause()"stream",然后在回调完成后.resume()完成:

    var fs = require('fs'),
        path = require('path'),
        XmlStream = require('xml-stream'),
        MongoClient = require('mongodb').MongoClient,
        url = 'mongodb://username:password@my.server:27017/mydatabase',
        amount = 0;
    
    MongoClient.connect(url, function(err, db) {
    
        var stream = fs.createReadStream(path.join(__dirname, 'motor.xml'));
        var xml = new XmlStream(stream);
    
        var docs = [];
        //xml.collect('ns:Statistik');
    
        // This is your event for the element matches
        xml.on('endElement: ns:Statistik', function(item) {
            docs.push(item);           // collect to array for insertMany
            amount++;
    
            if ( amount % 1000 === 0 ) { 
              xml.pause();             // pause the stream events
              db.collection('vehicles').insertMany(docs, function(err, result) {
                if (err) throw err;
                docs = [];             // clear the array
                xml.resume();          // resume the stream events
              });
            }
        });
    
        // End stream handler - insert remaining and close connection
        xml.on("end",function() {
          if ( amount % 1000 !== 0 ) {
            db.collection('vehicles').insertMany(docs, function(err, result) {
              if (err) throw err;
              db.close();
            });
          } else {
            db.close();
          }
        });
    
    });
    

    甚至对其进行一些现代化改造:

    const fs = require('fs'),
          path = require('path'),
          XmlStream = require('xml-stream'),
          MongoClient = require('mongodb').MongoClient;
    
    const uri = 'mongodb://username:password@my.server:27017/mydatabase';
    
    (async function() {
    
      let amount = 0,
          docs = [],
          db;
    
      try {
    
        db = await MongoClient.connect(uri);
    
        const stream = fs.createReadStream(path.join(__dirname, 'motor.xml')),
              xml = new XmlStream(stream);
    
        await Promise((resolve,reject) => {
          xml.on('endElement: ns:Statistik', async (item) => {
            docs.push(item);
            amount++;
    
            if ( amount % 1000 === 0 ) {
              try {
                xml.pause();
                await db.collection('vehicle').insertMany(docs);
                docs = [];
                xml.resume();
              } catch(e) {
                reject(e)
              }
            }
    
          });
    
          xml.on('end',resolve);
    
          xml.on('error',reject);
        });
    
        if ( amount % 1000 !== 0 ) {
          await db.collection('vehicle').insertMany(docs);
        }
    
      } catch(e) {
        console.error(e);
      } finally {
        db.close();
      }
    
    })();
    

    请注意,MongoClient 连接实际上包装了所有其他操作。您只想连接一次,其他操作发生在"stream" 的事件处理程序上。

    因此,对于您的XMLStream,事件处理程序会在表达式匹配时触发,并将数据提取并收集到一个数组中。每 1000 个项目调用.insertMany() 以插入文档,在“异步”调用上“暂停”和“恢复”。

    一旦完成,就会在"stream" 上触发“结束”事件。这是你关闭数据库连接的地方,事件循环将被释放并结束程序。

    虽然通过允许同时发生各种.insertMany() 调用(并且通常为“池大小”以免超出调用堆栈),有可能获得某种程度的“并行性”,但这基本上是过程的方式通过在等待其他异步 I/O 完成时简单地暂停来查看最简单的形式。

    注意:根据 follow up question 从原始代码中注释掉 .collect() 方法,这似乎没有必要,实际上是在内存中保留确实应该丢弃的节点每次写入数据库后。

    【讨论】:

    • 天哪,它看起来有效!我试图完成自己的工作,基本上做你做的事情,但我无法在开放的连接上解决问题。我的问题是,它给了我非常不一致的结果。如果我插入 1000 条记录,它实际上只会在数据库中显示 300 条(大约)。可能是因为我只是在连接完成之前随机关闭连接。非常感谢,尼尔!
    • 另一方面:你有什么线索为什么它真的开始了!大约 75000 次插入后速度变慢?当数据库为空时,我们说的是 1000/秒,但当我达到 75000 左右时,可能是 100-200/秒。
    • @MortenMoulder 您应该会看到使用.insertMany() 的显着改进,但至于吞吐量的一般差异取决于数据量,这是一个完全不同且非常广泛的主题。没有具体细节需要考虑的因素太多,例如有哪些索引(如果有的话)、可用内存、写入分布和基本硬件。如果您有其他问题,通常最好Ask a new Question,在那里您可以清楚地表达详细信息。
    • 没关系。我已经看到精度有了很大的提高,尽管如此。我现在从垃圾收集器那里得到一个“JavaScript 堆内存不足”,这可能暗示出了什么问题哈哈。谢谢!
    • @MortenMoulder 从上面的代码?因为我在清单中告诉您要做的是通过异步调用对.pause().resume() 进行流处理。因此,如果您在上面的列表中看到任何此类内容,那么唯一可能的区域是 xml-stream 实际上在某处泄漏。快速浏览 repo 表明它只是对 Expat 的一个轻包装,所以我发现这不太可能。
    猜你喜欢
    • 1970-01-01
    • 2023-04-08
    • 1970-01-01
    • 2017-10-16
    • 2012-12-20
    • 1970-01-01
    • 1970-01-01
    • 2017-05-13
    • 1970-01-01
    相关资源
    最近更新 更多