【问题标题】:How to implement long pooling stream using rxjs?如何使用 rxjs 实现长池化流?
【发布时间】:2018-03-13 11:00:11
【问题描述】:
  1. 我打电话给angularService.method1 获取id,我需要下一个请求
  2. 我想每秒钟拨打一次angularService.method2(id),直到收到回复{success: true},或者:

    {success: false, errors: [{errorCode: 'code', errorMessage: 'some error message'}]
    

经过几个小时的尝试,我来到了这个版本,实际上它可以工作。但是我有一些问题

  1. 当主流达到“完成”时,内部的 Observable 会被销毁吗?
  2. 这是一个正确的实现还是我在这里有一些问题?

    this.angularService
    .method1(data)
    .flatMap((res) => {
        return Observable
            .interval(1000)
            .flatMap(() => this.angularService.method2(res.id))
            .concatMap(resp => {
                return (resp.success) ? Observable.of(resp, null) : Observable.of(resp);
            })
            .takeWhile(resp => resp);
    })
    .switchMap((res) => {
        return (!res.success && res.errors) ? Observable.throw(res.errors) : Observable.of(res);
    })
    .subscribe(
        (next) => console.log(next),
        (err) => console.log(err),
        () => console.log('finished')
    );
    

【问题讨论】:

    标签: angular rxjs


    【解决方案1】:

    我了解到您想调用method2,直到您收到第一个非空或未定义的响应,然后关闭流。

    如果这是正确的,那么我将开始编写一个接收id 参数的函数,然后每秒调用一次method2,直到得到答案,类似于

    callMethod2EverySecond(id) {
      return Observable.interval(1000)
             .mergeMap(() => this.angularService.method2(id))
             .filter(resp => resp !== null)
             .take(1)
    }
    

    此函数执行以下操作

    • 每隔一秒调用一次method2,以id为参数——假设 method2 返回一个对象(带有正确的响应或 错误)或 null - 认为 mergeMap 是新的 flatMap的名字
    • 通过过滤器过滤掉所有保持为空的事件 operator
    • 获取第一个不为空的事件并关闭流 (这是由take 操作员执行的)

    然后您可以在外部逻辑中使用该功能,例如

    this.angularService
    .method1(data)
    .switchMap(res => callMethod2EverySecond(res.id))
    .subscribe(
        res => {
           if (!res.success && res.errors) {console.error(res)}
           else {console.log(res)}
        },
        (err) => console.log(err),
        () => console.log('finished')
    );
    

    考虑以下几点:

    • switchMap 和 mergeMap(或 flatMap)在这种情况下可能会产生 结果相同,但不一样 (read this for more details)
    • 因为你的错误情况可以通过分析内容来推断 响应,您可以直接在定义的第一个函数中执行此操作 作为subscribe 的参数,您不需要抛出 Observables 作为错误

    模拟 angularService.method1 和 angularService.method2 的工作代码示例

    以下是根据上述假设的代码的工作示例。

    已经模拟了angularService的method1method2

    import {Observable} from 'rxjs';
    
    method1('123')
    .switchMap(res => callMethod2EverySecond(res.id))
    .subscribe(
        res => {
           if (!res.success && res.errors) {console.error(res)}
           else {console.log('subscription processing', res)}
        },
        (err) => console.log(err),
        () => console.log('finished')
    );
    
    function callMethod2EverySecond(id) {
        return Observable.interval(10)
               .mergeMap(data => method2(id, data))
               .do(resp => console.log('resp', resp))
               .filter(resp => resp !== null)
               .take(1)
    }
    
    function method1(data: string) {
        return Observable.of({id: data});
    }
    
    function method2(id: string, interval: number) {
        const ret = randomIntInc(0,1) === 0 ? null : {success: true, errors: null, interval, id};
        return Observable
                .of(ret)
                .delay(randomIntInc(0,2000));
    }
    
    function randomIntInc(low, high) {
        return Math.floor(Math.random() * (high - low + 1) + low);
    }
    

    【讨论】:

    • Alexander 想要达到的目标还不是很清楚,至少对我来说是这样。你想调用 method2 直到你得到第一个不是 null 或未定义的响应,然后关闭流,这应该可以工作(除非我犯了一个错误 - 我无法从我所在的位置测试代码)
    • @Picci 不完全是。我只想调用一次method1,结果我收到了一些id。然后我想调用 method2 并重复调用 method2 直到我得到一些结果。它可能在第一次通话或第 7 次通话时发生,没关系
    • @Alexander Ponomarev - 我已经用工作代码示例编辑了我的答案,其中我模拟了方法 1 和方法 2 - 示例工作正常,至少根据我对您希望行为方式的理解
    猜你喜欢
    • 1970-01-01
    • 2017-10-12
    • 2018-07-12
    • 2021-09-08
    • 2017-10-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-08-02
    相关资源
    最近更新 更多