Skip to content

Instantly share code, notes, and snippets.

Embed
What would you like to do?
Subscribe to multiple inner observables concurrently and emit results in order
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