【问题标题】:How to use fibers with streams如何在流中使用纤维
【发布时间】:2014-10-04 14:08:42
【问题描述】:

我正在尝试将纤维与流一起使用:

var Fiber = require('fibers');
var Future = require('fibers/future');
var fs = require('fs');

function sleepForMs(ms) {
  var fiber = Fiber.current;
  setTimeout(function() {
    fiber.run();
  }, ms);
  Fiber.yield();
}

function catchError(f, onError) {
  return function () {
    var args = arguments;
    var run = function () {
      try {
        var ret = f.apply(null, args);
      }
      catch (e) {
        onError(e);
      }
      return ret;
    };
    if (Fiber.current) {
      return run();
    }
    else {
      return Fiber(run).run();
    }
  }
}

function processFile(callback) {
  var count, finished, onData, onException, onIgnoredEntry;
  count = 0;
  finished = false;
  onException = function (error) {
    if (finished) {
      console.error("Exception thrown after already finished:", error.stack || error);
    }
    if (finished) {
      return;
    }
    finished = true;
    return callback(error);
  };
  onData = function(data) {
    console.log("onData");
    if (finished) {
      return;
    }
    console.log("before sleep");
    sleepForMs(500);
    console.log("after sleep");
    throw new Error("test");
  };
  return fs.createReadStream('test.js').on('data', catchError(onData, onException)).on('end', function() {
    console.log("end");
    if (finished) {
      return;
    }
    finished = true;
    return callback(null, count);
  }).on('error', function(error) {
    console.log("error", error);
    if (finished) {
      return;
    }
    finished = true;
    return callback(error);
  });
};

Fiber(function () {
  console.log("Calling processFile");
  Future.wrap(processFile)().wait();
  console.log("processFile returned");
}).run();
console.log("back in main");

但它并没有真正起作用。数据回调在回调内部的光纤完成之前完成。所以上面的代码输出:

Calling processFile
back in main
onData
before sleep
end
processFile returned
after sleep
Exception thrown after already finished: Error: test

实际上它应该更像是:

Calling processFile
back in main
onData
before sleep
after sleep
end
processFile returned
Error: test

【问题讨论】:

  • 只需删除sleepForMs(500)
  • Sleep 只是任何其他可能产生的启用光纤的功能的一个示例。
  • 问题是onExceptiononData在单独的Fiber中执行(见return Fiber(run).run()块在if (Fiber.current)条件catchError函数),因为所有.on回调出现在当前Fiber之外。
  • 是的,如何让它们在一根光纤内工作?
  • 不使用光纤? v0.11 的节点有自己的yield 函数(实际上JS有):检查node --harmony_generators...

标签: javascript node.js node-fibers


【解决方案1】:

看起来没有人知道如何做你所要求的。

在这种情况下,您可以以某种传统的异步方式处理您的流,将您的屈服函数应用于结果。

这里有一些如何做到这一点的例子。


raw-body收集所有流数据

一种解决方案是在处理任何流数据之前收集所有流数据。使用raw-body module 可以轻松完成:

var rawBody = require('raw-body');

function processData(data) {
  console.log("before sleep");
  sleepForMs(500);
  console.log("after sleep");
}

function processFile(callback) {
  var stream = fs.createReadStream('fiber.js');
  rawBody(stream, function(err, data) {
    if (err) return callback(err);
    Fiber(processData).run(data); // process your data
    callback();
  });
}

使用此示例,您将:

  1. 等待所有块到达
  2. 开始处理您在Fiber 中的流数据
  3. processData返回主线程
  4. 流数据将在未来某个时间点进行处理

如果需要,您可以添加try ... catch 或任何其他异常处理,以防止processData 破坏您的应用。


使用智能作业队列处理系列中的所有块

但如果你真的想在所有数据块到达的那一刻处理它们,你可以使用一些智能控制流模块。这是使用async module 中的queue feature 的示例:

function processChunk(data, next) {
  return function() {
    console.log("before sleep");
    sleepForMs(500);
    console.log("after sleep");
    next();
  }
}

