Subscribe to multiple inner observables concurrently and emit results in order
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
import { | |
from, | |
BehaviorSubject, | |
Observable, | |
ObservableInput, | |
ObservedValueOf, | |
OperatorFunction, | |
MonoTypeOperatorFunction, | |
} from 'rxjs'; | |
import { mergeMap, delayWhen, finalize, skipWhile, take } from 'rxjs/operators'; | |
const delayUntilIsActiveIndex = <O extends ObservableInput<any>>( | |
activeIndex$: Observable<number>, | |
innerIndex: number | |
): MonoTypeOperatorFunction<O> => | |
delayWhen(() => | |
activeIndex$.pipe( | |
skipWhile(activeIdx => activeIdx !== innerIndex), | |
take(1) | |
) | |
); | |
const bumpIndexOnComplete = <O extends ObservableInput<any>>( | |
activeIndex$: BehaviorSubject<number>, | |
innerIndex: number | |
): MonoTypeOperatorFunction<O> => finalize(() => activeIndex$.next(innerIndex + 1)); | |
export function orderedMergeMap<T, O extends ObservableInput<any>>( | |
project: (value: T, index: number) => O, | |
concurrent = 5 | |
): OperatorFunction<T, ObservedValueOf<O>> { | |
const activeIndex$ = new BehaviorSubject(0); | |
return mergeMap( | |
(value: T, index: number) => | |
from(project(value, index)).pipe( | |
delayUntilIsActiveIndex(activeIndex$, index), | |
bumpIndexOnComplete(activeIndex$, index) | |
), | |
concurrent | |
); | |
} |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment