【问题标题】:"No fiber running" while trying to parse a csv line by line尝试逐行解析csv时“没有光纤运行”
【发布时间】:2016-05-06 01:23:35
【问题描述】:

我试图了解如何在 fast-csv 中使用 Fiber 来制作逐行阅读器(单用户命令行脚本),该阅读器在每一行暂停读取/处理,直到该行完成各种的异步调用。 (不滚动我自己的 csv 代码,我想使用一些已经弄清楚关于 csv 格式的问题的东西)

如果我这样做

var csv = require("fast-csv");

var CSV_STRING = 'a,b\n' +
'a1,b1\n' +
'a2,b2\n';

csv
.fromString(CSV_STRING, {headers: false})
.on("record", function (data) {
    console.log("line="+JSON.stringify(data));
    setTimeout(function(){
        console.log("timeout");
    },2000);
})
.on("end", function () {
    console.log("done parsing CSV records");
});
console.log("done initializing csv parse");

我得到了我的期望:

done initializing csv parse
line=["a","b"]
line=["a1","b1"]
line=["a2","b2"]
done parsing CSV records
timeout
timeout
timeout

如果我尝试在每条记录后使用 Fiber 来产生

Fiber(
    function () {
        var fiber = Fiber.current;

        csv
            .fromString(CSV_STRING, {headers: false})
            .on("record", function (data) {
                console.log("line="+JSON.stringify(data));
                setTimeout(function(){
                    console.log("timeout");
                    fiber.run();
                },2000);
                Fiber.yield();
            })
            .on("end", function () {
                console.log("done parsing CSV records");
            });
        console.log("done initializing csv parse");
    }).run();

我明白了

done initializing csv parse
line=["a","b"]
events.js:141
      throw er; // Unhandled 'error' event
      ^

Error: yield() called with no fiber running

我想我明白发生了什么,Fiber().run() 中的代码完成了,所以它在调用 yield 之前离开了 Fiber,所以当它达到 yield 时不再有 Fiber。 (因此巧妙的错误消息“没有光纤运行”)

对我来说,在完成解析之前保持光纤运行的合适方法是什么?

似乎是一个如此简单的问题,但我没有看到明显的答案?起初我想在它离开 Future().run() 之前设置一个 yield,但这不起作用,因为第一个 fiber.run() 会让它再次离开 Fiber。

我想要的流程是这样的:

done initializing csv parse
line=["a","b"]
timeout
line=["a1","b1"]
timeout
line=["a2","b2"]
timeout
done parsing CSV records

但如果不重新设计 fast-csv 的内部,这可能是不可能的,因为它控制了每个记录的事件何时触发。我目前的想法是,必须让每个事件在 fast-csv 中被触发,并让处理 csv.on("record") 中的事件的用户将控制权交还给快速解析 csv 的循环-csv。

