【发布时间】: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