我认为文档中给出的示例是为仅发出一次然后完成的 Observable 编写的,例如 http get。假设如果您想获得更多数据,那么您将再次订阅,这将重置genericRetryStrategy 内的计数器。但是,如果您现在想将相同的策略应用到一个长时间运行的 observable,其流不会完成,除非它给出错误(例如您使用 interval()),那么您需要修改 genericRetryStrategy()在需要重置计数器时被告知。
这可以通过多种方式完成,我在StackBlitz 中给出了一个简单的例子,基于你所说的你想要完成的事情。请注意,我还稍微更改了您的逻辑,以更符合您所说的您正在尝试做的事情,即“2 次成功尝试,然后 2 次不成功尝试”。重要的一点是修改被抛出到genericRetryStrategy() 的错误对象,以传达当前失败尝试的计数,以便它能够做出适当的反应。
为了完整起见,这里是复制的代码:
import { timer, interval, Observable, throwError } from 'rxjs';
import { map, switchMap, tap, retryWhen, delayWhen, mergeMap, shareReplay, finalize, catchError } from 'rxjs/operators';
console.clear();
interface Err {
status?: number;
msg?: string;
int: number;
}
export const genericRetryStrategy = ({
maxRetryAttempts = 3,
scalingDuration = 1000,
excludedStatusCodes = []
}: {
maxRetryAttempts?: number,
scalingDuration?: number,
excludedStatusCodes?: number[]
} = {}) => (attempts: Observable<any>) => {
return attempts.pipe(
mergeMap((error: Err) => {
// i here does not reset and continues to increment?
const retryAttempt = error.int;
// if maximum number of retries have been met
// or response is a status code we don't wish to retry, throw error
if (
retryAttempt > maxRetryAttempts ||
excludedStatusCodes.find(e => e === error.status)
) {
return throwError(error);
}
console.log(
`Attempt ${retryAttempt}: retrying in ${retryAttempt *
scalingDuration}ms`
);
// retry after 1s, 2s, etc...
return timer(retryAttempt * scalingDuration);
}),
finalize(() => console.log('We are done!'))
);
};
let int = 0;
let err: Err = {int: 0};
//emit value every 1s
interval(1000).pipe(
map((val) => {
if (val > 1) {
//error will be picked up by retryWhen
int++;
err.msg = "equals 1";
err.int = int;
throw err;
}
if (val === 0 && int === 1) {
err.msg = "greater than 2";
err.int = 2;
int=0;
throw err;
}
return val;
}),
retryWhen(genericRetryStrategy({
maxRetryAttempts: 3,
scalingDuration: 1000,
excludedStatusCodes: [],
}))
).subscribe(val => {
console.log(val)
});
对我来说,这仍然是非常必要的,但是如果不了解您要更深入地解决的问题,我目前想不出更具声明性的方法...