【问题标题】:Rxjs: updating values in observable stream with data from another observable, returning a single observable streamRxjs:使用来自另一个可观察对象的数据更新可观察流中的值,返回单个可观察流
【发布时间】:2016-06-03 09:05:36
【问题描述】:

背景

我正在尝试从拉取请求的Stash Rest Api 构造一个可观察的值流。不幸的是,PR 是否存在合并冲突的信息在合并列表的不同端点上可用。

打开的拉取请求列表可见于http://my.stash.com/rest/api/1.0/projects/myproject/repos/myrepo/pull-requests

对于每个 PR,关于合并冲突的数据在 http://my.stash.com/rest/api/1.0/projects/myproject/repos/myrepo/pull-requests/[PR-ID]/merge 可见

使用atlas-stash 包,我可以创建和订阅可观察的拉取请求流(每秒更新一次):

let pullRequestsObs = Rx.Observable.create(function(o) {
    stash.pullRequests(project, repo)
        .on('error', function(error) {o.onError(error)})
        .on('allPages', function(data) {
            o.onNext(data);
            o.onCompleted();
        });
    });

let pullRequestStream = pullRequestsObs
    .take(1)
    .merge(
        Rx.Observable
            .interval(1000)
            .flatMapLatest(pullRequestsObs)
    );

pullRequestsStream.subscribe(
    (data) => {
        console.log(data)
        // do something with data
    },
    (error) => log.error(error),
    () => log.info('done')
);

这可以按我的意愿和期望工作。最后,pullRequestsStream 是一个 observable,其值是 JSON 对象的列表。

我的目标

我希望更新 pullRequestsStream 值,以便列表的每个元素都包含来自 [PR-ID]/merge api 的信息。

我认为这可以在pullRequestsStream 上使用map 来实现,但我没有成功。

let pullRequestWithMergeStream = pullRequestStream.map(function(prlist) {
    _.map(prlist, function(pr) {
        let mergeObs = Rx.Observable.create(function(o) {
            stash.pullRequestMerge(project, repo, pr['id'])
                .on('error', function(error) {o.onError(error)})
                .on('newPage', function(data) {
                    o.onNext(data);
                    o.onCompleted();
                }).take(1);
        });

        mergeObs.subscribe(
            (data) => {
                pr['merge'] = data;
                return pr; // this definitely isn't right
            },
            (error) => log.error(error),
            () => log.info('done')
        );
    });
});

通过一些日志记录,我可以看到拉取请求和合并 api 都被正确命中,但是当我订阅 pullRequestWithMergeStream 时,我 获取未定义的值。

map 内的subscribe 步骤中使用return 不起作用(而且似乎不应该),但我无法弄清楚什么模式/习语可以实现我想要的。

有正确的方法吗?我是不是完全走错了路?

tl;博士

我可以使用来自不同可观察对象的信息更新来自 Rxjs.Observable 的值吗?

【问题讨论】:

  • 对任何尝试做类似事情的人的快速说明:user3743222 下面的回答符合我的要求,但现在我在工作中看到它,我发现我无论如何都试图做一件坏事。这造成了一个很大的瓶颈,因为一次创建了几十个 observables,并且一次都轮询 Stash REST API。这严重影响了 Stash 服务器。回到绘图板。
  • 我更新了代码以限制同时通话的数量。无需重绘:-)

标签: javascript rxjs observable


【解决方案1】:

您可以使用flatMapconcatMap 让一个任务触发另一个任务。您可以使用forkJoin 并行请求合并并将结果收集到一个地方。它没有经过测试,但应该是这样的:

pullRequestStream.concatMap(function (prlist){
  var arrayRequestMerge = prlist.map(function(pr){
    return Rx.Observable.create(function(o) {...same as your code});
  });
  return Rx.Observable.forkJoin(arrayRequestMerge)
         .do(function(arrayData){
               prlist.map(function(pr, index){pr['merge']=arrayData[index]
             })})
         .map(function(){return prlist})
})

PS:我认为prlist 是一个数组。

更新 根据您的评论,这是一个仅并行运行 maxConcurrent 调用的版本。

pullRequestStream.concatMap(function (prlist){
  var arrayRequestMerge = prlist.map(function(pr, index){
    return Rx.Observable.create(function(o) {
        stash.pullRequestMerge(project, repo, pr['id'])
            .on('error', function(error) {o.onError(error)})
            .on('newPage', function(data) {
                o.onNext({data: data, index : index});
                o.onCompleted();
            }).take(1);
    });
  });
  var maxConcurrent = 2;
  Rx.Observable.from(arrayRequestMerge)
    .merge(maxConcurrent)
    .do(function(obj){
               prlist[obj.index]['merge'] = obj.data
             })})
    .map(function(){return prlist})
})

【讨论】:

  • 正是我想要的。非常感谢。我一直在尝试使用flatMap,但它似乎对我不起作用。我不知道concatMapforkJoin。给出的代码也可以开箱即用。
猜你喜欢
  • 1970-01-01
  • 2023-04-08
  • 1970-01-01
  • 2017-07-12
  • 1970-01-01
  • 1970-01-01
  • 2019-02-18
  • 2019-10-05
  • 1970-01-01
相关资源
最近更新 更多