【问题标题】:Node.js Streams Readable to TransformNode.js 流可读可转换
【发布时间】:2015-10-27 02:30:54
【问题描述】:

我一直在尝试使用可读和转换流来处理一个非常大的文件。我似乎遇到的问题是,如果我不把可写流放在最后,程序似乎在返回结果之前就终止了。

例如:rstream.pipe(split()).pipe(tstream)

我的tstream 有一个发射器,当计数器达到阈值时会发射。当该阈值设置为较低的数字时,我会得到一个结果,但是当它很高时,它不会返回任何内容。如果我将它传递给文件编写器,它总是返回一个结果。我错过了什么明显的东西吗?

代码:

// Dependencies
var fs = require('fs');
var rstream = fs.createReadStream('file');
var wstream = fs.createWriteStream('output');
var split = require('split'); // used for separating stream by new line
var QTransformStream = require('./transform');

var qtransformstream = new QTransformStream();
qtransformstream.on('completed', function(result) {
    console.log('Result: ' + result);
});
exports.getQ = function getQ(filename, callback) {

    // THIS WORKS if i have a low counter for qtransformstream, 
    // but when it's high, I do not get a result
    //   rstream.pipe(split()).pipe(qtransformstream);

    // this always works
    rstream.pipe(split()).pipe(qtransformstream).pipe(wstream);

};

这是Qtransformstream的代码

// Dependencies
var Transform = require('stream').Transform,
    util = require('util');
// Constructor, takes in the Quser as an input
var TransformStream = function(Quser) {
    // Create this as a Transform Stream
    Transform.call(this, {
        objectMode: true
    });
    // Default the Qbase to 32 as an assumption
    this.Qbase = 32;
    if (Quser) {
        this.Quser = Quser;
    } else {
        this.Quser = 20;
    }
    this.Qpass = this.Quser + this.Qbase;
    this.Counter = 0;
    // Variables used as intermediates
    this.Qmin = 120;
    this.Qmax = 0;
};
// Extend the transform object
util.inherits(TransformStream, Transform);
// The Transformation to get the Qbase and Qpass
TransformStream.prototype._transform = function(chunk, encoding, callback) {
    var Qmin = this.Qmin;
    var Qmax = this.Qmax;
    var Qbase = this.Qbase;
    var Quser = this.Quser;
    this.Counter++;
    // Stop the stream after 100 reads and emit the data
    if (this.Counter === 100) {
        this.emit('completed', this.Qbase, this.Quser);
    }
    // do some calcs on this.Qbase

    this.push('something not important');
    callback();
};
// export the object
module.exports = TransformStream;

【问题讨论】:

  • 你能贴出QTransformStream实现的代码吗?
  • 输入文件中有多少行以及在这种情况下的最大计数器值是多少。如果计数器值大于行号,则不会发出 completed 事件。您还需要推送null 来结束流。不确定something not important 中有什么,但在某些时候应该有一个null
  • def的行数比计数器少,大约7000行。当我将它通过管道传输到写入流时,它确实有效。转换流是否需要 push(null) 才能工作?
  • 你是对的,它不是。可能是别的东西。

标签: node.js stream node.js-stream


【解决方案1】:

编辑:

另外,我不知道你的计数器有多高,但如果你填满缓冲区,它将停止将数据传递到转换流,在这种情况下,completed 永远不会真正命中,因为你永远不会达到计数器限制。尝试更改您的highwatermark

编辑 2:更好的解释

众所周知,transform stream 是双工流,这基本上意味着它可以接受来自源的数据,也可以将数据发送到目的地。这通常分别称为读和写。 transform stream 继承自 Node.js 实现的 read streamwrite stream。不过有一个警告,transform stream 不必实现 _read 或 _write 函数。从这个意义上说,您可以将其视为鲜为人知的 passthrough stream

如果您考虑transform stream 实现write stream 的事实,您还必须考虑写入流始终具有转储其内容的目的地这一事实。 您遇到的问题是,当您创建transform stream 时,您无法指定发送内容的位置。 将数据完全通过转换流传递的唯一方法是将其通过管道传输到写入流,否则,实质上您的流会被备份并且无法接受更多数据,因为数据无处可去.

这就是为什么当您通过管道传输到写入流时它总是有效的原因。写入流通过将数据发送到目的地来减轻数据备份,因此您的所有数据都将通过管道传输并发出完成事件。

当样本量较小时,您的代码在没有写入流的情况下工作的原因是您没有填满您的流,因此转换流可以接受足够的数据以允许达到完整的事件/阈值。随着阈值的增加,您的流可以在不将其发送到另一个地方(写入流)的情况下接受的数据量保持不变。这会导致您的流得到备份,并且它不能再接受数据,这意味着完成的事件将永远不会被发出。

我敢说,如果您为转换流增加highwatermark,您将能够增加您的阈值并且仍然可以使代码正常工作。这种方法虽然是不正确的。将您的流传输到一个写入流,该写入流会将数据发送到 dev/null 创建该写入流的方式是:

var writer = fs.createWriteStream('/dev/null');

buffering 上的 Node.js 文档中的部分解释了您遇到的错误。

【讨论】:

  • node 中的流并不像看起来那么简单。我很想看到这些微妙之处的详细解释。
  • 我试图做一个更好的解释,如果有部分不清楚,请告诉我。
【解决方案2】:

您不会中断 _transform 并且进程会走得很远。试试:

this.emit('completed', ...);
this.end();

这就是为什么'程序似乎在返回结果之前终止'

并且不要输出任何无用的数据:

var wstream = fs.createWriteStream('/dev/null');

祝你好运)

【讨论】:

    【解决方案3】:

    我建议使用Writable 而不是转换流。 然后将_transform 重命名为_write,如果您通过管道传输到该流,您的代码将使用该流。正如@Bradgnar 已经指出的那样,转换流需要一个消费者,否则它将stop the readable 流向其缓冲区推送更多数据。

    【讨论】:

      猜你喜欢
      • 2016-06-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-03-13
      • 2023-03-05
      • 1970-01-01
      • 2020-12-06
      • 2014-10-29
      相关资源
      最近更新 更多