【问题标题】:Node.js: Streaming from MongoDB to fileNode.js:从 MongoDB 流式传输到文件
【发布时间】:2015-12-31 14:15:47
【问题描述】:

我对 javascript/node.js 相当陌生,并且正在尝试让一个非常基本的场景工作:连接到 MongoDB,将 JSON 响应转换为 CSV,将其写入文件。我已经尝试如下

fs = require('fs');
var MongoClient = require('mongodb').MongoClient;
var Db = require('mongodb').Db;
var Server = require('mongodb').Server;
var Json2csvStream = require('json2csv-stream');
var Stream = require('stream');
var JSONStream = require('JSONStream');
var es = require('event-stream');
var csv = require('csv');

var fields = ['execAmendTime', 'execTime', 'execType', 'lastMkt', 'manualExecFlag', 'orderId', 'riskTrade', 'rootOrdId', 'salesCommissionRate', 'salesCommissionType', 'theoPov20Px',
'theoPov20BL', 'tradeFlags', 'tradeNotes', 'transactTime', 'version', 'book.bookName', 'businessUnit', 'commissionRate', 'commissionSource', 'commissionType', 'counterBook.bookName', 
'counterParty.name', 'createTime', 'currency', 'direction', 'execQuantity', 'fxRate','orderQuantity', 'positionTrader.name', 'price', 'primaryTrader.name','rootSystem', 'source',
'sectorGicsLevel1', 'salesTrader.name', 'tradedPrice', 'isCRB', 'clientCategory', 'tradeId', 'tradeDate', 'instrument.instrumentRic', 'notionalUSD','commissionUSD', 'region'];

// Connect to the db
MongoClient.connect("mongodb://*****", function (err, db) {
    if (err) { return console.dir(err); }

    if (!err) {
        console.log("We are connected");
    }

    db.open(function (err, db) {
        if (err) { return console.dir(err); }
        var newDb = db.db("test_db");

        var collection = newDb.collection('test', function (err, collection) {
            if (err) { return console.dir(err); }
            var parser = new Json2csvStream();
            var writer = fs.createWriteStream('out.csv');
            var stream = collection.find({ tradeDate: new Date('2015-12-29T00:00:00.000Z') }).stream();

            stream.pipe(parser).pipe(writer);

            stream.on("data", function (item) {
                console.log(item);
            });

            stream.on('end', function () {
                console.log("ended");
            });

            stream.on("end", function () {
                newDb.close();
                db.close();
            });
        });
    });
});

我收到如下错误。

我尝试使用 JSON.stringify 等添加转换,但我的尝试都没有奏效。看来我需要等到 Mongo 的查询流完成后再开始将其输入 json2csv 转换器?

有什么想法吗?我在这里做错了什么吗?

非常感谢!

输出:

We are connected
D:\WebTrial\MongoProject\node_modules\mongodb\lib\utils.js:98
    process.nextTick(function() { throw err; });
                                ^

TypeError: Invalid non-string/buffer chunk
    at validChunk (_stream_writable.js:178:14)
    at Writable.write (_stream_writable.js:205:12)
    at ondata (_stream_readable.js:525:20)
    at emitOne (events.js:82:20)
    at emit (events.js:169:7)
    at readableAddChunk (_stream_readable.js:146:16)
    at Readable.push (_stream_readable.js:110:10)
    at D:\WebTrial\MongoProject\node_modules\mongodb\lib\cursor.js:1102:10
    at handleCallback (D:\WebTrial\MongoProject\node_modules\mongodb\lib\utils.j
s:96:12)
    at D:\WebTrial\MongoProject\node_modules\mongodb\lib\cursor.js:673:5

【问题讨论】:

  • 你能指出堆栈跟踪的最后三行所指的源代码中的哪几行吗?

标签: javascript node.js mongodb csv streaming


【解决方案1】:

这是因为 find().stream() 流式传输对象,而 Json2csvStream 需要字符串。 event-stream 可以帮助您对对象进行字符串化。我还简化了你的代码,有不必要的东西:

var fs = require('fs');
var MongoClient = require('mongodb').MongoClient;
var es = require('event-stream');
var Json2csvStream = require('json2csv-stream');

// var Db = require('mongodb').Db;
// var Server = require('mongodb').Server;
// var Stream = require('stream');
// var JSONStream = require('JSONStream');
// var csv = require('csv');

var fields = ['execAmendTime', 'execTime', 'commissionUSD', 'region'];

// Connect to the db
// you can put the db name in the url
MongoClient.connect("mongodb://localhost:27017/test_db", function (err, db) {
    if (err) {
        return console.dir(err);
    } else {
        console.log("We are connected");
    }

    // without strict: true, err is always null
    // in strict mode, there is an err if the collection doesn't exist
    db.collection('test', { strict: true }, function (err, collection) {
        if (err) {
            return console.dir(err);
        }

        var json2csv = new Json2csvStream();
        var writer = fs.createWriteStream('out.csv');

        var mongoStream = collection.find(
            { tradeDate: new Date('2015-12-29T00:00:00.000Z') }
        ).stream();

        var stream = mongoStream
            .pipe(es.map(function (doc, next) {
                doc = JSON.stringify(doc);
                // console.log(doc);
                next(null, doc);
            })).pipe(json2csv).pipe(writer).on('close', function () {
                console.log('done...');
                db.close();
            });
    });
});

【讨论】:

  • 非常感谢您的回复!我会在星期一试试看!
  • 嗨,我今天试过了。现在出现不同的错误消息(见下文)。
  • D:\WebTrial\MongoProject\node_modules\mongodb\lib\utils.js:98 process.nextTick(function() { throw err; }); ^ SyntaxError: MyStream.writeHeader (D:\WebTrial\MongoProject\node_modules\json2csv-stre am\index.js:115:17) 在 MyStream._transform (D:\WebTrial) 的 Object.parse (native) 的输入意外结束\MongoProject\node_modules\json2csv-strea m\index.js:95:55) 在 Transform._read (_stream_transform.js:167:10) ...
  • 我可能需要添加一些东西来处理空值或空值吗?
  • @JenniferTenzer 我尝试了一个复杂的文档,我得到了同样的错误,你的文档有子文档吗? Json2csvStream 进行非常基本的检查以查找 JSON(它只查找 { 和我会尝试找到另一个模块。
猜你喜欢
  • 2012-06-18
  • 2013-12-02
  • 2014-08-07
  • 2016-09-21
  • 2013-12-24
  • 1970-01-01
  • 2014-12-06
  • 2017-05-09
  • 2017-07-09
相关资源
最近更新 更多