【问题标题】:Why flatMap has no output when one stream has error?为什么当一个流有错误时 flatMap 没有输出?
【发布时间】:2017-09-18 14:10:22
【问题描述】:

我尝试用highland.js编写一个程序来下载几个文件,解压缩并解析成对象,然后通过flatMap将对象流合并为一个流并打印出来。

function download(url) {
    return _(request(url))
        .through(zlib.createGunzip())
        .errors((err) => console.log('Error in gunzip', err))
        .through(toObjParser)
        .errors((err) => console.log('Error in OsmToObj', err));
}  

const urlList = ['url_1', 'url_2', 'url_3'];

_(urlList)
    .flatMap(download)
    .each(console.log);

当所有 URL 都有效时,它可以正常工作。如果 URL 无效且没有下载文件,则 gunzip 报告错误。我怀疑发生错误时流会关闭。我希望flatMap 将继续使用其他流,但是程序不会下载其他文件并且没有打印任何内容。

处理流错误的正确方法是什么?如何使flatMap在一个流出现错误后不停止?

在命令式编程中,我可以添加调试日志来跟踪错误发生的位置。如何调试流式代码?

PS。 toObjParser 是一个节点转换流。它采用可读的 OSM XML 流并输出与 Overpass OSM JSON 兼容的对象流。见https://www.npmjs.com/package/osm2obj

2017-12-19 更新:

我尝试按照@amsross 的建议在errors 中调用push。为了验证push 是否真的有效,我推送了一个 XML 文档,它被以下解析器解析,我从输出中看到了它。但是,stream 仍然停止并且 url_3 没有下载。

function download(url) {
    console.log('download', url);
    return _(request(url))
        .through(zlib.createGunzip())
        .errors((err, push) => {
            console.log('Error in gunzip', err);
            push(null, Buffer.from(`<?xml version='1.0' encoding='UTF-8'?>
<osmChange version="0.6">
<delete>
<node id="1" version="2" timestamp="2008-10-15T10:06:55Z" uid="5553" user="foo" changeset="1" lat="30.2719406" lon="120.1663723"/>
</delete>
</osmChange>`));
        })
        .through(new OsmToObj())
        .errors((err) => console.log('Error in OsmToObj', err));
}

const urlList = ['url_1_correct', 'url_2_wrong', 'url_3_correct'];

_(urlList)
    .flatMap(download)
    .each(console.log);

【问题讨论】:

    标签: javascript node.js stream highland.js


    【解决方案1】:

    2017 年 12 月 19 日更新: 好的,所以我不能给你一个很好的为什么,但我可以告诉你,从使用 sequence 中的 download 产生的流切换到 merge'ing 它们在一起可能会给你你想要的结果。不幸的是(或不是?),您将无法再按任何规定的顺序获得结果。

    const request = require('request')
    const zlib = require('zlib')
    const h = require('highland')
    
    // just so you can see there isn't some sort of race
    const rnd = (min, max) => Math.floor((Math.random() * (max - min))) + min
    const delay = ms => x => h(push => setTimeout(() => {
      push(null, x)
      push(null, h.nil)
    }, ms))
    
    const download = url => h(request(url))
      .flatMap(delay(rnd(0, 2000)))
      .through(zlib.createGunzip())
    
    h(['urlh1hcorrect', 'urlh2hwrong', 'urlh3hcorrect'])
      .map(download).merge()
      // vs .flatMap(download) or .map(download).sequence()
      .errors(err => h.log(err))
      .each(h.log)
    

    2017 年 12 月 3 日更新: 当在流上遇到错误时,它会结束该流。为避免这种情况,您需要处理错误。你目前使用errors报错,但不处理。您可以执行以下操作来移动到流中的下一个值:

    .errors((err, push) => {
      console.log(err)
      push(null) // push no error forward
    })
    

    原文: 不知道toObjParser的输入输出类型是很难回答的。

    因为through 将值流传递给提供的函数并期望返回值流,所以您的问题可能在于toObjParser 具有Stream -&gt; ObjectStream -&gt; Stream Object 之类的签名,其中错误是发生在内部流上,在消耗之前不会发出任何错误。

    .each(console.log) 的输出是什么?如果它正在记录流,那很可能是您的问题。

    【讨论】:

    • toObjParser 是一个节点转换流。它接受一个可读的 OSM XML 流并输出一个对象流。 npmjs.com/package/osm2obj
    • 噢噢噢噢!我误解了这个问题。我会更新我的答案,以反映我认为你可以做些什么来推动直播。
    • 我试过push,但还是不行。查看我有问题的更新。
    • 我添加了一个额外的语句和代码,与并行而不是按顺序使用值相关。这些似乎解决了您遇到的问题,但是对于为什么会出现这种情况的最佳解释是“只是因为”。对不起。
    猜你喜欢
    • 1970-01-01
    • 2012-05-23
    • 1970-01-01
    • 2011-11-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多