【问题讨论】:

    标签: javascript node.js csv node-fibers


    【解决方案1】:

    流是可暂停/可恢复的:

    var csv = require("fast-csv");
    
    var CSV_STRING = 'a,b\n' +
        'a1,b1\n' +
        'a2,b2\n';
    
    var stream = csv.fromString(CSV_STRING, { headers: false })
        .on("data", function (data) {
            // pause the stream
            stream.pause();
            console.log("line: " + JSON.stringify(data));
            setTimeout(function () {
                // all async stuff are done, resume the stream
                stream.resume();
                console.log("timeout");
            }, 2000);
        }).on("end", function () {
            console.log("done parsing CSV records");
        });
    

    控制台输出几乎正是您想要的:

    /*
    line: ["a","b"]
    timeout
    line: ["a1","b1"]
    timeout
    line: ["a2","b2"]
    done parsing CSV records
    timeout
    */
    

    我能问一下为什么你绝对需要同步读取你的 csv 吗?

    【讨论】:

    • 我尝试了您的建议,但遇到了暂停不能保证暂停的问题。事实上,在我的情况下,在它暂停之前触发了 20 多个事件。然后在文档中注意到“请注意,这不会立即暂停事件流。调用暂停后可能会发出几个事件,包括行。”
    • 至于为什么:也许这是我的老派,但是当我创建单个用户脚本时,在很多情况下我不希望并行执行。从调试的角度来看,我发现它是批处理类型作业的噩梦。在某些情况下,如果有问题,我需要在一张唱片上完成大量工作,然后再继续。事实上,我可能有 100,000 个事件将所有冲击资源排入队列。我只想一次处理一个请求。生成器解决方案似乎正是我想要的。它似乎可以很好地处理数据库查询、网络调用和其他异步内容。感谢您的意见!
    • 最后一点。我根本不是反异步的,只是在某些情况下我真的更喜欢同步风格的执行。
    • @sday 我明白了,但是 nodejs 是单线程的,没有并行执行。在像流这样的异步流上强制执行同步对我来说似乎适得其反。它还会使执行时间更长。事件队列不会爆炸,如有必要,一切都会暂停。我尝试在具有 1GB RAM 的虚拟机中使用一个 2GB 的大 csv 文件,内存很快被完全使用,处理速度变慢,但节点没有崩溃。
    • 也许我误解了正在发生的事情。如果我从 csv 库接收到 100 个事件,提示我进行 100 个 DB 写入,每个写入需要 100 毫秒,那么所有 100 个事件和 100 个 DB 写入调用都会在 50 毫秒内返回,现在我们等待 DB 回调表明它们已完成,对吗? Node 可能已经以单线程方式执行了它的调用,但它们本质上调用了在它们自己的线程中工作的较低级别的例程,该线程与执行的节点线程并行工作。这不正确吗?
    【解决方案2】:

    节点:v5.4.0

    嗯,这是获得这种行为的一种方法。我使用 es6 生成器逐行读取原始文件,然后使用 fast-csv 库上的生成器从逐行读取中解析原始字符串,这导致非异步执行流程和类似于旧单的输出用户命令行脚本。

    'use strict';
    var csv = require("fast-csv");
    var sfs = require('./sfs');
    
    function parse(line) {
        csv
            .fromString(line, {headers: false})
            .on("record", function (data) {
                it.next(data);
            });
    }
    
    function *main() {
        // Make sure to initialize with a max buffer big enough to span any possible line length.  Otherwise undefined
        var fs = new sfs(it, 4096);
        var result=yield fs.open("data.csv");
    
        var line;
    
        while((line=yield fs.readLine()) != null) {
            console.log("line="+line);
    
            var csvData=yield parse(line);
            console.log("value1="+csvData[0]+" value2="+csvData[1]);
        }
    
        console.log("DONE");
    }
    
    var it = main();
    it.next(); // get it all started
    

    连同一个 quacky (quick and hacky) 类来包装我需要的 fs 东西。我确信有一种更好的方法来做我所做的事情,但它可以满足我的需要。

    sfs.js

    'use strict';
    var fs=require('fs')
    
    class sfs {
        constructor(it, maxbufsize) {
            this.MAX_BUF=maxbufsize;
            this.it=it;
            this.fd=null;
            this.lineBuf="";
            this.buffer=new Buffer(this.MAX_BUF);
            this.buflen=0;
        }
    
        open(file) {
            var parent=this;
            fs.open(file,'r',function(err,fd){
                parent.fd=fd;
                var parent2=parent;
                fs.fstat(fd,function(err, stats){
                    parent2.stats=stats;
                    parent2.it.next(stats);
                })
            })
    
        }
    
        readLine(){
            var parent = this;
            var i=0
            var s=this.stats.size
            var line="";
            var index=this.MAX_BUF-this.buflen;
    
            // read data into buffer, buffer may already have data from previous read that was shifted left over extracted line
            fs.read(this.fd,this.buffer,this.MAX_BUF-index,index,null,function(err,len,buf){
                var expectedReadLen=parent.MAX_BUF-parent.buflen;
                if(len < expectedReadLen) {  // If we didn't read enough to backfill buffer, lets make sure the string is terminated
                    // as it shifts left so we don't try interpret older records to the right
                    parent.buffer.fill(' ',parent.buflen+len,parent.MAX_BUF);
                }
                parent.buflen+=len; // whatever was in buffer has more now
    
                index=parent.buffer.indexOf('\n');
    
                if(index > -1) {
                    line=parent.buffer.toString('utf8',0,index);
                    buf.copy(parent.buffer,0,index+1,parent.buflen); // shift unused data left
                    parent.buflen-=(index+1); // buffer left over after removing /n terminated line
                    if(len<expectedReadLen) {  // If we didn't read enough to backfill buffer, lets make sure we erase old data
                        parent.buffer.fill(' ',parent.buflen,parent.MAX_BUF);
                    }
                } else {
                    if(parent.buflen > 0) {
                        line=parent.buffer.toString('utf8',0,parent.buflen);
                        parent.buflen=0;
                    } else {
                        line=null;
                    }
                }
                parent.it.next(line);
            });
        }
    
        close() {
            fs.close(this.fd);
        }
    }
    
    module.exports=sfs;
    

    【讨论】:

      猜你喜欢
      • 2016-12-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多