RxJS: операторы
Higher-Order Observable Mapping
При higher-order mapping значение источника отображается не в плоское значение, а в другой Observable. Результат — Observable высшего порядка, значения которого сами являются Observable.
| Оператор | Поведение |
|---|---|
concatMap | Не подписывается на следующий Observable, пока не завершится предыдущий (по порядку). |
mergeMap | Не ждёт завершения предыдущего — несколько внутренних Observable работают параллельно. |
switchMap | Если пришло новое значение до завершения предыдущего Observable — отписывается от него, значения больше не учитываются. |
exhaustMap | Отображает во внутренний Observable, игнорирует новые значения, пока текущий не завершится. |
Комбинирование потоков
merge— объединяет потоки, эмитит значения по мере поступления.concat— последовательно, следующий поток после завершения предыдущего.zip— попарно объединяет значения по индексу.forkJoin— ждёт завершения всех, отдаёт последние значения каждого.combineLatest— эмитит последние значения всех потоков при каждом новом значении любого.withLatestFrom— берёт последнее значение из других потоков в момент эмита основного.
partition
Разделяет поток на две части по условию.
const [evens, odds] = partition(
of(1, 2, 3, 4).pipe(map((n) => n * 3)),
(n) => n % 2 === 0
);takeUntil
Эмитит значения источника, пока не эмитит переданный notifier$.
const stop$ = new Subject<void>();
source$.pipe(takeUntil(stop$)).subscribe();
function stop() {
stop$.next();
stop$.complete();
}finalize
Вызывает функцию, когда источник завершается (complete/error) или подписчик отписывается — независимо от причины.
const example = interval(1000).pipe(
take(3),
finalize(() => console.log("Sequence complete"))
);
// 0, 1, 2, 'Sequence complete'ignoreElements
Подавляет все эмитированные значения, но пропускает уведомление о завершении (complete/error).
of(1, 6, 5, 10).pipe(ignoreElements()).subscribe({
next: (x) => console.log(x), // не вызовется
complete: () => console.log("done"), // done
});