【问题标题】:how to pause and unpause node object stream while processing its output如何在处理其输出时暂停和取消暂停节点对象流
【发布时间】:2015-06-25 15:24:08
【问题描述】:

我目前正在逐行处理文件流,方法是通过发出'line' 事件的转换流运行它。我希望能够在发现当前行符合某些条件时暂停输入文件流,开始处理新流,完成后,继续逐行处理原始流。我已将其浓缩为下面的一个最小示例:

test.coffee:

fs = require 'fs'    
TestTransform = require './test-transform'
inStream = new TestTransform
fs.createReadStream("./test.coffee").pipe(inStream)
inStream.on 'line', (line) ->
  process.stdout.write "-->"
  if line.match /line\.match/g
    process.stdout.write line
    console.error "PAUSE"
    inStream.pause()
    fs.createReadStream("./test.coffee").pipe(process.stdout).on 'end', ->
      console.error "UNPAUSE"
      inStream.resume()
  else
    process.stdout.write line

test-transform.coffee:

Transform = require('stream').Transform

module.exports =
class TestTransform extends Transform
  constructor: ->
    Transform.call @, readableObjectMode: true
    @buffer = ""

  pushLines: ->
    newlineIndex = @buffer.indexOf "\n"
    while newlineIndex isnt -1
      @push @buffer.substr(0, newlineIndex + 1)
      @emit 'line', @buffer.substr(0, newlineIndex + 1)
      @buffer = @buffer.substr(newlineIndex + 1)
      newlineIndex = @buffer.indexOf "\n"

  _transform: (chunk, enc, cb) ->
    @buffer = @buffer + chunk.toString()
    @pushLines()
    cb?()

  _flush: (cb) ->
    @pushLines()
    @buffer += "\n"             # ending newline
    @push @buffer
    @emit 'line', @buffer       # push last line
    @buffer = ""
    cb?()

(不要太担心Transform流,这只是一个例子。)不管怎样,coffee test.coffee 的输出看起来像:

-->fs = require 'fs'
-->
-->TestTransform = require './test-transform'
-->
-->inStream = new TestTransform
-->
-->fs.createReadStream("./test.coffee").pipe(inStream)
-->
-->inStream.on 'line', (line) ->
-->  process.stdout.write "-->"
-->  if line.match /line\.match/g
PAUSE
-->    process.stdout.write line
-->    console.error "PAUSE"
-->    inStream.pause()
-->    fs.createReadStream("./test.coffee").pipe(process.stdout).on 'end', ->
-->      console.error "UNPAUSE"
-->      inStream.unpause()
-->  else
-->    process.stdout.write line
-->
fs = require 'fs'

TestTransform = require './test-transform'

inStream = new TestTransform

fs.createReadStream("./test.coffee").pipe(inStream)

inStream.on 'line', (line) ->
  process.stdout.write "-->"
  if line.match /line\.match/g
    process.stdout.write line
    console.error "PAUSE"
    inStream.pause()
    fs.createReadStream("./test.coffee").pipe(process.stdout).on 'end', ->
      console.error "UNPAUSE"
      inStream.unpause()
  else
    process.stdout.write line

很明显,管道并没有暂停,它只是一直持续到完成(即使PAUSE 正在按预期运行),并且由于"UNPAUSE" 也永远不会被写出,所以'end'回调永远不会触发。将流从转换流切换到暂停/取消暂停到 readStream 似乎也不起作用。我从这种行为中假设节点流在某种程度上不尊重事件回调中的暂停/取消暂停。

可能还有另一种方法可以在不调用暂停/取消暂停的情况下完成此操作;如果有某种方法可以等待流的结束并暂停当前的执行线程,那将有效地完成我想要做的事情。

