RxJS: Subjects и Schedulers

Subject

Subject — особый объект RxJS, который одновременно является Observable (может отправлять данные) и Observer (может подписываться на поток).

BehaviorSubject

Хранит последнее отправленное значение. Каждый новый подписчик в момент subscribe() сразу получает это значение.

const sbj = new BehaviorSubject<number>(5);
sbj.subscribe((vl) => console.log(`1st: ${vl}`));
sbj.subscribe((vl) => console.log(`2nd: ${vl}`));
sbj.next(7);
// 1st: 5, 2nd: 5, 1st: 7, 2nd: 7

ReplaySubject

Хранит заданное количество последних значений (задаётся в конструкторе). Новые подписчики сразу получают все n сохранённых значений.

const sbj = new ReplaySubject(2);
sbj.next(5);
sbj.subscribe((vl) => console.log(`1st: ${vl}`)); // 5
sbj.next(6);
sbj.next(7);
sbj.subscribe((vl) => console.log(`2nd: ${vl}`)); // 6, 7

AsyncSubject

Передаёт подписчикам только последнее значение и только после complete().

const sbj = new AsyncSubject();
sbj.subscribe((vl) => console.log(`Async: ${vl}`));
sbj.next(7);
sbj.next(8);
sbj.next(9);
sbj.complete();
// Async: 9

Свойства closed / isStopped

СитуацияclosedisStopped
На errorfalsetrue
На completionfalsetrue
После unsubscribetruetrue

Schedulers (планировщики)

Scheduler определяет, в каком контексте выполнения Observable доставляет уведомления Observer’у. Соответствуют очередям событийного цикла:

SchedulerКонтекст
queueSchedulerсинхронно (call stack)
asapSchedulerмикрозадача (Promise)
asyncSchedulerмакрозадача (setTimeout)
animationFrameSchedulerперед перерисовкой (requestAnimationFrame)
  • observeOn — планирует контекст, в котором эмитятся значения (next/error/complete).
  • subscribeOn — планирует контекст, в котором происходит сам вызов subscribe().