← Back to list

Чекати, але не вічно: пишемо свій RxJS оператор

У вас є дві важливі поїздки в різні частини міста. Один автобус — А, інший — Б. Ви не знаєте, який приїде швидше, чи взагалі приїдуть…

Eugene Rusakov · 2025-05-10 16:55 · 63 claps · 3.5 min read
#rxjs #custom-operator
Open on Medium ↗

Чекати, але не вічно: пишемо свій RxJS оператор

У вас є дві важливі поїздки в різні частини міста. Один автобус — А, інший — Б. Ви не знаєте, який приїде швидше, чи взагалі приїдуть обидва. Але у вас є максимум скажімо 10 хвилин, щоб чекати — далі треба йти.

Пишемо кастомний оператор який буде робити наступне:

  • Слухає, чи приїхав автобус А.
  • Слухає, чи приїхав автобус Б.
  • Чекає не більше ніж N часу.
  • Видає чіткий звіт: хто приїхав, а хто ні.
{
  busASuccess: true,   // або false
  busBSuccess: false   // або true
}

Погнали!

1. Створимо новий Observable

Розпочнемо з того, що створимо новий Observable, який поверне результат після виконання наших умов:

return new Observable<{ busASuccess: boolean; busBSuccess: boolean }>((subscriber) => {}

Цей Observable, буде видавати об’єкт з двома булевими значеннями — busASuccess та busBSuccessпісля виконання необхідної логіки.

2. Змінні, щоб памʼятати, хто вже приїхав

// Змінна для фіксації, чи приїхав Bus A
let busASuccess: boolean | undefined; 

// Змінна для фіксації, чи приїхав Bus B
let busBSuccess: boolean | undefined;

// Маркер, який вказує, чи був вже надісланий результат (еміт)
let hasEmitted = false;

3. Створюємо функцію checkAndEmit

Ця функція відповідає за те, щоб вирішити, коли вже час видавати результат:

const checkAndEmit = (forceEmit = false): void => {
  // Якщо результат вже було емічено — нічого не робимо
  if (hasEmitted) return

  // Чи отримали ми результат від busA$
  const hasA = busASuccess !== undefined

  // Чи отримали ми результат від busB$
  const hasB = busBSuccess !== undefined

  if ((hasA && hasB) || forceEmit) {
    // Виставляємо маркер, щоб не емітити повторно
    hasEmitted = true

    // Видаємо результат (усі значення, що не прийшли — замінюємо на false)
    subscriber.next({
      busASuccess: busASuccess ?? false,
      busBSuccess: busBSuccess ?? false
    })

    // Завершуємо стрім
    subscriber.complete()
  }
}

4. Підписуємось на busA$ і busB$

// Підписуємось на потік busA$ (умовно — "автобус А приїхав чи ні")
const busASub = busA$.subscribe({
  next: (value: boolean | undefined) => {
    // Зберігаємо отримане значення в локальну змінну
    busASuccess = value
    // Перевіряємо, чи час емітити результат
    checkAndEmit()
  },
  error: (err: any) => {
    // Якщо busA$ викинув помилку — передаємо її далі
    subscriber.error(err)
  }
})

// Підписуємось на потік busB$ (умовно — "автобус B приїхав чи ні")
const busBSub = busB$.subscribe({
  next: (value: boolean | undefined) => {
    // Зберігаємо отримане значення в локальну змінну
    busBSuccess = value
    // Перевіряємо, чи час емітити результат
    checkAndEmit()
  },
  error: (err: any) => {
    // Якщо busB$ викинув помилку — передаємо її далі
    subscriber.error(err)
  }
})

5. Додаємо таймер

const timeoutId = setTimeout(() => {
  checkAndEmit(true)
}, timeoutMs)

Після завершення все відписується й чиститься

return () => {
  clearTimeout(timeoutId)
  busASub.unsubscribe()
  busBSub.unsubscribe()
}

Повний код нашого кастомного оператора з прикладом:

import { Observable, of, timer } from 'rxjs'
import { delay } from 'rxjs/operators'

/**
 * Чекає на значення з двох потоків (наприклад, автобусів A і B) з таймаутом.
 * Якщо обидва потоки видадуть значення — повертає їх.
 * Якщо ні — після таймауту повертає те, що встигло прийти (або false, якщо нічого не прийшло).
 */
export function waitForBooleanStatusWithTimeout(
  busA$: Observable<boolean>,
  busB$: Observable<boolean>,
  timeoutMs: number = 3000
): Observable<{ busASuccess: boolean; busBSuccess: boolean }> {
  return new Observable<{ busASuccess: boolean; busBSuccess: boolean }>((subscriber) => {
    let busASuccess: boolean | undefined
    let busBSuccess: boolean | undefined
    let hasEmitted = false

    /**
     * Перевіряє чи можна емітити результат, або форсує завершення після таймауту
     * @param forceEmit - чи форсувати викид навіть якщо не всі потоки відповіли
     */
    const checkAndEmit = (forceEmit = false): void => {
      if (hasEmitted) return

      const hasA = busASuccess !== undefined
      const hasB = busBSuccess !== undefined

      if ((hasA && hasB) || forceEmit) {
        hasEmitted = true
        subscriber.next({
          busASuccess: busASuccess ?? false,
          busBSuccess: busBSuccess ?? false
        })
        subscriber.complete()
      }
    }

    // Підписка на потік busA$
    const busASub = busA$.subscribe({
      next: (value: boolean | undefined) => {
        busASuccess = value
        checkAndEmit()
      },
      error: (err: any) => subscriber.error(err)
    })

    // Підписка на потік busB$
    const busBSub = busB$.subscribe({
      next: (value: boolean | undefined) => {
        busBSuccess = value
        checkAndEmit()
      },
      error: (err: any) => subscriber.error(err)
    })

    // Запускаємо таймер: якщо ніхто або лише хтось один відповів — форсуємо завершення
    const timeoutId = setTimeout(() => {
      checkAndEmit(true)
    }, timeoutMs)

    // Очищення всіх ресурсів після завершення
    return () => {
      clearTimeout(timeoutId)
      busASub.unsubscribe()
      busBSub.unsubscribe()
    }
  })
}

// Приклад використання:
const busA$ = of(true).pipe(delay(1000))    // автобус A приїде через 1 сек
const busB$ = of(false).pipe(delay(5000))   // автобус B не встигне (таймаут — 3 сек)

waitForBooleanStatusWithTimeout(busA$, busB$, 3000).subscribe({
  next: ({ busASuccess, busBSuccess }) => {
    console.log('Результат:', { busASuccess, busBSuccess })
    // Очікується: { busASuccess: true, busBSuccess: false }
  },
  complete: () => {
    console.log('Чекання завершено.')
  }
})

Висновок

Досить часто потрібно одночасно чекати на кілька відповідей — наприклад, чи працюють обидва потрібні сервіси і якщо один з них не відповів, ми чекаємо якийсь час і завершуємо нашу логіку.

Це рішення стане в пригоді, коли ви маєте два Subject'и, але стандартні RxJS-оператори — combineLatest, merge чи race — не підходять. Наприклад, якщо один із потоків може взагалі не емітити, і вам потрібно завершити очікування самостійно через таймаут, зі спецефічною логікою обробки.

Наш оператор waitForTwoBusesOrTimeout вийшов як зручна абстракція для кейсів, де треба дочекатись кількох подій, але не чекати їх вічно.

[embed]


메타데이터
post_id
b0a08bb3b8cb
slug
чекати-але-не-вічно-пишемо-свій-rxjs-оператор-b0a08bb3b8cb
url
https://medium.com/@eurusik/%D1%87%D0%B5%D0%BA%D0%B0%D1%82%D0%B8-%D0%B0%D0%BB%D0%B5-%D0%BD%D0%B5-%D0%B2%D1%96%D1%87%D0%BD%D0%BE-%D0%BF%D0%B8%D1%88%D0%B5%D0%BC%D0%BE-%D1%81%D0%B2%D1%96%D0%B9-rxjs-%D0%BE%D0%BF%D0%B5%D1%80%D0%B0%D1%82%D0%BE%D1%80-b0a08bb3b8cb
canonical_url
https://medium.com/@eurusik/%D1%87%D0%B5%D0%BA%D0%B0%D1%82%D0%B8-%D0%B0%D0%BB%D0%B5-%D0%BD%D0%B5-%D0%B2%D1%96%D1%87%D0%BD%D0%BE-%D0%BF%D0%B8%D1%88%D0%B5%D0%BC%D0%BE-%D1%81%D0%B2%D1%96%D0%B9-rxjs-%D0%BE%D0%BF%D0%B5%D1%80%D0%B0%D1%82%D0%BE%D1%80-b0a08bb3b8cb
author_url
https://medium.com/@eurusik
status
ok
fetched_at
2026-08-02 03:19:08