【问题讨论】:

  • 是否必须完成处理才能再次开始读取流?启动新的处理作业并继续从流中读取还不够吗? Node 擅长异步处理。
  • @Interrobang 是的,我正在尝试将两个输入流通过管道传输到同一个输出流,重要的是在输入第一个流的其余部分之前将第二个流完全读取到输出。我不希望这两个流穿插在输出中。
  • 如果不散布它们就足够了,您可以使用像concat-stream 这样的缓冲区样式流。否则,您需要在流之上进行抽象。一种有趣的方法是使用类似Dust.js 的东西,它可以原生地交错流。
  • 我在想类似的事情。由于我不希望处理长度为千兆字节的流,因此我可以想象将其全部通过管道传输到缓冲区中,然后在其他流完成时对其进行处理。不过,我宁愿不必一次将整个流保存在内存中。我去看看灰尘,我以前没见过。

标签: javascript node.js coffeescript stream


【解决方案1】:

如果我正确理解了这个问题,这里有一个使用 Dust.js 的简单 Node 应用程序可以解决问题。

Dust 是一个模板引擎,但它最好的特性之一是它对 Node Streams 的原生理解。此示例使用 Dust 2.7.0。

我使用node-byline 代替您的转换流,但它做同样的事情——逐行读取流。

var fs = require('fs'),
    byline = require('byline'),
    dust = require('dustjs-linkedin');

var stream = byline(fs.createReadStream('./test.txt', { encoding: 'utf8' }));

var template = dust.loadSource(dust.compile('{#byline}--> {.|s}{~n}{match}{/byline}'));

dust.stream(template, {
  byline: stream,
  match: function(chunk, context) {
    var currentLine = context.current();

    if(currentLine.match(/line\.match/g)) {
      return fs.createReadStream('./test.txt', 'utf8');
    }
    return chunk;
  }
}).pipe(process.stdout);

这是我的程序的输出:

$ node index.js
--> fs = require 'fs'
--> TestTransform = require './test-transform'
--> inStream = new TestTransform
--> fs.createReadStream("./test.coffee").pipe(inStream)
--> inStream.on 'line', (line) ->
-->   process.stdout.write "-->"
-->   if line.match /line\.match/g
fs = require 'fs'
TestTransform = require './test-transform'
inStream = new TestTransform
fs.createReadStream("./test.coffee").pipe(inStream)
inStream.on 'line', (line) ->
  process.stdout.write "-->"
  if line.match /line\.match/g
    process.stdout.write line
    console.error "PAUSE"
    inStream.pause()
    fs.createReadStream("./test.coffee").pipe(process.stdout).on 'end', ->
      console.error "UNPAUSE"
      inStream.resume()
  else
    process.stdout.write line

-->     process.stdout.write line
-->     console.error "PAUSE"
-->     inStream.pause()
-->     fs.createReadStream("./test.coffee").pipe(process.stdout).on 'end', ->
-->       console.error "UNPAUSE"
-->       inStream.resume()
-->   else
-->     process.stdout.write line

如您所见,它正确地交错了输出。如果我可以进一步详细说明 Dust 部分的工作原理,请告诉我。

编辑:这里专门解释了 Dust 模板。

