Estoy tratando de usar RxJS para procesar un flujo de elementos, y me gustaría que vuelva a intentar cualquier falla, preferiblemente con un retraso (y una reducción exponencial si es posible), pero necesito garantizar el pedido, así que efectivamente quiero bloquear la transmisión y no quiero volver a procesar ningún elemento que ya se haya procesado
Así que estaba tratando de jugar con retryWhen , siguiendo este ejemplo:
const { interval, timer } = Rx; const { take, map, retryWhen, delayWhen, tap } = RxOperators; const source = take(5)(interval(1000)); source.pipe( map(val => { if (val >= 3) { throw val; } return val; }), retryWhen(errors => errors.pipe( delayWhen(val => timer(1000)) ) ) );Pero la transmisión se reinicia al principio, no solo vuelve a intentar la última:
¿Es posible lograr lo que quiero? También probé con otros operadores de docs, sin suerte. ¿Sería un poco en contra de la filosofía RxJS de alguna manera?
El retryWhen debe moverse a un Observable interno para manejar solo los valores fallidos y mantener el Observable principal en funcionamiento.
Prueba algo como lo siguiente:
// import { timer, interval, of } from 'rxjs'; // import { concatMap, delayWhen, map, retryWhen, take } from 'rxjs/operators'; const source = interval(1000).pipe(take(5)); source.pipe( concatMap((value) => of(value).pipe( map((val) => { if (val >= 3) { throw val; } return val; }), retryWhen((errors) => errors.pipe(delayWhen(() => timer(1000)))) ) ) );