【发布时间】: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请求。 - 两个整数值
initialDelay和repeatDelay。
所以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(() => loop(id)).startWith(response)之类的回复启动重试流来完成您需要的一切。通常,您无法真正摆脱递归,因为您可能需要无限期地重试并将每次尝试平面映射到前一个尝试中,并且大多数重试运算符仅提供 void 信号。如果您愿意处理错误,您可以获得的最接近的是retryWhen。