{#byline} {! look for the context variable named `byline` !}
{! okay, it's a stream. For each `data` event, output this stuff once !}
-->
{.|s} {! output the current `data`. Use |s to turn off HTML escaping !}
{~n} {! a newline !}
{match} {! look up the variable called `match` !}
{! okay, it's a function. Run it and insert the result !}
{! if the result is a stream, stream it in. !}
{/byline} {! done looping !}

【讨论】:

  • 这似乎有道理! (在阅读了灰尘语法大声笑之后)我一直在寻找一个没有外部依赖的解决方案,但这似乎很轻量级。在给dust.stream的匹配函数中,为什么写出了if line.match /line\.match/g这一行?似乎灰尘只会返回 fs.createReadStream 而不是块本身,并且该行会丢失。
  • 使用{#match/} 每行调用一次匹配函数。如果当前行 (context.current()) 匹配,则函数流入test.txt 的内容中。如果没有,它只返回当前的chunk,这允许流继续。
  • 这是有道理的。模板字符串中的{.|s} 部分是什么?我假设这是告诉它从“流”属性(以 s 开头)读取,或者如果流不存在,则只是任何属性,但这可能完全不符合要求。
  • {.} 表示“当前上下文”,|s 表示“不要 HTML 转义”。
  • 我更新了答案,对模板进行了更全面的解释。
【解决方案2】:

我实际上也找到了一个单独的答案;不那么漂亮,但也可以。

基本上,pause() 只暂停管道流的输出(在“流动”模式下);因为我在听'line' 事件,它没有流动,所以pause 当然什么也没做。所以第一个解决方案是使用removeListener 而不是pause,这确实有效地停止了流式传输。该文件现在看起来像:

fs = require 'fs'
TestTransform = require './test-transform'
inStream = new TestTransform
fs.createReadStream("./test.coffee").pipe(inStream)
c = (line) ->
  process.stdout.write "-->"
  if line.match /line\.match/g
    process.stdout.write line
    console.error "PAUSE"
    inStream.removeListener 'line', c
    f = fs.createReadStream("./test.coffee")
    f.on 'end', ->
      console.error "UNPAUSE"
      inStream.on 'line', c
    f.pipe(process.stdout)
  else
    process.stdout.write line
inStream.on 'line', c

这会产生几乎有效的输出:

-->fs = require 'fs'
-->TestTransform = require './test-transform'
-->inStream = new TestTransform
-->fs.createReadStream("./test.coffee").pipe(inStream)
-->c = (line) ->
-->  process.stdout.write "-->"
-->  if line.match /line\.match/g
PAUSE
fs = require 'fs'
TestTransform = require './test-transform'
inStream = new TestTransform
fs.createReadStream("./test.coffee").pipe(inStream)
c = (line) ->
  process.stdout.write "-->"
  if line.match /line\.match/g
    process.stdout.write line
    console.error "PAUSE"
    inStream.removeListener 'line', c
    f = fs.createReadStream("./test.coffee")
    f.on 'end', ->
      console.error "UNPAUSE"
      inStream.on 'line', c
    f.pipe(process.stdout)
  else
    process.stdout.write line
inStream.on 'line', c
UNPAUSE

但是,当我移除监听器时,看起来原始的可读流刚刚停止;这具有某种扭曲的意义(我猜节点垃圾会在所有侦听器都被删除后收集其可读流)。因此,我发现的最终工作解决方案依赖于管道。由于我在上面显示的 Transform 流也将其输出逐行推送到任何 'data' 侦听器,pause() 可以在这里有效地用于其原始目标,而不仅仅是杀死流。最终输出:

fs = require 'fs'
TestTransform = require './test-transform'
inStream = new TestTransform
fs.createReadStream("./test.coffee").pipe(inStream)
inStream.on 'data', (chunk) ->
  line = chunk.toString()
  process.stdout.write "-->#{line}"
  if line.match /line\.match/g
    inStream.pause()
    f = fs.createReadStream("./test.coffee")
    f.on 'end', ->
      inStream.resume()
    f.pipe(process.stdout)

带输出:

-->fs = require 'fs'
-->TestTransform = require './test-transform'
-->inStream = new TestTransform
-->fs.createReadStream("./test.coffee").pipe(inStream)
-->inStream.on 'data', (chunk) ->
-->  line = chunk.toString()
-->  process.stdout.write "-->#{line}"
-->  if line.match /line\.match/g
fs = require 'fs'
TestTransform = require './test-transform'
inStream = new TestTransform
fs.createReadStream("./test.coffee").pipe(inStream)
inStream.on 'data', (chunk) ->
  line = chunk.toString()
  process.stdout.write "-->#{line}"
  if line.match /line\.match/g
    inStream.pause()
    f = fs.createReadStream("./test.coffee")
    f.on 'end', ->
      inStream.resume()
    f.pipe(process.stdout)
-->    inStream.pause()
-->    f = fs.createReadStream("./test.coffee")
-->    f.on 'end', ->
-->      inStream.resume()
-->    f.pipe(process.stdout)
-->

这是预期的结果。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-06-17
    • 1970-01-01
    • 1970-01-01
    • 2012-01-13
    • 2016-12-17
    • 1970-01-01
    • 2012-03-25
    • 1970-01-01
    相关资源
    最近更新 更多