【问题标题】:Repeat request until condition is met and return intermediate values重复请求直到满足条件并返回中间值
【发布时间】:2018-04-19 20:07:39
【问题描述】:

我正在尝试编写一个(通用)函数run<ID, ENTITY>(…): Observable<ENTITY>,它采用以下参数:

  • 一个函数init: () => Observable<ID>,它是一个用于启动后端进程的初始化请求。
  • status: (id: ID) => Observable<ENTITY> 函数获取生成的 ID 并在后端查询其状态。
  • 一个函数repeat: (status: ENTITY) => boolean 确定是否必须重复status 请求。
  • 两个整数值initialDelayrepeatDelay

所以run 应该执行init,然后等待initialDelay 秒。从现在开始,它应该每隔repeatDelay 秒运行一次status,直到repeat() 返回false

但是,有两件重要的事情需要解决:

  • repeatDelay 只应在 status 发出其值时开始计算,以避免在 status 花费的时间超过 repeatDelay 时出现竞争条件
  • status 的调用发出的中间值也必须发送给调用者。

除了我提到的最后一件事之外,以下(不是很漂亮)版本完成了所有操作:它在重试 status 之前不等待网络响应。

run<ID, ENTITY>(…): Observable<ENTITY> {
    let finished = false;
    return init().mergeMap(id => {
        return Observable.timer(initialDelay, repeatDelay)
            .switchMap(() => {
                if (finished) return Observable.of(null);
                return status(id);
            })
            .takeWhile(response => {
                if (repeat(response)) return true;
                if (finished) return false;

                finished = true;
                return true;
            });
    });
}

我的第二个版本是这样的,它同样适用于除一个细节之外的所有细节:status 调用的中间值不会发出,但我确实需要它们在调用者中显示进度:

run<ID, ENTITY>(…): Observable<ENTITY> {
    const loop = id => {
        return status(id).switchMap(response => {
            return repeat(response)
                ? Observable.timer(repeatDelay).switchMap(() => loop(id))
                : Observable.of(response);
        });
    };

    return init()
        .mergeMap(id => Observable.timer(initialDelay).switchMap(() => loop(id)));
}

诚然,后者也有点杂乱无章。我确信 rxjs 可以以更简洁的方式解决这个问题(更重要的是,完全解决它),但我似乎无法弄清楚如何。

【问题讨论】:

  • 循环在第二个例子中变得递归
  • 您的第二个版本可以通过使用Observable.timer(repeatDelay).switchMap(() =&gt; loop(id)).startWith(response) 之类的回复启动重试流来完成您需要的一切。通常,您无法真正摆脱递归,因为您可能需要无限期地重试并将每次尝试平面映射到前一个尝试中,并且大多数重试运算符仅提供 void 信号。如果您愿意处理错误,您可以获得的最接近的是retryWhen

标签: rxjs rxjs5


【解决方案1】:

更新:Observable 原生支持expand 的递归,@IngoBürk 的回答也显示了这一点。这让我们可以更简洁地编写递归:

function run<ENTITY>(/* ... */): Observable<ENTITY> {
  return init().delay(initialDelay).flatMap(id =>
    status(id).expand(s => 
      repeat(s) ? Observable.of(null).delay(repeatDelay).flatMap(_ => status(id)) : Observable.empty()
    )
  )
}

Fiddle.


如果递归是可以接受的,那么你可以更简洁地做事:

function run(/* ... */): Observable<ENTITY> {
  function recurse(id: number): Observable<ENTITY> {
    const status$ = status(id).share();
    const tail$ = status$.delay(repeatDelay)
                         .flatMap(status => repeat(status) ? recurse(id, repeatDelay) : Observable.empty());
    return status$.merge(tail$);
  }
  return init().delay(initialDelay).flatMap(id => recurse(id));
}

试试the fiddle

【讨论】:

  • 我认为这并不正确:它首先等待,然后依次调用 init 和 status。我需要的是调用init,等待initialDelay,调用状态,从现在开始总是等待repeatDelay + status。
  • @IngoBürk 啊,那时我的答案都不正确。我都更新了。
  • 是否可以将status() 调用限制为某个最大调用次数?然后发出错误值?
【解决方案2】:

repeatWhen 运算符本身看起来很诱人,但它只提供onComplete 通知的无效流,因此在没有外部帮助的情况下,您无法根据值决定重复。在这里,BehaviorSubject 可能是您最好的选择:

function run(/* ... */): Observable<ENTITY> {
  const last_value$ = new BehaviorSubject();
  const delayed$ = init()
    .delay(initialDelay)
    .flatMap(id => 
      status(id).repeatWhen(completions => 
         completions.delay(repeatDelay)
                    .takeWhile(_ => repeat(last_value$.getValue()))
      )
    ).share();
  delayed$.subscribe(last_value$);
  return delayed$;
}

这里,repeatWhen 仅在最后一个请求的最后一个值是表示重复的状态时才重新订阅请求的冷源。

Try the fiddle.

警告:当repeatDelay 较小时,通知程序(传递给repeatWhen 的函数)执行与接收相应状态的BehaviorSubject 之间可能存在竞争条件风险。对于给定的观察者,我们保证所有onNext 通知都发生在onCompleted 通知之前,但这里我们的BehaviorSubjectrepeatWhen 是分开的。从the source 一目了然,看起来onNext 通知直接通过repeatWhen 传递。我怀疑concat 也是如此,但我不确定。

【讨论】:

  • 谢谢!这看起来不错,但我实际上也想出了一个将事物链接起来的解决方案。如果您查看它,看看您是否认为它有问题,我将不胜感激!
【解决方案3】:

所以我想出了这个似乎可行的方法。它使用expand 来创建无限的重试序列,并使用takeWhile 来确定从中花费多长时间。

takeWhile hack 可以通过编写带有谓词函数的自定义运算符 takeUntil 来删除。这目前不存在,请参阅rxjs#2420

小提琴:https://jsfiddle.net/2mhcvnog/1/

run<ID, ENTITY>(...): Observable<ENTITY> {
    /* This little hack is needed to also emit the final item in the takeWhile() loop. */
    let finished = false;

    const delayStatus = (id, delay) => {
        return Observable.of(null)
            .delay(delay)
            .switchMap(() => status(id))
            .map(status => [id, status]);
    };

    return init()
        .mergeMap(id => delayStatus(id, initialDelay))
        .expand(([id]) => {
            if (finished) {
                return Observable.of([id, null]);
            }

            return delayStatus(id, repeatDelay);
        })
        .takeWhile(([_, status]) => {
            if (repeat(status)) {
                return true;
            }

            if (finished) {
                return false;
            }

            finished = true;
            return true;
        })
        .map(([_, status]) => status);
}

【讨论】:

  • 我认为 expand 可以在 like so 处自行完成所有工作。
  • @concat 是的,看起来好多了。您想将其发布为答案吗?这就是我要选择的,因此我想标记为已接受。
  • 感谢您的所有意见,非常感谢!
猜你喜欢
  • 2017-09-10
  • 1970-01-01
  • 2015-07-28
  • 1970-01-01
  • 1970-01-01
  • 2019-10-07
  • 2013-12-28
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多