【问题标题】:NodeJS - Seeing data events from readable stream without corresponding pause from writable streamNodeJS - 从可读流中查看数据事件,而没有从可写流中相应暂停
【发布时间】:2015-05-13 18:15:58
【问题描述】:

我们在生产中看到一些流的内存使用率极高。这些文件存储在 S3 中,我们在 S3 对象上打开一个可读流,然后我们将该数据通过管道传输到我们本地文件系统(在我们的 EC2 实例上)上的一个文件。我们的一些客户有非常大的文件。在一个例子中,他们有一个超过 6GB 的文件,处理这个文件的节点进程使用了​​太多的内存,以至于我们几乎耗尽了所有的交换空间,机器减速到爬行。显然,某处存在一些内存泄漏,这正是我试图追踪的。

与此同时,当我们从流中看到某些事件时,我稍微增加了代码以记录。我有下面的代码和来自日志的一些示例输出以及一个小测试文件。令我困惑的是可读流接收到暂停事件,然后继续发出数据和暂停事件没有可写流发出排水事件。我在这里完全错过了什么吗?一旦可读流暂停,它如何在接收到耗尽之前继续发出数据事件?可写流还没有表明它已经准备好了,所以可读流不应该发送任何东西……对吧?

但看看输出。前 3 个事件对我来说很有意义:数据、暂停、耗尽。然后接下来的 3 个很好:数据、数据、暂停。但是 THEN 它发出另一个数据和另一个暂停事件,然后最终作为第 9 个事件排出。我不明白为什么会发生事件 7 和 8,因为直到第 9 个事件才发生排水。然后在第 9 个事件之后再次出现一堆数据/暂停对,没有任何相应的消耗。为什么?我期望的是一些数据事件,然后是暂停,然后 NOTHING 直到发生耗尽事件——此时数据事件可能再次发生。在我看来,一旦发生暂停,在触发排水事件之前根本不应该发生任何数据事件。也许我仍然从根本上误解了 Node 流?

更新:文档没有提到任何关于可读流发出的暂停事件,但他们确实提到了暂停功能可用。大概当可写流返回false时会调用它,我假设暂停函数会发出暂停事件。无论如何,如果调用 pause(),文档似乎与我对世界的看法不谋而合。见https://nodejs.org/docs/v0.10.30/api/stream.html#stream_class_stream_readable

此方法将导致流模式下的流停止发送数据 事件。任何可用的数据都将保留在内部 缓冲区。

这个测试是在我的开发机器上运行的(Ubuntu 14.04 和 Node v0.10.37)。我们在 prod 中的 EC2 实例几乎相同。我认为他们现在运行的是 v0.10.30。

S3Service.prototype.getFile = function(bucket, key, fileName) {
  var deferred = Q.defer(),
    self = this,
    s3 = self.newS3(),
    fstream = fs.createWriteStream(fileName),
    shortname = _.last(fileName.split('/'));

  logger.debug('Get file from S3 at [%s] and write to [%s]', key, fileName);

  // create a readable stream that will retrieve the file from S3
  var request = s3.getObject({
    Bucket: bucket,
    Key: key
  }).createReadStream();

  // if network request errors out then we need to reject
  request.on('error', function(err) {
      logger.error(err, 'Error encountered on S3 network request');
      deferred.reject(err);
    })
    .on('data', function() {
      logger.info('data event from readable stream for [%s]', shortname);
    })
    .on('pause', function() {
      logger.info('pause event from readable stream for [%s]', shortname);
    });

  // resolve when our writable stream closes, or reject if we get some error
  fstream.on('close', function() {
      logger.info('close event from writable stream for [%s] -- done writing file', shortname);
      deferred.resolve();
    })
    .on('error', function(err) {
      logger.error(err, 'Error encountered writing stream to [%s]', fileName);
      deferred.reject(err);
    })
    .on('drain', function() {
      logger.info('drain event from writable stream for [%s]', shortname);
    });

  // pipe the S3 request stream into a writable file stream
  request.pipe(fstream);

  return deferred.promise;
};

[2015-05-13T17:21:00.427Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.427Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.427Z] INFO: worker/7525 on bdmlinux: drain event from writable stream for [FeedItem.csv] [2015-05-13T17:21:00.507Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.514Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.515Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.515Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.515Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.515Z] INFO: worker/7525 on bdmlinux: drain event from writable stream for [FeedItem.csv] [2015-05-13T17:21:00.595Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.596Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.596Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.596Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.597Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.597Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.597Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.597Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.598Z] INFO: worker/7525 on bdmlinux: drain event from writable stream for [FeedItem.csv] [2015-05-13T17:21:00.601Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.602Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.602Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.602Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.603Z] INFO: worker/7525 on bdmlinux: drain event from writable stream for [FeedItem.csv] [2015-05-13T17:21:00.627Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.627Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.627Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.628Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.628Z] INFO: worker/7525 on bdmlinux: drain event from writable stream for [FeedItem.csv] [2015-05-13T17:21:00.688Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.689Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.689Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.689Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.690Z] INFO: worker/7525 on bdmlinux: data event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.690Z] INFO: worker/7525 on bdmlinux: pause event from readable stream for [FeedItem.csv] [2015-05-13T17:21:00.691Z] INFO: worker/7525 on bdmlinux: close event from writable stream for [FeedItem.csv] -- done writing file

【问题讨论】:

    标签: node.js amazon-s3 stream


    【解决方案1】:

    您可能在这里遇到一些类似量子的“观察现象会改变结果”的情况。节点introduced v0.10 中的一种新的流式传输方式。来自docs

    如果您附加一个数据事件侦听器,那么它将把流切换到流动模式,并且一旦数据可用,就会将数据传递给您的处理程序。

    即,附加数据侦听器会将流恢复为经典流模式。这可能就是为什么您的行为与您在其余文档中阅读的内容不一致的原因。为了不打扰地观察事物,您可以尝试删除您的 on('data') 并使用 through 在其间插入您自己的流,如下所示:

    var through = require('through');
    
    var observer = through(function write(data) {
        console.log('Data!');
        this.queue(data);
    }, function end() {
        this.queue(null);
    });
    
    request.pipe(observer).pipe(fstream);
    

    【讨论】:

    • 太棒了,我认为你成功了。我在某处读到附加数据事件会将流切换到“旧模式”,但我没有意识到在 v0.10 之前,暂停事件只是 advisory 而不是保证。这似乎是这里的秘诀。所以我认为你是完全正确的。附加数据事件侦听器会将流切换到“旧模式”以实现向后兼容性。在旧模式下, pause() 方法只是建议性的。
    猜你喜欢
    • 1970-01-01
    • 2021-03-03
    • 1970-01-01
    • 2023-03-05
    • 2021-11-07
    • 2020-10-30
    • 1970-01-01
    • 2016-06-04
    • 1970-01-01
    相关资源
    最近更新 更多