Чекати, але не вічно: пишемо свій RxJS оператор
У вас є дві важливі поїздки в різні частини міста. Один автобус — А, інший — Б. Ви не знаєте, який приїде швидше, чи взагалі приїдуть…
Чекати, але не вічно: пишемо свій 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