function processFile(callback) {
  var q = async.queue(function(data, next) {
    Fiber(processChunk(data, next)).run();
  }, 1);
  fs.createReadStream('fiber.js').on('data', function(data) {
    q.push(data);
  }).on('error', function(err) {
    callback(err);
  }).on('end', function() {
    callback(); // not waiting to queue to drain
  })
}

使用此示例,您将:

  1. 开始收听streampush将每个新块发送到处理队列
  2. stream 关闭时从processData 返回主线程,不等待处理数据块
  3. 所有数据块都将在某个时间点以严格的顺序进行处理

我知道这不是你要求的,但我希望它会帮助你。

【讨论】:

    【解决方案2】:

    减少睡眠时间并为其他块设置一些优先级或计时器。 以便在一定的时间限制后根据优先级显示。 这就是您如何以您想要的方式获得输出。

    【讨论】:

    • 睡眠时间只是一个可以产生的函数的例子。
    【解决方案3】:

    这是一个使用 wait.for 的实现(Fibers 的包装器)https://github.com/luciotato/waitfor

    在这个实现中,为每个数据块启动一个纤程,因此并行启动“n”个任务。在所有纤程完成之前,ProcessFile 不会“返回”。

    这是一个演示,展示了如何使用 Fibers 和 wait.for 执行此操作,但当然,在生产中使用它之前,您应该将模块级变量和所有函数封装在一个类中。

    var wait = require('wait.for');
    var fs = require('fs');
    
    var tasksLaunched=0;
    var finalCallback;
    var callbackDone=false;
    var dataArr=[]
    
    function sleepForMs(ms,sleepCallback) {
      setTimeout(function() {
        return sleepCallback();
      }, ms);
    }
    
    function resultReady(err,data){
    
        if (err){
          callbackDone = true;
          return finalCallback(err);
        }
    
        dataArr.push(data);
        if (dataArr.length>=tasksLaunched && !callbackDone) {
          callbackDone = true;
          return finalCallback(null,dataArr);
        }
    }
    
    function processChunk(data,callback) {
        var ms=Math.floor(Math.random()*1000);
        console.log('waiting',ms);
        wait.for(sleepForMs,ms);
        console.log(data.length,"chars");
        return callback(null,data.length);
    }
    
    function processFile(filename,callback) {
      var count, onData, onException, onIgnoredEntry;
      count = 0;
      finalCallback = callback;
    
      onException = function (error) {
        if (!callbackDone){
          callbackDone = true;
          return callback(error);
        }
      };
    
      onData = function(data) {
        console.log("onData");
        tasksLaunched++;
        wait.launchFiber(processChunk,data,resultReady);
      };
    
      fs.createReadStream(filename)
        .on('data', onData)
        .on('end', function() {
            console.log("end");
        })
        .on('error', function(error) {
            console.log("error", error);
            if (!callbackDone) {
                callbackDone = true;
                return callback(error);
              }
        });
    };
    
    function mainFiber() {
      console.log("Calling processFile");
      var data = wait.for(processFile,'/bin/bash');
      console.log(data.length,"results");
      console.log("processFile returned");
    };
    
    //MAIN
    wait.launchFiber(mainFiber);
    console.log("back in main");
    

    【讨论】:

    • 我认为问题的重点是使用一些 yieldind 函数顺序处理到达的数据块。
    • 鉴于“期望的输出”,他似乎希望在 processFile 返回之前完成这两个任务
    • 是的。 sleep 这里只是一个例子。想象一下可以产生的任何其他支持光纤的功能。我希望能够在每个数据块的回调中调用它,而不必担心其他块同时在做什么。这样事情的行为就会像您期望的阻塞调用一样。
    • 我已经改变了答案,为每个块启动一个光纤
    • 如果sleepForMs 是一个产生光纤的函数而不是一个基于回调的函数,它会起作用吗?
    猜你喜欢
    • 1970-01-01
    • 2012-12-04
    • 2013-09-13
    • 1970-01-01
    • 1970-01-01
    • 2021-08-09
    • 1970-01-01
    • 2016-07-09
    • 2021-07-10
    相关资源
    最近更新 更多