Pytania rekrutacyjne RxJS

47 pytań — 20 odpowiedzi dostępnych od razu

Darmowe pytania rekrutacyjne RxJS pomagają przygotować się do rozmowy o programowaniu reaktywnym i pracy ze strumieniami danych. Ten podgląd prezentuje przykładowe pytania z zestawu 47 zagadnień.

Materiał obejmuje Observable, operatory tworzące, transformujące, filtrujące i łączące. Dzięki fiszkom możesz przećwiczyć dobór operatorów oraz wyjaśnić, jak zarządzać przepływem zdarzeń w aplikacji.

Darmowy start

Pierwsze 20 odpowiedzi są dostępne od razu

Przeglądaj pełną listę 47 pytań. Pełne odpowiedzi, przykłady kodu i materiały do nauki odblokujesz w panelu.

Odblokuj 27 odpowiedzi

Podstawy RxJS

Czym jest RxJS i jakie problemy rozwiązuje w aplikacjach JavaScript?

Odpowiedź w 30 sekund: RxJS (Reactive Extensions for JavaScript) to biblioteka do programowania reaktywnego wykorzystująca strumienie danych (Observable). Rozwiązuje problemy związane z zarządzaniem asynchronicznością, złożonymi przepływami danych i koordynacją wielu źródeł zdarzeń w aplikacjach JavaScript.

Odpowiedź w 2 minuty: RxJS to potężna biblioteka implementująca wzorzec Observer i programowanie reaktywne w JavaScript. Pozwala traktować wszystkie asynchroniczne operacje - od kliknięć użytkownika, przez żądania HTTP, po timery - jako ujednolicone strumienie danych, które można komponować, transformować i łączyć.

Główne problemy, które rozwiązuje RxJS to: zarządzanie callback hell poprzez deklaratywne operatory, łatwe anulowanie asynchronicznych operacji (subscriptions), elegancka obsługa błędów w złożonych przepływach danych oraz zaawansowane scenariusze jak debouncing, throttling, retry logic czy łączenie wielu źródeł danych. RxJS jest szczególnie popularna w ekosystemie Angular, ale sprawdza się w każdej aplikacji JavaScript wymagającej zaawansowanego zarządzania asynchronicznością.

Biblioteka oferuje ponad 100 operatorów pozwalających na transformację, filtrowanie, łączenie i kontrolę strumieni danych w sposób funkcyjny i deklaratywny. Dzięki temu kod staje się bardziej czytelny, testowalny i łatwiejszy w utrzymaniu niż tradycyjne podejścia z callbackami czy Promise.

Przykład kodu:

import { fromEvent } from 'rxjs';
import { debounceTime, map, distinctUntilChanged } from 'rxjs/operators';

// Pole wyszukiwania z opóźnionym wysyłaniem zapytań
const searchInput = document.getElementById('search');

fromEvent(searchInput, 'input')
  .pipe(
    map((event: any) => event.target.value), // Wyciągnij wartość
    debounceTime(300),                         // Czekaj 300ms po ostatnim wpisie
    distinctUntilChanged()                     // Ignoruj jeśli wartość się nie zmieniła
  )
  .subscribe(searchTerm => {
    console.log('Szukam:', searchTerm);
    // Tutaj wywołaj API z wyszukiwaniem
  });
graph LR
    A[Zdarzenia input] --> B[map - wyciągnij wartość]
    B --> C[debounceTime - opóźnienie]
    C --> D[distinctUntilChanged - unikalne]
    D --> E[subscribe - wywołaj API]

Materiały

Czym jest Observable i czym różni się od Promise?

Odpowiedź w 30 sekund: Observable to strumień danych, który może emitować wiele wartości w czasie, obsługuje leniwą ewaluację i pozwala na anulowanie. Promise reprezentuje pojedynczą wartość asynchroniczną, wykonuje się od razu i nie można go anulować.

Odpowiedź w 2 minuty: Observable to podstawowy typ w RxJS reprezentujący strumień danych, który może emitować zero, jedną lub wiele wartości w czasie. Observable jest "leniwy" (lazy) - kod wewnątrz nie wykonuje się dopóki ktoś się nie zasubskrybuje. Każda subskrypcja tworzy niezależne wykonanie (dla cold Observable).

Kluczowe różnice względem Promise:

  1. Ilość wartości: Promise zwraca dokładnie jedną wartość (lub błąd), Observable może emitować wiele wartości przez cały okres życia
  2. Ewaluacja: Promise rozpoczyna wykonanie natychmiast po utworzeniu (eager), Observable tylko po subskrypcji (lazy)
  3. Anulowanie: Promise nie może być anulowany po rozpoczęciu, Observable można w każdej chwili anulować poprzez unsubscribe
  4. Operatory: Observable oferuje bogaty zestaw operatorów do transformacji (map, filter, merge itp.), Promise ma tylko then/catch/finally

Observable świetnie sprawdza się do obsługi WebSocketów, zdarzeń UI, intervalów czy dowolnych strumieni danych zmieniających się w czasie. Promise lepiej pasuje do pojedynczych operacji HTTP gdzie potrzebujesz tylko jednej odpowiedzi.

Przykład kodu:

// PROMISE - wykonuje się natychmiast, jedna wartość
const promise = new Promise((resolve) => {
  console.log('Promise wykonany!'); // Loguje się od razu
  setTimeout(() => resolve('Wynik'), 1000);
});

promise.then(value => console.log(value));
// Nie można anulować!

// OBSERVABLE - leniwy, wiele wartości, można anulować
import { Observable } from 'rxjs';

const observable = new Observable(subscriber => {
  console.log('Observable wykonany!'); // Loguje się tylko po subscribe
  let count = 0;
  const interval = setInterval(() => {
    subscriber.next(count++); // Emituj wiele wartości
  }, 1000);

  // Funkcja czyszcząca - wykonana przy unsubscribe
  return () => {
    console.log('Anulowano!');
    clearInterval(interval);
  };
});

const subscription = observable.subscribe(value => console.log(value));

// Anuluj po 5 sekundach
setTimeout(() => subscription.unsubscribe(), 5000);
graph TD
    A[Promise] --> B[Eager - wykonuje się od razu]
    A --> C[Jedna wartość]
    A --> D[Nie można anulować]

    E[Observable] --> F[Lazy - wykonuje się po subscribe]
    E --> G[Wiele wartości w czasie]
    E --> H[Można anulować - unsubscribe]

Materiały

Jak działa subskrypcja (Subscription) i dlaczego ważne jest jej anulowanie?

Odpowiedź w 30 sekund: Subscription reprezentuje wykonanie Observable i zwraca obiekt z metodą unsubscribe(). Anulowanie subskrypcji jest kluczowe aby uniknąć wycieków pamięci, niepotrzebnych operacji i błędów w aplikacji gdy komponenty/obiekty przestają istnieć.

Odpowiedź w 2 minuty: Subscription to obiekt zwracany przez metodę subscribe() reprezentujący trwające wykonanie Observable. Zawiera metodę unsubscribe(), która pozwala przerwać strumień danych i wykonać cleanup. Gdy wywołasz subscribe(), Observable rozpoczyna emitowanie wartości do przekazanych funkcji callback (next, error, complete).

Anulowanie subskrypcji jest krytyczne z kilku powodów:

  1. Wycieki pamięci: Nieodsubskrybowane Observable (np. interwały, WebSockety) działają w tle zajmując pamięć nawet gdy nie są już potrzebne
  2. Niepotrzebne operacje: Zapytania HTTP, timery czy event listenery działają niepotrzebnie zużywając zasoby
  3. Błędy w aplikacji: Aktualizacje usuniętych komponentów (w React/Angular) prowadzą do błędów "cannot set state of unmounted component"
  4. Problemy z wydajnością: Wielokrotne subskrypcje bez cleanup mogą drastycznie spowolnić aplikację

W Angular możesz użyć AsyncPipe który automatycznie anuluje subskrypcje. W innych frameworkach używaj wzorca przechowywania subskrypcji i anulowania ich w lifecycle hooks (useEffect cleanup, ngOnDestroy, componentWillUnmount).

Przykład kodu:

import { interval, Subscription } from 'rxjs';

// Pojedyncza subskrypcja
const subscription = interval(1000).subscribe(value => {
  console.log('Timer:', value);
});

// Anuluj po 5 sekundach
setTimeout(() => {
  subscription.unsubscribe();
  console.log('Subskrypcja anulowana');
}, 5000);

// Zarządzanie wieloma subskrypcjami
class UserDashboard {
  private subscriptions = new Subscription();

  init() {
    // Dodaj wiele subskrypcji do jednego kontenera
    this.subscriptions.add(
      interval(1000).subscribe(x => console.log('Timer 1:', x))
    );

    this.subscriptions.add(
      interval(2000).subscribe(x => console.log('Timer 2:', x))
    );
  }

  destroy() {
    // Anuluj wszystkie subskrypcje jednocześnie
    this.subscriptions.unsubscribe();
    console.log('Wszystkie subskrypcje anulowane');
  }
}

// Przykład z React Hook
import { useEffect } from 'react';

function MyComponent() {
  useEffect(() => {
    const subscription = interval(1000).subscribe(x => {
      console.log('Wartość:', x);
    });

    // Cleanup - wykonany gdy komponent się odmontuje
    return () => subscription.unsubscribe();
  }, []);

  return <div>Komponent z Observable</div>;
}
sequenceDiagram
    participant Component
    participant Observable
    participant Subscription

    Component->>Observable: subscribe()
    Observable->>Subscription: return Subscription
    Observable->>Component: next(value1)
    Observable->>Component: next(value2)
    Component->>Subscription: unsubscribe()
    Subscription->>Observable: cleanup()
    Note over Observable: Zatrzymuje emitowanie

Materiały

Czym jest Observer i jakie metody udostępnia (next, error, complete)?

Odpowiedź w 30 sekund: Observer to obiekt z metodami callback definiującymi jak reagować na emitowane wartości. Zawiera trzy metody: next(value) - obsługa nowej wartości, error(err) - obsługa błędu, complete() - zakończenie strumienia.

Odpowiedź w 2 minuty: Observer to obiekt implementujący interfejs z trzema opcjonalnymi metodami callback, które określają jak konsument reaguje na notyfikacje z Observable. Jest to implementacja wzorca Observer, gdzie Observable to subject (podmiot), a Observer to obserwator reagujący na zmiany.

Trzy metody Observer:

  1. next(value): Wywoływana za każdym razem gdy Observable emituje wartość. Może być wywołana 0 lub więcej razy. To główna metoda obsługująca strumień danych.
  2. error(error): Wywoływana gdy wystąpi błąd w Observable. Po wywołaniu error strumień się kończy i nie będzie więcej emisji (ani next, ani complete). Wywołana maksymalnie raz.
  3. complete(): Sygnalizuje zakończenie strumienia bez błędu. Po complete nie będzie już żadnych emisji. Wywołana maksymalnie raz. Niektóre Observable nigdy się nie kończą (np. interval).

Możesz przekazać Observer jako obiekt lub jako osobne funkcje do subscribe(). Wszystkie metody są opcjonalne - możesz np. obsłużyć tylko next ignorując error i complete.

Przykład kodu:

import { Observable } from 'rxjs';

// Observable emitujący wartości
const observable = new Observable(subscriber => {
  subscriber.next(1);
  subscriber.next(2);
  subscriber.next(3);

  // Symulacja błędu (odkomentuj aby przetestować)
  // subscriber.error(new Error('Coś poszło nie tak!'));

  subscriber.complete();

  // To nie zostanie wysłane - po complete nic nie emituje
  subscriber.next(4);
});

// Observer jako obiekt
const observer = {
  next: (value: number) => {
    console.log('Otrzymano wartość:', value);
  },
  error: (err: Error) => {
    console.error('Wystąpił błąd:', err.message);
  },
  complete: () => {
    console.log('Strumień zakończony!');
  }
};

observable.subscribe(observer);

// Alternatywnie - przekaż funkcje bezpośrednio
observable.subscribe(
  value => console.log('Next:', value),        // next
  err => console.error('Error:', err),         // error
  () => console.log('Complete!')               // complete
);

// Lub tylko next (pozostałe opcjonalne)
observable.subscribe(value => console.log(value));

// Praktyczny przykład - zapytanie HTTP
import { ajax } from 'rxjs/ajax';

ajax.getJSON('https://api.example.com/users').subscribe({
  next: users => {
    console.log('Pobrano użytkowników:', users);
    // Zaktualizuj UI
  },
  error: err => {
    console.error('Błąd pobierania:', err);
    // Pokaż komunikat błędu użytkownikowi
  },
  complete: () => {
    console.log('Zapytanie zakończone');
    // Ukryj loader
  }
});
stateDiagram-v2
    [*] --> Active: subscribe()
    Active --> Active: next(value)
    Active --> Error: error(err)
    Active --> Complete: complete()
    Error --> [*]
    Complete --> [*]

    note right of Active
        Może emitować wiele
        wartości przez next()
    end note

    note right of Error
        Terminal state
        Kończy strumień
    end note

    note right of Complete
        Terminal state
        Kończy strumień
    end note

Materiały

Jaka jest różnica między cold i hot Observable?

Odpowiedź w 30 sekund: Cold Observable tworzy nowe, niezależne wykonanie dla każdej subskrypcji (unicast) - każdy subscriber dostaje własny strumień od początku. Hot Observable współdzieli jedno wykonanie między wszystkich subscribers (multicast) - wszyscy dostają te same wartości w tym samym czasie.

Odpowiedź w 2 minuty: Cold i Hot Observable różnią się sposobem tworzenia i współdzielenia strumienia danych między subskrybentami.

Cold Observable (unicast):

  • Tworzy nowe wykonanie dla każdej subskrypcji
  • Każdy subscriber otrzymuje własną, niezależną kopię wartości
  • Strumień rozpoczyna się od początku dla każdego nowego subscribera
  • Przykłady: HTTP requests, timery, of(), from(), interval()
  • Producent danych jest tworzony wewnątrz Observable

Hot Observable (multicast):

  • Współdzieli jedno wykonanie między wszystkich subscribers
  • Wszyscy subskrybenci otrzymują te same wartości w tym samym czasie
  • Nowy subscriber otrzymuje tylko wartości emitowane PO jego subskrypcji
  • Przykłady: DOM events, WebSocket streams, Subject
  • Producent danych istnieje niezależnie od Observable

Możesz przekształcić Cold Observable w Hot używając operatorów share(), shareReplay() lub Subject. Jest to przydatne gdy chcesz uniknąć wielokrotnego wykonywania kosztownych operacji (np. HTTP requests) dla każdego subscribera.

Przykład kodu:

import { Observable, interval } from 'rxjs';
import { share, take } from 'rxjs/operators';

// COLD OBSERVABLE - każdy subscriber dostaje własny strumień
const coldObservable = new Observable(subscriber => {
  console.log('Tworzę nowego producenta!');
  const random = Math.random();
  subscriber.next(random);
  subscriber.complete();
});

console.log('=== COLD OBSERVABLE ===');
coldObservable.subscribe(x => console.log('Subscriber 1:', x));
coldObservable.subscribe(x => console.log('Subscriber 2:', x));
// Wyświetli dwie różne losowe liczby - każdy ma własny strumień

// HOT OBSERVABLE - wszyscy współdzielą ten sam strumień
const hotObservable = coldObservable.pipe(share());

console.log('=== HOT OBSERVABLE ===');
hotObservable.subscribe(x => console.log('Subscriber 1:', x));
hotObservable.subscribe(x => console.log('Subscriber 2:', x));
// Wyświetli tę samą losową liczbę - współdzielony strumień

// Praktyczny przykład - HTTP request
import { ajax } from 'rxjs/ajax';
import { shareReplay } from 'rxjs/operators';

// Cold - każdy subscribe wywoła nowe HTTP request
const coldHttp = ajax.getJSON('https://api.example.com/users');

coldHttp.subscribe(users => console.log('Request 1:', users.length));
coldHttp.subscribe(users => console.log('Request 2:', users.length));
// Wysłane 2 oddzielne requesty!

// Hot z cache - jeden request, współdzielony wynik
const hotHttp = ajax.getJSON('https://api.example.com/users').pipe(
  shareReplay(1) // Cache ostatniej wartości dla nowych subscribers
);

hotHttp.subscribe(users => console.log('Request 1:', users.length));
hotHttp.subscribe(users => console.log('Request 2:', users.length));
// Tylko 1 request! Drugi subscriber dostaje cache'owaną wartość

// Subject - naturalnie hot
import { Subject } from 'rxjs';

const subject = new Subject();

subject.subscribe(x => console.log('Sub 1:', x));
subject.next(1); // Oba dostaną 1

subject.subscribe(x => console.log('Sub 2:', x));
subject.next(2); // Oba dostaną 2
subject.next(3); // Oba dostaną 3
graph TD
    subgraph "Cold Observable - Unicast"
        A[Observable] --> B[Subscribe 1]
        A --> C[Subscribe 2]
        B --> D[Wykonanie 1: values 1,2,3]
        C --> E[Wykonanie 2: values 1,2,3]
    end

    subgraph "Hot Observable - Multicast"
        F[Observable] --> G[Wspólne wykonanie]
        G --> H[Subscribe 1]
        G --> I[Subscribe 2]
        H --> J[Dostaje: values 1,2,3]
        I --> K[Dostaje: values 2,3 - po subskrypcji]
    end

Materiały

Operatory Tworzące w RxJS

6. Jak używać operatora of() do tworzenia Observable z wartości?

Odpowiedź w 30 sekund: Operator of() tworzy Observable, który emituje podane wartości po kolei, a następnie od razu się kończy (complete). Jest idealny do tworzenia prostych strumieni z wcześniej znanych wartości lub do testowania.

Odpowiedź w 2 minuty: Operator of() to jeden z najprostszych operatorów tworzących w RxJS. Przyjmuje dowolną liczbę argumentów i tworzy Observable, który emituje każdą z tych wartości synchronicznie, jedna po drugiej, a następnie wywołuje complete(). Jest to bardzo przydatne gdy chcemy opakować znane wartości w Observable, np. do testowania, zwracania stałych wartości z funkcji, lub łączenia z innymi operatorami.

W praktyce of() często używamy gdy musimy zwrócić Observable z funkcji, ale wartość jest już znana i nie wymaga asynchronicznego pobierania. Możemy przekazać wartości różnych typów - liczby, stringi, obiekty, tablice, funkcje. Każdy argument będzie osobną emisją. Jeśli chcemy wyemitować tablicę jako pojedynczą wartość, przekazujemy ją jako jeden argument.

Różnica między of([1, 2, 3]) a of(1, 2, 3) jest kluczowa: pierwsza wersja wyemituje jedną wartość (tablicę), druga wyemituje trzy oddzielne wartości. Operator of() jest synchroniczny - wszystkie emisje następują natychmiast, w tym samym cyklu zdarzeń.

Częste zastosowania to mockowanie danych w testach, tworzenie prostych strumieni do demonstracji działania operatorów, oraz sytuacje gdy chcemy ujednolicić API funkcji zwracających zarówno synchroniczne jak i asynchroniczne wartości.

Przykład kodu:

import { of } from 'rxjs';

// Podstawowe użycie - emituje 3 wartości i kończy
of(1, 2, 3).subscribe({
  next: (value) => console.log('Wartość:', value),
  complete: () => console.log('Zakończone!')
});
// Output: Wartość: 1, Wartość: 2, Wartość: 3, Zakończone!

// Emitowanie różnych typów danych
of('tekst', 42, { name: 'Jan' }, [1, 2, 3]).subscribe(
  value => console.log(typeof value, value)
);

// Emitowanie tablicy jako pojedynczej wartości
of([1, 2, 3]).subscribe(arr => {
  console.log('Tablica:', arr); // Output: Tablica: [1, 2, 3]
});

// Praktyczne użycie - mockowanie API
function getUserById(id: number): Observable<User> {
  if (id === 1) {
    // Zwracamy mockowe dane
    return of({ id: 1, name: 'Jan Kowalski' });
  }
  // W rzeczywistości tutaj byłoby HTTP request
  return ajax.getJSON(`/api/users/${id}`);
}

// Użycie w pipe do transformacji
of(5, 10, 15)
  .pipe(
    map(x => x * 2)
  )
  .subscribe(x => console.log(x)); // Output: 10, 20, 30

Materiały

7. Jak działa operator from() i jakie typy danych może konwertować?

Odpowiedź w 30 sekund: Operator from() konwertuje różne typy danych do Observable: tablice, Promise, iteratory, Observable-like obiekty. W przeciwieństwie do of(), który traktuje tablicę jako jedną wartość, from() emituje każdy element tablicy osobno.

Odpowiedź w 2 minuty: Operator from() jest bardzo uniwersalnym operatorem tworzącym, który potrafi przekonwertować wiele różnych typów danych na Observable. Najczęściej używamy go do konwersji tablic (emitując każdy element osobno), Promise (konwertując do Observable, który emituje wynik Promise), oraz obiektów iterowalnych jak Map, Set, String czy generatory.

Gdy przekażemy tablicę do from([1, 2, 3]), otrzymamy Observable emitujący trzy osobne wartości: 1, 2, 3, a następnie complete. To fundamentalna różnica w porównaniu do of([1, 2, 3]), który wyemitowałby całą tablicę jako jedną wartość. Dzięki temu from() świetnie nadaje się do przetwarzania kolekcji element po elemencie.

Promise są kolejnym częstym przypadkiem użycia. from(fetch('/api/data')) przekonwertuje Promise zwracane przez fetch na Observable. Gdy Promise się rozwiąże, Observable wyemituje wartość i zakończy się. Jeśli Promise zostanie odrzucone, Observable wyemituje błąd. To pozwala na łączenie kodu opartego na Promise z operatorami RxJS.

Operator from() obsługuje również obiekty Observable-like (z metodą subscribe), iteratory (obiekty z metodą Symbol.iterator), oraz inne struktury danych. Może obsłużyć string (emitując każdy znak osobno) oraz async iteratory. Ta uniwersalność czyni from() kluczowym narzędziem do integracji RxJS z różnymi źródłami danych w aplikacji.

Przykład kodu:

import { from } from 'rxjs';
import { map, take } from 'rxjs/operators';

// Konwersja tablicy - każdy element osobno
from([10, 20, 30, 40]).subscribe(
  value => console.log('Element:', value)
);
// Output: Element: 10, Element: 20, Element: 30, Element: 40

// Konwersja Promise
const promise = new Promise(resolve => {
  setTimeout(() => resolve('Dane z Promise'), 1000);
});

from(promise).subscribe(
  value => console.log('Promise zwrócił:', value)
);
// Output (po 1s): Promise zwrócił: Dane z Promise

// Konwersja Map
const mapa = new Map([
  ['klucz1', 'wartość1'],
  ['klucz2', 'wartość2']
]);

from(mapa).subscribe(
  ([key, value]) => console.log(`${key}: ${value}`)
);
// Output: klucz1: wartość1, klucz2: wartość2

// Konwersja Set
from(new Set([1, 2, 2, 3, 3, 3])).subscribe(
  value => console.log('Unikalna wartość:', value)
);
// Output: Unikalna wartość: 1, 2, 3

// Konwersja stringa
from('RxJS').subscribe(
  char => console.log('Znak:', char)
);
// Output: Znak: R, Znak: x, Znak: J, Znak: S

// Generator function
function* numberGenerator() {
  yield 1;
  yield 2;
  yield 3;
}

from(numberGenerator()).subscribe(
  num => console.log('Z generatora:', num)
);

// Praktyczne użycie - fetch API z RxJS
from(fetch('https://api.example.com/data'))
  .pipe(
    switchMap(response => from(response.json())),
    map(data => data.items)
  )
  .subscribe(items => console.log('Pobrane dane:', items));

// Różnica między from() i of()
console.log('--- Używając from() ---');
from([1, 2, 3]).subscribe(x => console.log(x)); // 1, 2, 3

console.log('--- Używając of() ---');
of([1, 2, 3]).subscribe(x => console.log(x)); // [1, 2, 3]

Materiały

8. Czym jest operator fromEvent() i jak obsługiwać zdarzenia DOM?

Odpowiedź w 30 sekund: Operator fromEvent() konwertuje zdarzenia DOM (i inne event emitters) na Observable. Automatycznie dodaje event listener, emituje każde zdarzenie, i usuwa listener przy unsubscribe. Idealny do obsługi kliknięć, ruchów myszy, czy inputów użytkownika.

Odpowiedź w 2 minuty: fromEvent() to kluczowy operator do pracy ze zdarzeniami w przeglądarce. Pozwala on przekształcić dowolne zdarzenie DOM w strumień Observable, co umożliwia zastosowanie całego zestawu operatorów RxJS do obsługi interakcji użytkownika. Operator przyjmuje dwa główne parametry: target (element DOM, window, document, lub dowolny EventTarget) oraz nazwę zdarzenia (np. 'click', 'mousemove', 'keyup').

Wielką zaletą fromEvent() jest automatyczne zarządzanie cyklem życia listenera. Kiedy tworzysz Observable za pomocą fromEvent(), listener jest dodawany dopiero gdy ktoś zasubskrybuje ten Observable. Co ważniejsze, gdy wywołasz unsubscribe(), listener zostanie automatycznie usunięty, co zapobiega wyciekom pamięci - częstemu problemowi przy ręcznym zarządzaniu event listenerami.

Observable stworzony przez fromEvent() nigdy się nie kończy sam - będzie emitował zdarzenia dopóki nie wywołasz unsubscribe lub dopóki element nie zostanie usunięty z DOM. Każda emisja to obiekt Event z przeglądarki, zawierający wszystkie informacje o zdarzeniu (target, timestamp, dane specyficzne dla typu zdarzenia).

W praktyce fromEvent() świetnie współgra z operatorami RxJS. Możesz użyć debounceTime() do ograniczenia częstotliwości zdarzeń (np. przy input), throttleTime() do kontrolowania rate'u (np. przy scroll), map() do wyciągnięcia konkretnych danych ze zdarzenia, czy filter() do warunkowej obsługi. To czyni kod obsługi zdarzeń znacznie bardziej czytelnym i deklaratywnym niż tradycyjne podejście z callback'ami.

Przykład kodu:

import { fromEvent } from 'rxjs';
import { map, debounceTime, throttleTime, filter, tap } from 'rxjs/operators';

// Podstawowa obsługa kliknięcia
const button = document.querySelector('#myButton');
const clicks$ = fromEvent(button, 'click');

clicks$.subscribe(event => {
  console.log('Przycisk kliknięty!', event);
});

// Obsługa inputu z debounce (czeka aż użytkownik przestanie pisać)
const searchInput = document.querySelector('#searchInput');
const search$ = fromEvent(searchInput, 'input').pipe(
  debounceTime(500), // Czeka 500ms po ostatnim znaku
  map((event: Event) => (event.target as HTMLInputElement).value),
  filter(text => text.length >= 3), // Minimum 3 znaki
  tap(text => console.log('Szukam:', text))
);

search$.subscribe(searchTerm => {
  // Wykonaj wyszukiwanie
  performSearch(searchTerm);
});

// Obsługa ruchu myszy z throttle (limituje częstotliwość)
const mousemove$ = fromEvent(document, 'mousemove').pipe(
  throttleTime(100), // Maksymalnie co 100ms
  map((event: MouseEvent) => ({ x: event.clientX, y: event.clientY }))
);

mousemove$.subscribe(position => {
  console.log(`Pozycja myszy: X=${position.x}, Y=${position.y}`);
});

// Obsługa scroll z pozycją
const scroll$ = fromEvent(window, 'scroll').pipe(
  throttleTime(200),
  map(() => window.scrollY)
);

scroll$.subscribe(scrollPosition => {
  if (scrollPosition > 300) {
    // Pokaż przycisk "scroll to top"
    showScrollTopButton();
  }
});

// Obsługa klawiszy - tylko Enter
const input = document.querySelector('#messageInput');
const enterKey$ = fromEvent(input, 'keyup').pipe(
  filter((event: KeyboardEvent) => event.key === 'Enter'),
  map((event: Event) => (event.target as HTMLInputElement).value)
);

enterKey$.subscribe(message => {
  console.log('Wysłano wiadomość:', message);
  sendMessage(message);
});

// Łączenie wielu zdarzeń - drag and drop
const element = document.querySelector('#draggable');
const mousedown$ = fromEvent(element, 'mousedown');
const mousemove$ = fromEvent(document, 'mousemove');
const mouseup$ = fromEvent(document, 'mouseup');

const drag$ = mousedown$.pipe(
  switchMap(() => mousemove$.pipe(
    takeUntil(mouseup$),
    map((event: MouseEvent) => ({
      x: event.clientX,
      y: event.clientY
    }))
  ))
);

drag$.subscribe(position => {
  // Przesuń element
  element.style.left = position.x + 'px';
  element.style.top = position.y + 'px';
});

// Unsubscribe - ważne przy niszczeniu komponentów
const subscription = clicks$.subscribe(/* ... */);

// Później, np. w ngOnDestroy() w Angular
subscription.unsubscribe(); // Usuwa event listener

Diagram przepływu zdarzeń:

graph LR
    A[DOM Event] --> B[fromEvent]
    B --> C[Observable Stream]
    C --> D[debounceTime/throttleTime]
    D --> E[map/filter]
    E --> F[subscribe]
    F --> G[Handler Function]

    style A fill:#e1f5ff
    style C fill:#fff4e1
    style F fill:#e7ffe1

Materiały

9. Jak używać operatora interval() i timer() do tworzenia strumieni czasowych?

Odpowiedź w 30 sekund: interval(n) emituje liczby (0, 1, 2...) co n milisekund w nieskończoność. timer(delay, period) czeka delay milisekund, emituje 0, a potem opcjonalnie emituje kolejne wartości co period milisekund. Oba przydatne do polling, animacji i operacji czasowych.

Odpowiedź w 2 minuty: Operatory interval() i timer() służą do tworzenia Observable bazujących na czasie, co jest niezwykle przydatne w wielu scenariuszach aplikacji. interval(n) to prostszy z nich - tworzy nieskończony strumień emitujący kolejne liczby całkowite (zaczynając od 0) w regularnych odstępach czasu. Na przykład interval(1000) będzie emitował 0, 1, 2, 3... co sekundę, aż do momentu unsubscribe.

timer() jest bardziej elastyczny i ma dwa tryby działania. Jako timer(delay) działa jak setTimeout - czeka określony czas i emituje pojedynczą wartość 0, po czym się kończy. Jako timer(delay, period) działa jak kombinacja setTimeout i setInterval - czeka delay milisekund, emituje 0, a potem emituje kolejne wartości (1, 2, 3...) co period milisekund. Możesz też użyć timer(0, 1000), co zadziała podobnie do interval(1000), ale pierwsza emisja nastąpi natychmiast.

Oba operatory świetnie współpracują z innymi operatorami RxJS. Możesz użyć take(n) aby ograniczyć liczbę emisji, takeUntil() aby zakończyć strumień na podstawie warunku, czy switchMap() aby uruchomić inną operację przy każdym ticku. Częste zastosowania to polling API (sprawdzanie nowych danych co X sekund), animacje (aktualizacja interfejsu co frame), liczniki czasu, timeout'y, oraz operacje wymagające opóźnienia.

Ważne jest aby pamiętać o unsubscribe przy tych operatorach, ponieważ będą działać w nieskończoność i mogą powodować wycieki pamięci. W frameworkach jak Angular, najlepiej używać ich z operatorem takeUntil() podpiętym do lifecycle hook'a komponentu (np. ngOnDestroy), lub korzystać z async pipe, który automatycznie zarządza subskrypcją.

Przykład kodu:

import { interval, timer } from 'rxjs';
import { take, takeUntil, map, switchMap, tap } from 'rxjs/operators';

// INTERVAL - podstawowe użycie
console.log('Start interval');
const interval$ = interval(1000); // Co 1 sekundę

interval$.pipe(take(5)).subscribe(
  value => console.log('Interval emitował:', value)
);
// Output: 0, 1, 2, 3, 4 (co sekundę)

// TIMER - pojedyncza emisja (jak setTimeout)
console.log('Start timer - pojedyncza emisja');
timer(3000).subscribe(
  () => console.log('Timer zakończony po 3 sekundach!')
);

// TIMER - z okresowymi emisjami (jak setInterval)
console.log('Start timer - okresowe emisje');
timer(2000, 1000).pipe(take(5)).subscribe(
  value => console.log('Timer emitował:', value)
);
// Czeka 2s, potem emituje: 0, 1, 2, 3, 4 (co 1s)

// TIMER vs INTERVAL - różnica w pierwszej emisji
timer(0, 1000).pipe(take(3)).subscribe(
  v => console.log('Timer z delay 0:', v)
);
// Output natychmiast: 0, potem 1, 2 (co 1s)

interval(1000).pipe(take(3)).subscribe(
  v => console.log('Interval:', v)
);
// Output po 1s: 0, potem 1, 2 (co 1s)

// Praktyczne użycie 1: Polling API
const pollApi$ = interval(5000).pipe(
  switchMap(() => fetch('/api/status').then(r => r.json())),
  tap(data => console.log('Pobrano nowe dane:', data))
);

const pollSubscription = pollApi$.subscribe();

// Zatrzymaj polling po 30 sekundach
timer(30000).subscribe(() => {
  pollSubscription.unsubscribe();
  console.log('Polling zatrzymany');
});

// Praktyczne użycie 2: Licznik odliczający
const countdown$ = timer(0, 1000).pipe(
  map(n => 10 - n),
  take(11)
);

countdown$.subscribe(
  seconds => console.log(`Pozostało: ${seconds}s`),
  null,
  () => console.log('Odliczanie zakończone!')
);

// Praktyczne użycie 3: Automatyczne odświeżanie tokenu
const tokenRefresh$ = timer(3600000, 3600000).pipe(
  switchMap(() => refreshAuthToken())
);

tokenRefresh$.subscribe(
  token => console.log('Token odświeżony:', token)
);

// Praktyczne użycie 4: Animacja z interwałem
const animation$ = interval(16).pipe( // ~60 FPS
  take(100),
  map(frame => frame / 100) // 0.0 do 1.0
);

animation$.subscribe(progress => {
  const element = document.querySelector('#animated');
  element.style.opacity = progress.toString();
});

// Praktyczne użycie 5: Retry z opóźnieniem
function fetchWithRetry(url: string) {
  return fetch(url).pipe(
    catchError(error => {
      console.log('Błąd, retry za 3s...');
      return timer(3000).pipe(
        switchMap(() => fetch(url))
      );
    })
  );
}

// Użycie z takeUntil (Angular pattern)
import { Subject } from 'rxjs';

class MyComponent {
  private destroy$ = new Subject<void>();

  ngOnInit() {
    // Polling zatrzyma się automatycznie przy destroy
    interval(1000).pipe(
      takeUntil(this.destroy$),
      tap(n => console.log('Tick:', n))
    ).subscribe();
  }

  ngOnDestroy() {
    this.destroy$.next();
    this.destroy$.complete();
  }
}

Diagram czasowy:

gantt
    title Porównanie interval() i timer()
    dateFormat X
    axisFormat %L ms

    section interval(1000)
    0: milestone, 1000, 0
    1: milestone, 2000, 0
    2: milestone, 3000, 0
    3: milestone, 4000, 0

    section timer(2000, 1000)
    wait: 0, 2000
    0: milestone, 2000, 0
    1: milestone, 3000, 0
    2: milestone, 4000, 0

    section timer(3000)
    wait: 0, 3000
    0: milestone, 3000, 0

Materiały

10. Jak działa operator ajax() do wykonywania żądań HTTP?

Odpowiedź w 30 sekund: Operator ajax() z rxjs/ajax tworzy Observable wykonujący żądania HTTP (alternatywa dla fetch/XMLHttpRequest). Obsługuje GET, POST, PUT, DELETE, automatycznie parsuje JSON, pozwala na konfigurację headers, timeout, i świetnie współpracuje z operatorami RxJS jak retry, catchError, czy switchMap.

Odpowiedź w 2 minuty: ajax() to dedykowany operator RxJS do wykonywania żądań HTTP, będący wrapperem wokół XMLHttpRequest. Choć jest starszy niż Fetch API, oferuje doskonałą integrację z ekosystemem RxJS i kilka unikalnych zalet. Możesz użyć prostej formy ajax('/api/users') dla GET, lub pełnej konfiguracji z obiektem ustawień zawierającym url, method, headers, body, responseType i wiele więcej.

Największą zaletą ajax() w porównaniu do from(fetch()) jest natywna obsługa przez RxJS. Observable zwrócony przez ajax() można łatwo anulować przez unsubscribe (co przerywa żądanie HTTP), podczas gdy Promise z fetch nie da się anulować. ajax() zwraca obiekt AjaxResponse zawierający response, status, responseType, oraz oryginalne XMLHttpRequest, co daje pełną kontrolę nad odpowiedzią.

Operator oferuje pomocnicze metody: ajax.get(), ajax.post(), ajax.put(), ajax.delete(), ajax.patch() oraz ajax.getJSON() która automatycznie parsuje odpowiedź jako JSON i zwraca bezpośrednio dane (nie cały obiekt response). To bardzo wygodne dla typowych przypadków użycia API RESTful.

W praktyce ajax() świetnie współpracuje z operatorami obsługi błędów (catchError, retry, retryWhen), operatorami transformacji (map, pluck, switchMap), oraz zaawansowanymi wzorcami jak automatic retry z exponential backoff, caching, concurrent request limiting przez mergeMap(fn, concurrent), czy request deduplication przez shareReplay(). Jest to szczególnie wartościowe w złożonych aplikacjach gdzie potrzebujesz precyzyjnej kontroli nad przepływem żądań HTTP i obsługą błędów.

Przykład kodu:

import { ajax } from 'rxjs/ajax';
import { map, catchError, retry, switchMap, debounceTime } from 'rxjs/operators';
import { of } from 'rxjs';

// Proste GET request
ajax('/api/users').subscribe(
  response => console.log('Odpowiedź:', response),
  error => console.error('Błąd:', error)
);

// GET z automatycznym parsowaniem JSON
ajax.getJSON<User[]>('/api/users').subscribe(
  users => console.log('Użytkownicy:', users)
);

// POST request z danymi
const newUser = {
  name: 'Jan Kowalski',
  email: 'jan@example.com'
};

ajax.post('/api/users', newUser, {
  'Content-Type': 'application/json'
}).subscribe(
  response => console.log('Utworzono użytkownika:', response)
);

// Pełna konfiguracja żądania
ajax({
  url: '/api/data',
  method: 'POST',
  headers: {
    'Content-Type': 'application/json',
    'Authorization': 'Bearer token123'
  },
  body: {
    query: 'search term'
  },
  timeout: 5000 // Timeout po 5 sekundach
}).subscribe(
  response => console.log('Dane:', response.response),
  error => console.error('Błąd HTTP:', error)
);

// Obsługa błędów z fallback
ajax.getJSON('/api/users').pipe(
  catchError(error => {
    console.error('Błąd pobierania, używam cache:', error);
    return of(getCachedUsers()); // Zwróć dane z cache
  })
).subscribe(users => displayUsers(users));

// Automatyczne retry przy błędzie
ajax.getJSON('/api/unstable-endpoint').pipe(
  retry(3), // Spróbuj 3 razy przed błędem
  catchError(error => {
    console.error('Nie udało się po 3 próbach:', error);
    return of(null);
  })
).subscribe(data => console.log('Dane:', data));

// Retry z exponential backoff
import { retryWhen, delayWhen, tap, take } from 'rxjs/operators';
import { timer } from 'rxjs';

ajax.getJSON('/api/data').pipe(
  retryWhen(errors => errors.pipe(
    tap(err => console.log('Retry po błędzie:', err)),
    delayWhen((_, i) => timer(Math.pow(2, i) * 1000)), // 1s, 2s, 4s, 8s...
    take(4) // Maksymalnie 4 retry
  ))
).subscribe(
  data => console.log('Sukces:', data),
  error => console.error('Ostateczny błąd:', error)
);

// Praktyczne użycie: Autocomplete z debounce
const searchInput$ = fromEvent(input, 'input').pipe(
  debounceTime(300),
  map(event => (event.target as HTMLInputElement).value),
  filter(text => text.length >= 2),
  switchMap(searchTerm =>
    ajax.getJSON(`/api/search?q=${searchTerm}`).pipe(
      catchError(() => of([]))
    )
  )
);

searchInput$.subscribe(results => {
  displaySearchResults(results);
});

// Anulowanie żądania przez unsubscribe
const subscription = ajax.getJSON('/api/large-data').subscribe(
  data => console.log('Dane:', data)
);

// Po 1 sekundzie anuluj żądanie jeśli nadal trwa
setTimeout(() => {
  subscription.unsubscribe();
  console.log('Żądanie anulowane');
}, 1000);

// Progress monitoring (tylko XMLHttpRequest)
ajax({
  url: '/api/upload',
  method: 'POST',
  body: formData,
  progressSubscriber: {
    next: e => {
      if (e.type === 'upload_progress') {
        const percentComplete = (e.loaded / e.total) * 100;
        console.log(`Upload: ${percentComplete}%`);
      }
    }
  }
}).subscribe(
  response => console.log('Upload zakończony:', response)
);

// Łączenie wielu żądań
import { forkJoin } from 'rxjs';

forkJoin({
  users: ajax.getJSON('/api/users'),
  posts: ajax.getJSON('/api/posts'),
  comments: ajax.getJSON('/api/comments')
}).subscribe(({ users, posts, comments }) => {
  console.log('Wszystkie dane pobrane:', { users, posts, comments });
});

// Sekwencyjne żądania (jedno po drugim)
ajax.getJSON<User>('/api/user/1').pipe(
  switchMap(user =>
    ajax.getJSON(`/api/posts?userId=${user.id}`).pipe(
      map(posts => ({ user, posts }))
    )
  )
).subscribe(({ user, posts }) => {
  console.log(`Posty użytkownika ${user.name}:`, posts);
});

// CRUD operations helper
class UserService {
  private baseUrl = '/api/users';

  getAll() {
    return ajax.getJSON<User[]>(this.baseUrl);
  }

  getById(id: number) {
    return ajax.getJSON<User>(`${this.baseUrl}/${id}`);
  }

  create(user: Partial<User>) {
    return ajax.post(this.baseUrl, user, {
      'Content-Type': 'application/json'
    }).pipe(map(response => response.response));
  }

  update(id: number, user: Partial<User>) {
    return ajax.put(`${this.baseUrl}/${id}`, user, {
      'Content-Type': 'application/json'
    }).pipe(map(response => response.response));
  }

  delete(id: number) {
    return ajax.delete(`${this.baseUrl}/${id}`);
  }
}

// Użycie serwisu
const userService = new UserService();

userService.getAll().subscribe(
  users => console.log('Wszyscy użytkownicy:', users)
);

userService.create({ name: 'Anna', email: 'anna@example.com' }).subscribe(
  newUser => console.log('Utworzono:', newUser)
);

Porównanie ajax() vs fetch():

// Fetch API (Promise-based)
from(fetch('/api/data')
  .then(response => response.json()))
  .subscribe(data => console.log(data));

// RxJS ajax (Observable-based)
ajax.getJSON('/api/data')
  .subscribe(data => console.log(data));

// Zalety ajax():
// ✓ Można anulować przez unsubscribe
// ✓ Natywna integracja z operatorami RxJS
// ✓ Automatyczne parsowanie JSON przez getJSON()
// ✓ Progress events
// ✓ Request/response interceptors łatwiejsze do implementacji

// Zalety fetch():
// ✓ Nowoczesne API
// ✓ Szersze wsparcie przeglądarek (standardowe API)
// ✓ Lepsza obsługa CORS
// ✓ Service Worker support

Materiały

Operatory Transformacji RxJS

Jak działa operator map() i kiedy go używać?

Odpowiedź w 30 sekund: Operator map() transformuje każdą wartość emitowaną przez Observable, stosując do niej podaną funkcję. Jest to odpowiednik metody Array.map() w RxJS - pobiera wartość, przekształca ją i emituje wynik.

Odpowiedź w 2 minuty: Operator map() jest jednym z najbardziej podstawowych i najczęściej używanych operatorów transformacji w RxJS. Działa synchronicznie - dla każdej wartości emitowanej przez źródłowy Observable, map() natychmiast aplikuje podaną funkcję transformacji i emituje wynik.

Używamy map() gdy potrzebujemy przekształcić dane "jeden do jednego" - każda wartość wejściowa produkuje dokładnie jedną wartość wyjściową. Typowe przypadki użycia to: wyciąganie konkretnych pól z obiektów, konwersja typów danych, formatowanie wartości, wykonywanie obliczeń na wartościach czy mapowanie odpowiedzi HTTP do modeli domenowych.

Kluczowa różnica w porównaniu do operatorów wyższego rzędu (jak mergeMap, switchMap) polega na tym, że funkcja przekazana do map() musi zwracać zwykłą wartość, a nie Observable. Jeśli zwrócimy Observable, otrzymamy Observable zagnieżdżony wewnątrz Observable, co zazwyczaj nie jest tym czego chcemy.

Operator map() zachowuje strukturę strumienia - nie wpływa na timing emisji, nie dodaje ani nie usuwa wartości, jedynie je przekształca. Jest całkowicie synchroniczny i deterministyczny, co czyni go bezpiecznym i przewidywalnym w użyciu.

Przykład kodu:

import { of, fromEvent } from 'rxjs';
import { map } from 'rxjs/operators';

// Przykład 1: Transformacja prostych wartości
const numbers$ = of(1, 2, 3, 4, 5);
const doubled$ = numbers$.pipe(
  map(x => x * 2)
);
doubled$.subscribe(val => console.log(val)); // 2, 4, 6, 8, 10

// Przykład 2: Wyciąganie pól z obiektów
interface User {
  id: number;
  name: string;
  email: string;
}

const users$ = of<User>(
  { id: 1, name: 'Anna', email: 'anna@example.com' },
  { id: 2, name: 'Jan', email: 'jan@example.com' }
);

const userNames$ = users$.pipe(
  map(user => user.name)
);
userNames$.subscribe(name => console.log(name)); // 'Anna', 'Jan'

// Przykład 3: Transformacja zdarzeń DOM
const clicks$ = fromEvent<MouseEvent>(document, 'click');
const clickPositions$ = clicks$.pipe(
  map(event => ({ x: event.clientX, y: event.clientY }))
);
clickPositions$.subscribe(pos => console.log(`Kliknięto w: ${pos.x}, ${pos.y}`));

// Przykład 4: Mapowanie odpowiedzi HTTP
interface ApiResponse {
  data: { users: User[] };
  status: number;
}

const apiResponse$ = of<ApiResponse>({
  data: { users: [{ id: 1, name: 'Anna', email: 'anna@example.com' }] },
  status: 200
});

const extractedUsers$ = apiResponse$.pipe(
  map(response => response.data.users)
);

Materiały

Jaka jest różnica między map() a mergeMap() (flatMap)?

Odpowiedź w 30 sekund: map() transformuje wartości synchronicznie (wartość → wartość), podczas gdy mergeMap() obsługuje operacje asynchroniczne (wartość → Observable → wartość). mergeMap() automatycznie subskrybuje wewnętrzne Observable'e i spłaszcza wyniki do jednego strumienia, zarządzając wieloma równoczesnymi subskrypcjami.

Odpowiedź w 2 minuty: Główna różnica między map() a mergeMap() polega na tym, jak obsługują wartości zwracane przez funkcję transformacji. Operator map() oczekuje zwykłej wartości i emituje ją bezpośrednio. Jeśli zwrócimy Observable z funkcji przekazanej do map(), otrzymamy Observable zagnieżdżony w Observable (Observable<Observable>), co rzadko jest pożądane.

mergeMap() (alias flatMap()) został zaprojektowany specjalnie do obsługi operacji asynchronicznych. Gdy funkcja transformacji zwraca Observable, mergeMap() automatycznie subskrybuje ten wewnętrzny Observable, czeka na jego wartości i emituje je w głównym strumieniu. Proces ten nazywamy "spłaszczaniem" (flattening).

Kluczowa cecha mergeMap() to równoczesność - nie czeka na zakończenie poprzedniego wewnętrznego Observable przed subskrybowaniem kolejnego. Jeśli źródłowy Observable emituje wartości szybko, mergeMap() może zarządzać wieloma aktywnymi subskrypcjami jednocześnie. Można to kontrolować opcjonalnym parametrem concurrent określającym maksymalną liczbę równoczesnych subskrypcji.

Typowe przypadki użycia mergeMap() to: wywołania HTTP API (gdzie każda wartość źródłowa inicjuje request), operacje na bazie danych, operacje plikowe, czy jakiekolwiek inne asynchroniczne operacje gdzie kolejność nie ma znaczenia i chcemy maksymalnej wydajności poprzez przetwarzanie równoległe.

Przykład kodu:

import { of, interval, fromEvent } from 'rxjs';
import { map, mergeMap, take, delay } from 'rxjs/operators';

// PROBLEM: Używanie map() z Observable (NIEPRAWIDŁOWE)
const numbers$ = of(1, 2, 3);
const wrongWay$ = numbers$.pipe(
  map(n => of(n * 2)) // Zwraca Observable<Observable<number>>
);
wrongWay$.subscribe(obs => {
  console.log(obs); // Wypisze Observable, a nie wartość!
});

// ROZWIĄZANIE: Używanie mergeMap()
const rightWay$ = numbers$.pipe(
  mergeMap(n => of(n * 2))
);
rightWay$.subscribe(val => console.log(val)); // 2, 4, 6

// Przykład 2: Symulacja wywołań API
interface Product { id: number; name: string; }

function fetchProductDetails(id: number): Observable<Product> {
  // Symulacja opóźnienia API
  return of({ id, name: `Produkt ${id}` }).pipe(
    delay(Math.random() * 1000)
  );
}

const productIds$ = of(1, 2, 3, 4, 5);

// Z mergeMap - wszystkie requesty równocześnie
productIds$.pipe(
  mergeMap(id => fetchProductDetails(id))
).subscribe(product => console.log('Otrzymano:', product.name));
// Produkty mogą przyjść w dowolnej kolejności!

// Przykład 3: Obsługa kliknięć z wywołaniem API
const searchButton = document.getElementById('search-button');
const clicks$ = fromEvent(searchButton, 'click');

clicks$.pipe(
  mergeMap(() => {
    // Każde kliknięcie inicjuje nowe zapytanie API
    return fetch('/api/search').then(r => r.json());
  })
).subscribe(results => {
  console.log('Wyniki wyszukiwania:', results);
});

// Przykład 4: Kontrola równoczesności
productIds$.pipe(
  mergeMap(
    id => fetchProductDetails(id),
    3 // Maksymalnie 3 równoczesne requesty
  )
).subscribe(product => console.log(product));

// Przykład 5: Wizualizacja różnicy
console.log('=== map() z Observable ===');
of(1, 2).pipe(
  map(x => of(x * 10))
).subscribe(val => console.log(typeof val)); // 'object' (Observable)

console.log('=== mergeMap() z Observable ===');
of(1, 2).pipe(
  mergeMap(x => of(x * 10))
).subscribe(val => console.log(typeof val)); // 'number' (spłaszczona wartość)

Diagram:

graph TB
    subgraph "map() - Transformacja synchroniczna"
        A1[1] --> B1[map x => x*2]
        A2[2] --> B1
        A3[3] --> B1
        B1 --> C1[2]
        B1 --> C2[4]
        B1 --> C3[6]
    end

    subgraph "mergeMap() - Transformacja asynchroniczna"
        D1[1] --> E1[mergeMap x => API call]
        D2[2] --> E1
        D3[3] --> E1
        E1 --> F1[Observable 1]
        E1 --> F2[Observable 2]
        E1 --> F3[Observable 3]
        F1 -.-> G1[wynik]
        F2 -.-> G2[wynik]
        F3 -.-> G3[wynik]
        style F1 fill:#ffcccc
        style F2 fill:#ffcccc
        style F3 fill:#ffcccc
    end

Materiały

Kiedy używać switchMap() zamiast mergeMap()?

Odpowiedź w 30 sekund: Używaj switchMap() gdy chcesz anulować poprzednie operacje po nadejściu nowej wartości (np. autouzupełnianie, wyszukiwanie na żywo). mergeMap() nadaje się gdy wszystkie operacje powinny się zakończyć niezależnie od nowych wartości (np. zapisywanie logów, analytics). switchMap() zawsze emituje tylko z najnowszego wewnętrznego Observable.

Odpowiedź w 2 minuty: switchMap() jest wariantem operatorów spłaszczających z kluczową cechą: gdy źródłowy Observable emituje nową wartość, switchMap() automatycznie anuluje (unsubscribe) poprzedni wewnętrzny Observable i subskrybuje nowy. To zachowanie typu "przełączanie się" (switching) zapewnia, że w danym momencie tylko jedna wewnętrzna subskrypcja jest aktywna.

Najbardziej klasycznym przypadkiem użycia switchMap() jest implementacja wyszukiwania na żywo (type-ahead search). Gdy użytkownik wpisuje tekst, każda litera może potencjalnie wywołać zapytanie do API. Jeśli użytkownik wpisze "program" szybko, nie chcemy wysyłać 7 osobnych requestów i czekać na wszystkie odpowiedzi - interesuje nas tylko wynik dla ostatniego, kompletnego terminu. switchMap() automatycznie anuluje przestarzałe requesty.

Z kolei mergeMap() pozwala wszystkim wewnętrznym Observable'om działać równocześnie i zakończyć się niezależnie. Jest odpowiedni gdy każda operacja jest ważna i powinna się zakończyć, np. zapisywanie danych do bazy, wysyłanie eventów analytics, czy przetwarzanie kolejki zadań.

Kluczowa różnica w praktyce: switchMap() zapobiega problemowi "race condition" gdzie późniejszy request może zakończyć się przed wcześniejszym, nadpisując wyniki w niewłaściwej kolejności. Gwarantuje też, że interfejs użytkownika zawsze pokazuje dane odpowiadające najnowszemu stanowi, nie przestarzałe wyniki z poprzednich operacji.

Przykład kodu:

import { fromEvent, of, interval } from 'rxjs';
import {
  switchMap,
  mergeMap,
  map,
  debounceTime,
  distinctUntilChanged,
  delay
} from 'rxjs/operators';

// Przykład 1: Wyszukiwanie na żywo (PRAWIDŁOWE użycie switchMap)
const searchInput = document.getElementById('search') as HTMLInputElement;
const searchTerm$ = fromEvent(searchInput, 'input').pipe(
  map(event => (event.target as HTMLInputElement).value),
  debounceTime(300),
  distinctUntilChanged()
);

function searchAPI(term: string) {
  console.log(`Wysyłam zapytanie dla: "${term}"`);
  return of(`Wyniki dla: ${term}`).pipe(
    delay(1000) // Symulacja opóźnienia API
  );
}

// Z switchMap - anuluje poprzednie requesty
searchTerm$.pipe(
  switchMap(term => searchAPI(term))
).subscribe(results => {
  console.log('Otrzymano:', results);
});
// Jeśli wpiszesz "prog" a potem szybko "program",
// zobaczysz tylko "Wyniki dla: program"

// Przykład 2: Problem z mergeMap (NIEPRAWIDŁOWE dla wyszukiwania)
searchTerm$.pipe(
  mergeMap(term => searchAPI(term))
).subscribe(results => {
  console.log('Otrzymano:', results);
});
// Jeśli wpiszesz "prog" a potem "program",
// zobaczysz OBA wyniki, potencjalnie w niewłaściwej kolejności!

// Przykład 3: Przełączanie między źródłami danych
const dataSource$ = of('użytkownicy', 'produkty', 'zamówienia');

function fetchData(source: string) {
  return interval(1000).pipe(
    map(i => `${source}: rekord ${i}`)
  );
}

dataSource$.pipe(
  switchMap(source => fetchData(source))
).subscribe(data => console.log(data));
// Emituje tylko z najnowszego źródła

// Przykład 4: Nawigacja z anulowaniem poprzednich requestów
interface RouteParams { id: string; }

const route$ = of<RouteParams>(
  { id: 'user-1' },
  { id: 'user-2' },
  { id: 'user-3' }
);

function loadUserData(id: string) {
  console.log(`Ładuję dane użytkownika: ${id}`);
  return of(`Dane dla ${id}`).pipe(delay(2000));
}

route$.pipe(
  switchMap(params => loadUserData(params.id))
).subscribe(data => console.log('Załadowano:', data));
// Załaduje tylko dane dla 'user-3'

// Przykład 5: Kiedy NIE używać switchMap
const clicksToSave$ = fromEvent(document.getElementById('save-btn'), 'click');

function saveData(data: any) {
  console.log('Zapisuję dane...');
  return of('Zapisano!').pipe(delay(1000));
}

// ZŁE: switchMap anuluje poprzednie zapisy!
clicksToSave$.pipe(
  switchMap(() => saveData({ timestamp: Date.now() }))
).subscribe();

// DOBRE: mergeMap pozwala wszystkim zapisom się zakończyć
clicksToSave$.pipe(
  mergeMap(() => saveData({ timestamp: Date.now() }))
).subscribe();

Diagram porównawczy:

sequenceDiagram
    participant Source as Źródło
    participant Operator
    participant Inner1 as Wewn. Obs 1
    participant Inner2 as Wewn. Obs 2
    participant Output as Wyjście

    Note over Source,Output: switchMap()
    Source->>Operator: wartość A
    Operator->>Inner1: subskrybuj
    Inner1-->>Output: wynik A1
    Source->>Operator: wartość B
    Operator->>Inner1: ANULUJ ❌
    Operator->>Inner2: subskrybuj
    Inner2-->>Output: wynik B1
    Inner2-->>Output: wynik B2

    Note over Source,Output: mergeMap()
    Source->>Operator: wartość A
    Operator->>Inner1: subskrybuj
    Inner1-->>Output: wynik A1
    Source->>Operator: wartość B
    Operator->>Inner2: subskrybuj (równolegle)
    Inner1-->>Output: wynik A2 ✓
    Inner2-->>Output: wynik B1 ✓
    Inner1-->>Output: wynik A3 ✓
    Inner2-->>Output: wynik B2 ✓

Kiedy używać którego operatora:

graph TD
    A[Potrzebuję spłaszczenia Observable] --> B{Czy nowe wartości<br/>unieważniają poprzednie?}
    B -->|TAK| C[switchMap]
    B -->|NIE| D{Czy kolejność ma znaczenie?}
    D -->|TAK| E[concatMap]
    D -->|NIE| F{Czy wszystkie operacje<br/>muszą się zakończyć?}
    F -->|TAK| G[mergeMap]
    F -->|NIE| H[exhaustMap]

    C -.-> I[Wyszukiwanie<br/>Autouzupełnianie<br/>Nawigacja]
    E -.-> J[Kolejka zadań<br/>Sekwencyjne API calls]
    G -.-> K[Analytics<br/>Zapis logów<br/>Batch processing]
    H -.-> L[Debouncing akcji<br/>Login/Submit forms]

Materiały

Jak działa operator concatMap() i kiedy jest przydatny?

Odpowiedź w 30 sekund: concatMap() spłaszcza Observable'e zachowując ścisłą kolejność - czeka aż poprzedni wewnętrzny Observable się zakończy przed subskrybowaniem kolejnego. Używaj go gdy kolejność operacji jest krytyczna (np. sekwencyjne zapisy do bazy, transakcje, kolejka zadań).

Odpowiedź w 2 minuty: Operator concatMap() łączy cechy map() i operatora konkatenacji - transformuje wartości na Observable'e i subskrybuje je sekwencyjnie, jeden po drugim. Kluczowa właściwość: nawet jeśli źródłowy Observable emituje wartości szybko, concatMap() buforuje je i przetwarza po kolei, czekając aż każdy wewnętrzny Observable zakończy się całkowicie przed rozpoczęciem kolejnego.

To gwarantuje, że wyniki są emitowane w tej samej kolejności co wartości źródłowe, niezależnie od tego jak długo trwa każda operacja asynchroniczna. Jest to fundamentalna różnica w porównaniu do mergeMap(), gdzie szybsze operacje mogą "wyprzedzić" wolniejsze, zmieniając kolejność wyników.

concatMap() jest niezbędny w scenariuszach gdzie:

  1. Kolejność ma krytyczne znaczenie - np. seria operacji zapisujących dane, gdzie każdy krok zależy od poprzedniego
  2. Operacje nie mogą się nakładać - np. transakcje bankowe, aktualizacje stanu w systemach eventowych
  3. Potrzebujemy gwarancji FIFO (First-In-First-Out) - przetwarzanie kolejek zadań, sekwencyjne wykonywanie komend

Wadą concatMap() jest potencjalnie niższa wydajność - jeśli masz 100 operacji i każda trwa sekundę, concatMap() zajmie 100 sekund, podczas gdy mergeMap() mógłby wykonać wszystkie równocześnie w ~1 sekundę. Jednak ta "wada" jest często zamierzoną cechą, gdy równoległość mogłaby spowodować problemy z race conditions czy niespójnością danych.

Przykład kodu:

import { of, from, interval } from 'rxjs';
import { concatMap, mergeMap, map, delay, take } from 'rxjs/operators';

// Przykład 1: Różnica między concatMap a mergeMap
console.log('=== mergeMap - bez gwarancji kolejności ===');
of(3, 2, 1).pipe(
  mergeMap(n => of(n).pipe(
    delay(n * 1000), // Im większa liczba, tym dłużej czeka
    map(() => `Zakończono: ${n}`)
  ))
).subscribe(console.log);
// Wynik: "Zakończono: 1", "Zakończono: 2", "Zakończono: 3"
// (w odwrotnej kolejności!)

console.log('=== concatMap - gwarantowana kolejność ===');
of(3, 2, 1).pipe(
  concatMap(n => of(n).pipe(
    delay(n * 1000),
    map(() => `Zakończono: ${n}`)
  ))
).subscribe(console.log);
// Wynik: "Zakończono: 3", "Zakończono: 2", "Zakończono: 1"
// (w oryginalnej kolejności!)

// Przykład 2: Sekwencyjne zapisywanie do bazy danych
interface UserUpdate {
  userId: number;
  field: string;
  value: any;
}

const updates: UserUpdate[] = [
  { userId: 1, field: 'email', value: 'nowy@email.com' },
  { userId: 1, field: 'name', value: 'Nowa Nazwa' },
  { userId: 1, field: 'verified', value: true }
];

function saveToDatabase(update: UserUpdate) {
  console.log(`Zapisuję: ${update.field} = ${update.value}`);
  return of(`Zapisano ${update.field}`).pipe(
    delay(Math.random() * 1000)
  );
}

// PRAWIDŁOWE: concatMap zapewnia kolejność
from(updates).pipe(
  concatMap(update => saveToDatabase(update))
).subscribe(
  result => console.log(result),
  error => console.error(error),
  () => console.log('Wszystkie aktualizacje zakończone')
);

// Przykład 3: Przetwarzanie kolejki zadań
interface Task {
  id: number;
  action: string;
  priority: number;
}

const taskQueue: Task[] = [
  { id: 1, action: 'Inicjalizacja', priority: 1 },
  { id: 2, action: 'Załaduj konfigurację', priority: 2 },
  { id: 3, action: 'Połącz z bazą', priority: 3 },
  { id: 4, action: 'Uruchom serwis', priority: 4 }
];

function executeTask(task: Task) {
  console.log(`[${new Date().toISOString()}] Wykonuję: ${task.action}`);
  return of(`✓ ${task.action}`).pipe(
    delay(1000) // Każde zadanie trwa 1 sekundę
  );
}

from(taskQueue).pipe(
  concatMap(task => executeTask(task))
).subscribe(
  result => console.log(result),
  error => console.error('Błąd zadania:', error),
  () => console.log('🎉 Wszystkie zadania wykonane!')
);

// Przykład 4: Transakcje bankowe (muszą być sekwencyjne!)
interface Transaction {
  type: 'withdraw' | 'deposit';
  amount: number;
  accountId: string;
}

const transactions: Transaction[] = [
  { type: 'deposit', amount: 100, accountId: 'ACC001' },
  { type: 'withdraw', amount: 50, accountId: 'ACC001' },
  { type: 'deposit', amount: 200, accountId: 'ACC001' }
];

let balance = 1000;

function processTransaction(transaction: Transaction) {
  return of(transaction).pipe(
    delay(500),
    map(tx => {
      if (tx.type === 'deposit') {
        balance += tx.amount;
        return `Wpłata: +${tx.amount} zł. Saldo: ${balance} zł`;
      } else {
        if (balance >= tx.amount) {
          balance -= tx.amount;
          return `Wypłata: -${tx.amount} zł. Saldo: ${balance} zł`;
        } else {
          throw new Error('Niewystarczające środki!');
        }
      }
    })
  );
}

console.log(`Stan początkowy: ${balance} zł`);
from(transactions).pipe(
  concatMap(tx => processTransaction(tx))
).subscribe(
  result => console.log(result),
  error => console.error('Błąd transakcji:', error)
);

// Przykład 5: Porównanie wydajności
console.log('=== Test wydajności ===');
const items = [1, 2, 3, 4, 5];

console.time('concatMap - sekwencyjnie');
from(items).pipe(
  concatMap(i => of(i).pipe(delay(100)))
).subscribe({
  complete: () => console.timeEnd('concatMap - sekwencyjnie')
  // ~500ms (5 * 100ms)
});

console.time('mergeMap - równolegle');
from(items).pipe(
  mergeMap(i => of(i).pipe(delay(100)))
).subscribe({
  complete: () => console.timeEnd('mergeMap - równolegle')
  // ~100ms (wszystkie równocześnie)
});

// Przykład 6: HTTP requests w określonej kolejności
interface ApiEndpoint {
  url: string;
  method: string;
}

const apiCalls: ApiEndpoint[] = [
  { url: '/api/auth/login', method: 'POST' },
  { url: '/api/user/profile', method: 'GET' },
  { url: '/api/user/settings', method: 'GET' }
];

function makeApiCall(endpoint: ApiEndpoint) {
  console.log(`Wywołuję: ${endpoint.method} ${endpoint.url}`);
  return of({ url: endpoint.url, data: 'mock data' }).pipe(
    delay(500)
  );
}

// Login musi się zakończyć przed pobraniem profilu i ustawień
from(apiCalls).pipe(
  concatMap(endpoint => makeApiCall(endpoint))
).subscribe(response => {
  console.log(`Odpowiedź z: ${response.url}`, response.data);
});

Diagram działania:

sequenceDiagram
    participant Source as Źródło
    participant Buffer as Bufor concatMap
    participant Inner as Wewn. Observable
    participant Output as Wyjście

    Source->>Buffer: wartość A
    Buffer->>Inner: subskrybuj A
    activate Inner
    Source->>Buffer: wartość B (czeka w buforze)
    Source->>Buffer: wartość C (czeka w buforze)
    Inner-->>Output: A1
    Inner-->>Output: A2
    Inner-->>Output: A3 (complete)
    deactivate Inner

    Buffer->>Inner: subskrybuj B
    activate Inner
    Inner-->>Output: B1
    Inner-->>Output: B2 (complete)
    deactivate Inner

    Buffer->>Inner: subskrybuj C
    activate Inner
    Inner-->>Output: C1 (complete)
    deactivate Inner

    Note over Source,Output: Kolejność zachowana: A → B → C

Porównanie operatorów:

graph LR
    subgraph "concatMap - Sekwencyjnie"
        A1[A] --> B1[proces A]
        B1 --> C1[wynik A]
        C1 --> D1[proces B]
        D1 --> E1[wynik B]
        E1 --> F1[proces C]
        F1 --> G1[wynik C]
    end

    subgraph "mergeMap - Równolegle"
        A2[A] --> B2[proces A]
        A3[B] --> C2[proces B]
        A4[C] --> D2[proces C]
        B2 --> E2[wyniki w dowolnej kolejności]
        C2 --> E2
        D2 --> E2
    end

Materiały

Czym jest operator exhaustMap() i w jakich scenariuszach go stosować?

Odpowiedź w 30 sekund: exhaustMap() ignoruje nowe wartości źródłowe dopóki poprzedni wewnętrzny Observable nie zakończy się. Używaj go do zapobiegania wielokrotnemu wywoływaniu akcji (np. podwójne kliknięcie przycisku submit, wielokrotne wywołania API podczas gdy poprzednie jeszcze trwa).

Odpowiedź w 2 minuty: Operator exhaustMap() jest często pomijanym, ale bardzo przydatnym operatorem spłaszczającym z unikalną strategią zarządzania subskrypcjami: gdy wewnętrzny Observable jest aktywny, wszystkie nowe wartości ze źródła są całkowicie ignorowane. Dopiero gdy wewnętrzny Observable się zakończy, exhaustMap() zasubskrybuje kolejną wartość ze źródła.

To zachowanie typu "wyczerpanie" (exhausting) jest przeciwieństwem switchMap() - zamiast anulować poprzednią operację na rzecz nowej, exhaustMap() chroni trwającą operację przed przerwaniem i odrzuca próby rozpoczęcia nowych operacji.

Najczęstsze przypadki użycia to scenariusze gdzie chcemy zapobiec "spamowaniu" akcji:

  1. Formularze logowania/rejestracji - wielokrotne kliknięcia przycisku "Zaloguj" powinny wywołać tylko jeden request
  2. Rate limiting - ograniczanie częstotliwości wywołań API
  3. Operacje kosztowne - zapobieganie nakładaniu się długotrwałych operacji (np. eksport danych, generowanie raportów)
  4. Debouncing akcji użytkownika - ignorowanie szybkich powtórzeń tej samej akcji

Kluczowa zaleta exhaustMap() to prostsza logika w porównaniu do manualnego flagowania stanu "operacja w toku". Zamiast ręcznie zarządzać stanem isLoading i sprawdzać go przed każdą akcją, exhaustMap() automatycznie odrzuca nadmiarowe wywołania.

Przykład kodu:

import { fromEvent, interval, of, Subject } from 'rxjs';
import { exhaustMap, mergeMap, take, delay, tap } from 'rxjs/operators';

// Przykład 1: Problem - wielokrotne kliknięcia przycisku (ZŁE)
const loginButton = document.getElementById('login-btn');
const clicks$ = fromEvent(loginButton, 'click');

function loginRequest() {
  console.log('🔄 Wysyłam request logowania...');
  return of('Zalogowano!').pipe(
    delay(2000), // Symulacja opóźnienia API
    tap(() => console.log('✅ Request zakończony'))
  );
}

// ZŁE: mergeMap pozwala na wiele równoczesnych requestów
console.log('=== Z mergeMap (PROBLEM) ===');
clicks$.pipe(
  tap(() => console.log('👆 Kliknięcie')),
  mergeMap(() => loginRequest())
).subscribe(result => console.log(result));
// Szybkie 3 kliknięcia = 3 równoczesne requesty! ❌

// DOBRE: exhaustMap ignoruje kliknięcia podczas trwania requestu
console.log('=== Z exhaustMap (ROZWIĄZANIE) ===');
clicks$.pipe(
  tap(() => console.log('👆 Kliknięcie')),
  exhaustMap(() => loginRequest())
).subscribe(result => console.log(result));
// Szybkie 3 kliknięcia = tylko 1 request! ✅

// Przykład 2: Formularz zapisu z ochroną przed duplikatami
interface FormData {
  username: string;
  email: string;
}

const submitButton = document.getElementById('submit-form');
const formSubmit$ = fromEvent(submitButton, 'click');

function saveFormData(data: FormData) {
  console.log('💾 Zapisuję dane...', data);
  return of({ success: true, id: Math.random() }).pipe(
    delay(3000) // Wolny endpoint
  );
}

formSubmit$.pipe(
  exhaustMap(() => {
    const formData: FormData = {
      username: 'user123',
      email: 'user@example.com'
    };
    return saveFormData(formData);
  })
).subscribe(
  response => console.log('✅ Zapisano:', response),
  error => console.error('❌ Błąd:', error)
);
// Wielokrotne kliknięcia "Zapisz" podczas zapisywania są ignorowane

// Przykład 3: Refresh button z cooldownem
const refreshButton = document.getElementById('refresh-btn');
const refresh$ = fromEvent(refreshButton, 'click');

function refreshData() {
  console.log('🔄 Odświeżam dane...');
  return interval(1000).pipe(
    take(5),
    tap(i => console.log(`Krok ${i + 1}/5`))
  );
}

refresh$.pipe(
  tap(() => console.log('👆 Kliknięto refresh')),
  exhaustMap(() => refreshData())
).subscribe(
  value => console.log('Dane:', value),
  error => console.error(error),
  () => console.log('✅ Odświeżanie zakończone')
);
// Próby odświeżenia podczas trwającego procesu są ignorowane

// Przykład 4: Rate limiting dla API calls
const searchSubject = new Subject<string>();

function searchAPI(query: string) {
  console.log(`🔍 Szukam: "${query}"`);
  return of([`Wynik 1 dla ${query}`, `Wynik 2 dla ${query}`]).pipe(
    delay(2000)
  );
}

searchSubject.pipe(
  tap(query => console.log(`📝 Wprowadzono: "${query}"`)),
  exhaustMap(query => searchAPI(query))
).subscribe(results => {
  console.log('📊 Wyniki:', results);
});

// Symulacja szybkiego wpisywania
searchSubject.next('java');
setTimeout(() => searchSubject.next('javascript'), 100); // Ignorowane
setTimeout(() => searchSubject.next('js'), 200); // Ignorowane
setTimeout(() => searchSubject.next('jsx'), 300); // Ignorowane
setTimeout(() => searchSubject.next('typescript'), 2500); // Wykonane (po zakończeniu pierwszego)

// Przykład 5: Eksport danych - długotrwała operacja
const exportButton = document.getElementById('export-btn');
const export$ = fromEvent(exportButton, 'click');

function exportLargeDataset() {
  console.log('📤 Rozpoczynam eksport...');
  return interval(500).pipe(
    take(10),
    tap(progress => {
      console.log(`Postęp: ${(progress + 1) * 10}%`);
    }),
    delay(100)
  );
}

export$.pipe(
  tap(() => console.log('👆 Próba eksportu')),
  exhaustMap(() => exportLargeDataset())
).subscribe({
  complete: () => console.log('✅ Eksport zakończony!')
});
// Kliknięcia podczas trwającego eksportu są ignorowane

// Przykład 6: Porównanie wszystkich operatorów spłaszczających
console.log('\n=== PORÓWNANIE OPERATORÓW ===\n');

const source$ = interval(1000).pipe(take(3)); // Emituje: 0, 1, 2

function innerObservable(value: number) {
  return of(`Wartość: ${value}`).pipe(delay(2500));
}

// mergeMap - wszystkie równocześnie
console.log('mergeMap:');
source$.pipe(
  tap(v => console.log(`  Emitowano: ${v}`)),
  mergeMap(v => innerObservable(v))
).subscribe(result => console.log(`  Wynik: ${result}`));
// Emitowane: 0, 1, 2
// Wyniki: wszystkie 3 (równocześnie)

// concatMap - sekwencyjnie
console.log('\nconcatMap:');
source$.pipe(
  tap(v => console.log(`  Emitowano: ${v}`)),
  concatMap(v => innerObservable(v))
).subscribe(result => console.log(`  Wynik: ${result}`));
// Emitowane: 0, 1, 2
// Wyniki: 0 → 1 → 2 (po kolei)

// switchMap - anuluje poprzednie
console.log('\nswitchMap:');
source$.pipe(
  tap(v => console.log(`  Emitowano: ${v}`)),
  switchMap(v => innerObservable(v))
).subscribe(result => console.log(`  Wynik: ${result}`));
// Emitowane: 0, 1, 2
// Wynik: tylko 2 (poprzednie anulowane)

// exhaustMap - ignoruje nowe podczas trwania
console.log('\nexhaustMap:');
source$.pipe(
  tap(v => console.log(`  Emitowano: ${v}`)),
  exhaustMap(v => innerObservable(v))
).subscribe(result => console.log(`  Wynik: ${result}`));
// Emitowane: 0, 1, 2
// Wynik: tylko 0 (1 i 2 zignorowane)

// Przykład 7: Praktyczne użycie - ochrona przed double-click
interface PaymentRequest {
  amount: number;
  currency: string;
}

const payButton = document.getElementById('pay-btn');
const payment$ = fromEvent(payButton, 'click');

function processPayment(request: PaymentRequest) {
  console.log(`💳 Przetwarzam płatność: ${request.amount} ${request.currency}`);
  return of({ transactionId: Date.now(), status: 'success' }).pipe(
    delay(3000)
  );
}

payment$.pipe(
  exhaustMap(() => processPayment({ amount: 99.99, currency: 'PLN' }))
).subscribe(
  response => {
    console.log(`✅ Płatność zatwierdzona! ID: ${response.transactionId}`);
    // Można teraz zaktualizować UI, pokazać potwierdzenie, etc.
  },
  error => console.error('❌ Błąd płatności:', error)
);
// Double-click nie spowoduje podwójnej płatności!

Diagram porównawczy strategii:

sequenceDiagram
    participant User as Użytkownik
    participant Operator
    participant API

    Note over User,API: exhaustMap() - Ignoruje nowe podczas trwania
    User->>Operator: Kliknięcie 1
    Operator->>API: Request 1
    activate API
    User->>Operator: Kliknięcie 2 ❌ IGNOROWANE
    User->>Operator: Kliknięcie 3 ❌ IGNOROWANE
    API-->>Operator: Odpowiedź 1
    deactivate API
    User->>Operator: Kliknięcie 4 ✓
    Operator->>API: Request 2
    activate API
    API-->>Operator: Odpowiedź 2
    deactivate API

    Note over User,API: mergeMap() - Wszystkie wykonane
    User->>Operator: Kliknięcie 1
    Operator->>API: Request 1
    User->>Operator: Kliknięcie 2
    Operator->>API: Request 2 (równolegle)
    User->>Operator: Kliknięcie 3
    Operator->>API: Request 3 (równolegle)

Macierz decyzyjna:

graph TD
    A[Otrzymuję nową wartość podczas<br/>aktywnego wewnętrznego Observable] --> B{exhaustMap}
    A --> C{switchMap}
    A --> D{mergeMap}
    A --> E{concatMap}

    B --> B1[❌ IGNORUJ<br/>nową wartość]
    C --> C1[❌ ANULUJ poprzedni<br/>✓ Subskrybuj nowy]
    D --> D1[✓ Subskrybuj nowy<br/>równolegle]
    E --> E1[📦 Dodaj do kolejki<br/>✓ Subskrybuj później]

    style B1 fill:#ffcccc
    style C1 fill:#ffffcc
    style D1 fill:#ccffcc
    style E1 fill:#ccccff

Materiały

Jak używać operatora pluck() do wyciągania właściwości z obiektów?

Odpowiedź w 30 sekund: pluck() to skrótowa wersja map() do wyciągania zagnieżdżonych właściwości z obiektów. Zamiast map(x => x.user.name) możesz użyć pluck('user', 'name'). Od RxJS 8 operator jest deprecated na rzecz map() z opcjonalnym chainingiem.

Odpowiedź w 2 minuty: Operator pluck() upraszcza często spotykany wzorzec wyciągania właściwości z emitowanych obiektów. Przyjmuje listę kluczy (strings lub numbers) reprezentujących ścieżkę do zagnieżdżonej właściwości i zwraca Observable emitujący tylko te wartości.

Składnia jest prostsza i bardziej deklaratywna niż równoważne użycie map(). Zamiast pisać map(obj => obj.data.user.profile.name), możesz użyć pluck('data', 'user', 'profile', 'name'). To szczególnie przydatne przy głęboko zagnieżdżonych strukturach danych, typowych w odpowiedziach API.

Ważna uwaga o deprecation: Od RxJS 8, operator pluck() jest oznaczony jako deprecated i zaleca się używanie map() w połączeniu z opcjonalnym chainingiem (?.) dostępnym w nowoczesnym TypeScript/JavaScript. Powód deprecation: pluck() nie oferuje żadnych znaczących korzyści wydajnościowych, a nowoczesna składnia języka sprawia, że map() jest równie zwięzły i bardziej elastyczny.

Mimo deprecation, pluck() jest nadal powszechnie spotykany w starszym kodzie i warto go rozumieć. Jeśli wartość na ścieżce nie istnieje (undefined), pluck() emituje undefined bez rzucania błędu.

Przykład kodu:

import { of, fromEvent } from 'rxjs';
import { pluck, map } from 'rxjs/operators';

// Przykład 1: Podstawowe użycie pluck()
interface User {
  id: number;
  name: string;
  email: string;
}

const users$ = of<User>(
  { id: 1, name: 'Anna Kowalska', email: 'anna@example.com' },
  { id: 2, name: 'Jan Nowak', email: 'jan@example.com' }
);

// Stary sposób z map()
users$.pipe(
  map(user => user.name)
).subscribe(name => console.log(name));

// Sposób z pluck() (deprecated od RxJS 8)
users$.pipe(
  pluck('name')
).subscribe(name => console.log(name));

// Nowoczesny sposób (zalecany)
users$.pipe(
  map(user => user.name)
).subscribe(name => console.log(name));

// Przykład 2: Zagnieżdżone właściwości
interface ApiResponse {
  status: number;
  data: {
    user: {
      profile: {
        firstName: string;
        lastName: string;
        address: {
          city: string;
          country: string;
        };
      };
    };
  };
}

const apiResponse$ = of<ApiResponse>({
  status: 200,
  data: {
    user: {
      profile: {
        firstName: 'Anna',
        lastName: 'Kowalska',
        address: {
          city: 'Warszawa',
          country: 'Polska'
        }
      }
    }
  }
});

// Z map() - verbose
apiResponse$.pipe(
  map(response => response.data.user.profile.address.city)
).subscribe(city => console.log(city)); // 'Warszawa'

// Z pluck() - zwięzłe (deprecated)
apiResponse$.pipe(
  pluck('data', 'user', 'profile', 'address', 'city')
).subscribe(city => console.log(city)); // 'Warszawa'

// Nowoczesny sposób z optional chaining
apiResponse$.pipe(
  map(response => response.data?.user?.profile?.address?.city)
).subscribe(city => console.log(city)); // 'Warszawa'

// Przykład 3: Obsługa undefined wartości
interface Product {
  id: number;
  name: string;
  details?: {
    description?: string;
  };
}

const products$ = of<Product>(
  { id: 1, name: 'Laptop', details: { description: 'Szybki laptop' } },
  { id: 2, name: 'Mysz' }, // Brak details
  { id: 3, name: 'Klawiatura', details: {} } // Brak description
);

// pluck() emituje undefined gdy właściwość nie istnieje
products$.pipe(
  pluck('details', 'description')
).subscribe(desc => console.log(desc ?? 'Brak opisu'));
// 'Szybki laptop', 'Brak opisu', 'Brak opisu'

// Z map() i optional chaining - bardziej eksplicytne
products$.pipe(
  map(product => product.details?.description ?? 'Brak opisu')
).subscribe(desc => console.log(desc));

// Przykład 4: pluck() z tablicami (indeksy numeryczne)
interface ArrayData {
  items: string[];
}

const arrayData$ = of<ArrayData>(
  { items: ['pierwszy', 'drugi', 'trzeci'] }
);

// Wyciąganie konkretnego elementu tablicy
arrayData$.pipe(
  pluck('items', 0) // Pierwszy element
).subscribe(item => console.log(item)); // 'pierwszy'

// Przykład 5: Zdarzenia DOM
const input = document.getElementById('search-input') as HTMLInputElement;
const input$ = fromEvent<InputEvent>(input, 'input');

// Stary sposób z pluck()
input$.pipe(
  pluck('target', 'value')
).subscribe(value => console.log('Wpisano:', value));

// Nowoczesny sposób (type-safe!)
input$.pipe(
  map(event => (event.target as HTMLInputElement).value)
).subscribe(value => console.log('Wpisano:', value));

// Przykład 6: Kombinacja pluck() z innymi operatorami
interface BlogPost {
  id: number;
  title: string;
  author: {
    name: string;
    posts: number;
  };
  tags: string[];
}

const posts$ = of<BlogPost>(
  {
    id: 1,
    title: 'Wprowadzenie do RxJS',
    author: { name: 'Anna', posts: 15 },
    tags: ['rxjs', 'javascript', 'reactive']
  },
  {
    id: 2,
    title: 'Zaawansowane operatory',
    author: { name: 'Jan', posts: 23 },
    tags: ['rxjs', 'advanced']
  }
);

// Łańcuch operacji z pluck()
import { filter } from 'rxjs/operators';

posts$.pipe(
  pluck('author', 'name'),
  filter((name: string) => name.startsWith('A'))
).subscribe(name => console.log('Autor:', name)); // 'Anna'

// Przykład 7: Migracja z pluck() do map()
console.log('=== MIGRACJA Z pluck() ===\n');

interface Config {
  app: {
    settings: {
      theme: string;
      language: string;
    };
  };
}

const config$ = of<Config>({
  app: {
    settings: {
      theme: 'dark',
      language: 'pl'
    }
  }
});

// Przed (RxJS 7 i wcześniejsze)
config$.pipe(
  pluck('app', 'settings', 'theme')
).subscribe(theme => console.log('Motyw:', theme));

// Po (RxJS 8+, zalecane)
config$.pipe(
  map(config => config.app.settings.theme)
).subscribe(theme => console.log('Motyw:', theme));

// Z optional chaining dla bezpieczeństwa
config$.pipe(
  map(config => config.app?.settings?.theme ?? 'light')
).subscribe(theme => console.log('Motyw:', theme));

// Przykład 8: Dlaczego map() jest lepszy (Type Safety)
interface TypedData {
  count: number;
  values: string[];
}

const typed$ = of<TypedData>({ count: 3, values: ['a', 'b', 'c'] });

// pluck() - brak sprawdzania typów podczas kompilacji
typed$.pipe(
  pluck('nonexistent') // TypeScript nie wyłapie błędu!
).subscribe(val => console.log(val)); // undefined

// map() - pełne sprawdzanie typów
typed$.pipe(
  map(data => data.count) // TypeScript weryfikuje że 'count' istnieje
).subscribe(count => console.log(count)); // 3

// Przykład 9: Alternatywne podejścia w nowoczesnym JS/TS
const modernApproach$ = of<User>(
  { id: 1, name: 'Anna', email: 'anna@example.com' }
);

// Destrukturyzacja w map()
modernApproach$.pipe(
  map(({ name, email }) => ({ name, email }))
).subscribe(user => console.log(user));

// Property shorthand
modernApproach$.pipe(
  map(user => ({
    displayName: user.name,
    contact: user.email
  }))
).subscribe(formatted => console.log(formatted));

Porównanie składni:

// Przykład porównawczy
interface DeepObject {
  level1: {
    level2: {
      level3: {
        value: string;
      };
    };
  };
}

const deep$ = of<DeepObject>({
  level1: {
    level2: {
      level3: {
        value: 'głęboka wartość'
      }
    }
  }
});

// 1. pluck() - zwięzłe, ale deprecated
deep$.pipe(pluck('level1', 'level2', 'level3', 'value'));

// 2. map() klasyczny - bezpieczny typowo
deep$.pipe(map(obj => obj.level1.level2.level3.value));

// 3. map() z optional chaining - najbezpieczniejszy
deep$.pipe(map(obj => obj.level1?.level2?.level3?.value));

// 4. map() z destrukturyzacją - czytelny
deep$.pipe(
  map(({ level1: { level2: { level3: { value } } } }) => value)
);

Diagram migracji:

graph LR
    A[RxJS 7 i wcześniejsze] --> B[pluck 'user', 'name']
    A --> C[map x => x.user.name]

    D[RxJS 8+] --> E[map x => x.user.name ✓]
    D --> F[map x => x.user?.name ✓✓]
    D --> G[pluck 'user', 'name' ⚠️ deprecated]

    style B fill:#ffcccc
    style C fill:#ffffcc
    style E fill:#ccffcc
    style F fill:#ccffff
    style G fill:#ffcccc

Materiały

Jak działa operator scan() i czym różni się od reduce()?

Odpowiedź w 30 sekund: scan() emituje wartość akumulatora po każdej wartości źródłowej (podobnie jak Array.reduce() ale z emisją na każdym kroku), podczas gdy reduce() emituje tylko końcowy wynik po zakończeniu Observable. scan() jest używany do śledzenia stanu w czasie rzeczywistym.

Odpowiedź w 2 minuty: Operator scan() jest jednym z najpotężniejszych operatorów do zarządzania stanem w RxJS. Przyjmuje funkcję akumulującą (accumulator function) i opcjonalną wartość początkową (seed value), a następnie stosuje tę funkcję do każdej wartości emitowanej przez źródłowy Observable, emitując bieżącą wartość akumulatora po każdym kroku.

Analogia do Array API: jeśli Array.reduce() zwraca tylko końcowy wynik, to scan() to jakby "reduce z historią" - widzisz każdy pośredni krok akumulacji. To czyni scan() idealnym do implementacji:

  • Liczników i sum bieżących - np. suma kliknięć, licznik zdarzeń
  • Zarządzania stanem aplikacji - przechowywanie i aktualizowanie stanu w reaktywny sposób
  • Historii zmian - budowanie tablic z historią wartości
  • Sliding window calculations - obliczenia na oknach danych (np. średnia krocząca)

Kluczowa różnica między scan() a reduce():

  • scan() emituje przy każdej wartości źródłowej → użyteczny dla continuous streams
  • reduce() emituje tylko raz, gdy źródło się zakończy → użyteczny dla finite streams

scan() nie zakłada, że Observable się zakończy, co czyni go odpowiednim dla nieskończonych strumieni zdarzeń. reduce() czeka na complete() i dopiero wtedy emituje wynik, więc nie zadziała z nigdy niekończącymi się strumieniami.

Przykład kodu:

import { of, fromEvent, interval, Subject } from 'rxjs';
import { scan, reduce, map, take } from 'rxjs/operators';

// Przykład 1: Różnica między scan() a reduce()
const numbers$ = of(1, 2, 3, 4, 5);

console.log('=== scan() - emituje przy każdym kroku ===');
numbers$.pipe(
  scan((acc, value) => acc + value, 0)
).subscribe(result => console.log(result));
// Wynik: 1, 3, 6, 10, 15
// (suma bieżąca po każdej wartości)

console.log('\n=== reduce() - emituje tylko końcowy wynik ===');
numbers$.pipe(
  reduce((acc, value) => acc + value, 0)
).subscribe(result => console.log(result));
// Wynik: 15
// (tylko końcowa suma)

// Przykład 2: Licznik kliknięć
const button = document.getElementById('counter-btn');
const clicks$ = fromEvent(button, 'click');

// Prosty licznik
clicks$.pipe(
  scan(count => count + 1, 0)
).subscribe(count => {
  console.log(`Liczba kliknięć: ${count}`);
  document.getElementById('counter-display').textContent = `${count}`;
});

// Przykład 3: Zarządzanie stanem aplikacji
interface AppState {
  count: number;
  lastAction: string;
  history: string[];
}

type Action =
  | { type: 'INCREMENT' }
  | { type: 'DECREMENT' }
  | { type: 'RESET' };

const actions$ = new Subject<Action>();

const initialState: AppState = {
  count: 0,
  lastAction: 'INIT',
  history: []
};

const state$ = actions$.pipe(
  scan((state, action) => {
    switch (action.type) {
      case 'INCREMENT':
        return {
          count: state.count + 1,
          lastAction: 'INCREMENT',
          history: [...state.history, 'INCREMENT']
        };
      case 'DECREMENT':
        return {
          count: state.count - 1,
          lastAction: 'DECREMENT',
          history: [...state.history, 'DECREMENT']
        };
      case 'RESET':
        return {
          count: 0,
          lastAction: 'RESET',
          history: [...state.history, 'RESET']
        };
      default:
        return state;
    }
  }, initialState)
);

state$.subscribe(state => {
  console.log('Stan aplikacji:', state);
});

// Wysyłanie akcji
actions$.next({ type: 'INCREMENT' }); // count: 1
actions$.next({ type: 'INCREMENT' }); // count: 2
actions$.next({ type: 'DECREMENT' }); // count: 1
actions$.next({ type: 'RESET' });     // count: 0

// Przykład 4: Suma krocząca (running total)
const transactions$ = of(
  { type: 'wpłata', amount: 100 },
  { type: 'wypłata', amount: 50 },
  { type: 'wpłata', amount: 200 },
  { type: 'wypłata', amount: 75 }
);

transactions$.pipe(
  scan((balance, transaction) => {
    const newBalance = transaction.type === 'wpłata'
      ? balance + transaction.amount
      : balance - transaction.amount;
    console.log(
      `${transaction.type}: ${transaction.amount} zł | Saldo: ${newBalance} zł`
    );
    return newBalance;
  }, 1000) // Saldo początkowe: 1000 zł
).subscribe(
  finalBalance => console.log(`\nKońcowe saldo: ${finalBalance} zł`)
);

// Przykład 5: Budowanie tablicy z historią
const values$ = of('A', 'B', 'C', 'D');

values$.pipe(
  scan((acc, value) => [...acc, value], [] as string[])
).subscribe(array => console.log('Tablica:', array));
// ['A']
// ['A', 'B']
// ['A', 'B', 'C']
// ['A', 'B', 'C', 'D']

// Przykład 6: Średnia krocząca (moving average)
const dataPoints$ = of(10, 20, 30, 40, 50);

interface MovingAverage {
  values: number[];
  average: number;
}

dataPoints$.pipe(
  scan((acc, value) => {
    // Przechowujemy ostatnie 3 wartości
    const newValues = [...acc.values, value].slice(-3);
    const average = newValues.reduce((sum, v) => sum + v, 0) / newValues.length;
    return { values: newValues, average };
  }, { values: [], average: 0 } as MovingAverage)
).subscribe(result => {
  console.log(`Wartości: [${result.values}] | Średnia: ${result.average.toFixed(2)}`);
});
// Wartości: [10] | Średnia: 10.00
// Wartości: [10,20] | Średnia: 15.00
// Wartości: [10,20,30] | Średnia: 20.00
// Wartości: [20,30,40] | Średnia: 30.00
// Wartości: [30,40,50] | Średnia: 40.00

// Przykład 7: Maksimum i minimum dotychczasowe
const temps$ = of(15, 23, 18, 30, 12, 25);

interface MinMax {
  current: number;
  min: number;
  max: number;
}

temps$.pipe(
  scan((acc, temp) => ({
    current: temp,
    min: Math.min(acc.min, temp),
    max: Math.max(acc.max, temp)
  }), { current: 0, min: Infinity, max: -Infinity } as MinMax)
).subscribe(result => {
  console.log(
    `Temp: ${result.current}°C | Min: ${result.min}°C | Max: ${result.max}°C`
  );
});

// Przykład 8: Zliczanie wystąpień (frequency counter)
const words$ = of('kot', 'pies', 'kot', 'ptak', 'pies', 'kot');

words$.pipe(
  scan((counts, word) => ({
    ...counts,
    [word]: (counts[word] || 0) + 1
  }), {} as Record<string, number>)
).subscribe(counts => console.log('Licznik:', counts));
// { kot: 1 }
// { kot: 1, pies: 1 }
// { kot: 2, pies: 1 }
// { kot: 2, pies: 1, ptak: 1 }
// { kot: 2, pies: 2, ptak: 1 }
// { kot: 3, pies: 2, ptak: 1 }

// Przykład 9: Śledzenie czasu sesji
const userActions$ = interval(1000).pipe(take(10));

interface Session {
  startTime: number;
  lastActionTime: number;
  totalActions: number;
  duration: number; // w sekundach
}

userActions$.pipe(
  scan((session, _) => {
    const now = Date.now();
    return {
      startTime: session.startTime,
      lastActionTime: now,
      totalActions: session.totalActions + 1,
      duration: Math.floor((now - session.startTime) / 1000)
    };
  }, {
    startTime: Date.now(),
    lastActionTime: Date.now(),
    totalActions: 0,
    duration: 0
  } as Session)
).subscribe(session => {
  console.log(
    `Akcja #${session.totalActions} | ` +
    `Czas trwania sesji: ${session.duration}s`
  );
});

// Przykład 10: Problem z reduce() na nieskończonym strumieniu
console.log('=== reduce() z nieskończonym strumieniem ===');
const infinite$ = interval(1000).pipe(take(5));

infinite$.pipe(
  reduce((acc, val) => acc + val, 0)
).subscribe(
  result => console.log('reduce() wynik:', result),
  error => console.error(error),
  () => console.log('reduce() zakończone')
);
// Emituje tylko: "reduce() wynik: 10" po zakończeniu (0+1+2+3+4)

console.log('\n=== scan() z nieskończonym strumieniem ===');
infinite$.pipe(
  scan((acc, val) => acc + val, 0)
).subscribe(
  result => console.log('scan() wynik:', result),
  error => console.error(error),
  () => console.log('scan() zakończone')
);
// Emituje: 0, 1, 3, 6, 10 (suma bieżąca przy każdej wartości)

// Przykład 11: Indeks bez wartości początkowej
const letters$ = of('a', 'b', 'c');

// scan() bez seed - pierwszy element staje się acc
letters$.pipe(
  scan((acc, value, index) => {
    console.log(`Krok ${index}: acc="${acc}", value="${value}"`);
    return acc + value;
  })
).subscribe(result => console.log('Wynik:', result));
// Krok 1: acc="a", value="b" → Wynik: "ab"
// Krok 2: acc="ab", value="c" → Wynik: "abc"

Diagram porównawczy:

sequenceDiagram
    participant Source as Źródło [1, 2, 3, 4]
    participant Scan as scan((acc,v)=>acc+v, 0)
    participant Reduce as reduce((acc,v)=>acc+v, 0)
    participant Output1 as Output scan()
    participant Output2 as Output reduce()

    Source->>Scan: 1
    Scan->>Output1: 1 (0+1)
    Source->>Reduce: 1
    Note over Reduce: Akumuluje: 1

    Source->>Scan: 2
    Scan->>Output1: 3 (1+2)
    Source->>Reduce: 2
    Note over Reduce: Akumuluje: 3

    Source->>Scan: 3
    Scan->>Output1: 6 (3+3)
    Source->>Reduce: 3
    Note over Reduce: Akumuluje: 6

    Source->>Scan: 4
    Scan->>Output1: 10 (6+4)
    Source->>Reduce: 4
    Note over Reduce: Akumuluje: 10

    Note over Source: complete()
    Reduce->>Output2: 10 (tylko teraz!)

    Note over Output1: scan() emitował: 1, 3, 6, 10
    Note over Output2: reduce() emitował: 10

Przypadki użycia:

graph TD
    A[Potrzebuję akumulacji wartości] --> B{Czy Observable<br/>kiedykolwiek się zakończy?}
    B -->|NIE/Nie wiem| C[scan]
    B -->|TAK| D{Czy potrzebuję<br/>wyników pośrednich?}
    D -->|TAK| C
    D -->|NIE| E[reduce]

    C --> F[Use cases:<br/>- Liczniki w czasie rzeczywistym<br/>- Stan aplikacji<br/>- Streaming analytics<br/>- Live dashboards]
    E --> G[Use cases:<br/>- Końcowa suma/agregacja<br/>- Batch processing<br/>- Report generation]

    style C fill:#ccffcc
    style E fill:#ccccff

Materiały

Operatory Filtrujące w RxJS

18. Jak działa operator filter() w RxJS?

Odpowiedź w 30 sekund: Operator filter() w RxJS działa podobnie jak Array.prototype.filter() - przepuszcza tylko te wartości ze strumienia, które spełniają określony warunek (predykat). Wartości, które nie spełniają warunku, są odrzucane i nie trafiają do subskrybenta.

Odpowiedź w 2 minuty: Operator filter() jest jednym z podstawowych operatorów transformacyjnych w RxJS, który pozwala selektywnie przepuszczać wartości emitowane przez Observable. Przyjmuje funkcję predykatu, która jest wywoływana dla każdej emitowanej wartości i zwraca true (wartość przechodzi) lub false (wartość jest odrzucana).

Predykat otrzymuje trzy argumenty: aktualną wartość, jej indeks w sekwencji oraz źródłowy Observable. Dzięki temu można tworzyć złożone warunki filtrowania. Operator zachowuje kolejność wartości i nie modyfikuje samych wartości - jedynie decyduje, które z nich zostaną przekazane dalej.

filter() jest często używany do oczyszczania strumieni danych, walidacji wejścia użytkownika, lub wybierania określonych typów zdarzeń. Jest to operator synchroniczny - decyzja o przepuszczeniu wartości jest podejmowana natychmiast. W przypadku Observable, który się zakończy lub zwróci błąd, filter() przekaże to zakończenie lub błąd bez zmian.

Warto pamiętać, że filter() nie zmienia momentu emisji wartości - jeśli wartość przeszła przez filtr, jest emitowana w tym samym momencie, co w źródłowym Observable. To odróżnia go od operatorów czasowych jak debounceTime() czy throttleTime().

Przykład kodu:

import { fromEvent, map, filter } from 'rxjs';

// Filtrowanie zdarzeń kliknięcia tylko na przyciski
const clicks$ = fromEvent<MouseEvent>(document, 'click');

const buttonClicks$ = clicks$.pipe(
  filter(event => (event.target as HTMLElement).tagName === 'BUTTON')
);

buttonClicks$.subscribe(event => {
  console.log('Kliknięto przycisk:', event.target);
});

// Filtrowanie liczb parzystych
import { interval } from 'rxjs';

const numbers$ = interval(1000);

const evenNumbers$ = numbers$.pipe(
  filter(num => num % 2 === 0)
);

evenNumbers$.subscribe(num => {
  console.log('Liczba parzysta:', num); // 0, 2, 4, 6...
});

// Filtrowanie z użyciem indeksu
const firstFiveEvenNumbers$ = numbers$.pipe(
  filter((num, index) => num % 2 === 0 && index < 10)
);

// Filtrowanie obiektów (np. użytkowników pełnoletnich)
import { of } from 'rxjs';

interface User {
  name: string;
  age: number;
}

const users$ = of<User>(
  { name: 'Jan', age: 17 },
  { name: 'Anna', age: 25 },
  { name: 'Piotr', age: 30 }
);

const adults$ = users$.pipe(
  filter(user => user.age >= 18)
);

adults$.subscribe(user => {
  console.log('Użytkownik pełnoletni:', user.name);
  // Wynik: Anna, Piotr
});

Materiały:

19. Czym są operatory take(), takeUntil() i takeWhile()?

Odpowiedź w 30 sekund: take(n) pobiera tylko pierwsze n wartości i kończy strumień, takeUntil(notifier$) pobiera wartości aż do momentu, gdy inny Observable wyemituje wartość, a takeWhile(predicate) pobiera wartości dopóki spełniają warunek. Wszystkie te operatory automatycznie kończą subskrypcję.

Odpowiedź w 2 minuty: Operatory z rodziny take* służą do kontrolowanego pobierania wartości ze strumienia i automatycznego kończenia subskrypcji, co jest kluczowe dla zarządzania pamięcią i unikania wycieków.

take(count) emituje tylko określoną liczbę pierwszych wartości ze źródłowego Observable, a następnie kończy się. Jest to przydatne, gdy potrzebujemy tylko kilku pierwszych wartości lub chcemy ograniczyć liczbę wykonań operacji. Na przykład take(1) jest często używany do pobrania pojedynczej wartości, podobnie jak Promise.

takeUntil(notifier$) jest jednym z najważniejszych operatorów do zarządzania cyklem życia subskrypcji. Emituje wartości ze źródła dopóki Observable notifier$ nie wyemituje swojej pierwszej wartości - wtedy natychmiast kończy się i anuluje subskrypcję źródła. Jest powszechnie używany w komponentach Angular z wzorcem destroy subject do automatycznego czyszczenia subskrypcji przy niszczeniu komponentu.

takeWhile(predicate, inclusive?) emituje wartości dopóki spełniają one warunek określony w predykacie. Gdy tylko pojawi się wartość nie spełniająca warunku, operator kończy się. Parametr inclusive (domyślnie false) określa, czy ostatnia wartość, która nie spełniła warunku, powinna być jeszcze wyemitowana. Ten operator jest przydatny do pobierania wartości na podstawie ich zawartości, a nie tylko pozycji w sekwencji.

Przykład kodu:

import { interval, fromEvent, Subject } from 'rxjs';
import { take, takeUntil, takeWhile } from 'rxjs';

// 1. take() - pobierz tylko 5 pierwszych wartości
const numbers$ = interval(1000);

numbers$.pipe(
  take(5)
).subscribe({
  next: num => console.log('take:', num),
  complete: () => console.log('take zakończony')
});
// Wynik: 0, 1, 2, 3, 4, potem complete

// 2. takeUntil() - wzorzec z Angular do czyszczenia subskrypcji
class MyComponent {
  private destroy$ = new Subject<void>();

  ngOnInit() {
    // Wszystkie subskrypcje z takeUntil zostaną automatycznie anulowane
    interval(1000).pipe(
      takeUntil(this.destroy$)
    ).subscribe(num => {
      console.log('Licznik:', num);
    });

    // Inna subskrypcja również z takeUntil
    fromEvent(document, 'click').pipe(
      takeUntil(this.destroy$)
    ).subscribe(() => {
      console.log('Kliknięcie');
    });
  }

  ngOnDestroy() {
    // Jedno wywołanie kończy wszystkie subskrypcje
    this.destroy$.next();
    this.destroy$.complete();
  }
}

// 3. takeWhile() - pobieraj dopóki wartość < 5
numbers$.pipe(
  takeWhile(num => num < 5)
).subscribe({
  next: num => console.log('takeWhile:', num),
  complete: () => console.log('takeWhile zakończony')
});
// Wynik: 0, 1, 2, 3, 4, potem complete (5 już nie przejdzie)

// 4. takeWhile() z inclusive=true
numbers$.pipe(
  takeWhile(num => num < 5, true) // inclusive = true
).subscribe({
  next: num => console.log('takeWhile inclusive:', num),
  complete: () => console.log('zakończony')
});
// Wynik: 0, 1, 2, 3, 4, 5, potem complete (5 zostaje wyemitowane)

// 5. Praktyczny przykład - pobieranie danych dopóki użytkownik jest zalogowany
import { BehaviorSubject } from 'rxjs';

const isLoggedIn$ = new BehaviorSubject<boolean>(true);

interval(1000).pipe(
  takeWhile(() => isLoggedIn$.value)
).subscribe(num => {
  console.log('Dane dla zalogowanego użytkownika:', num);
});

// Po 5 sekundach "wyloguj" użytkownika
setTimeout(() => {
  isLoggedIn$.next(false);
  console.log('Użytkownik wylogowany - subskrypcja zakończona');
}, 5000);

// 6. Łączenie take() z innymi operatorami
import { map } from 'rxjs';

fromEvent(document, 'click').pipe(
  take(3), // tylko 3 pierwsze kliknięcia
  map(event => ({ x: (event as MouseEvent).clientX, y: (event as MouseEvent).clientY }))
).subscribe({
  next: coords => console.log('Pozycja kliknięcia:', coords),
  complete: () => console.log('Zapisano 3 kliknięcia')
});

Materiały:

20. Jak używać operatorów skip(), skipUntil() i skipWhile()?

Odpowiedź w 30 sekund: Operatory skip* są przeciwieństwem take* - pomijają wartości zamiast je pobierać. skip(n) pomija pierwsze n wartości, skipUntil(notifier$) pomija wartości aż inny Observable wyemituje, a skipWhile(predicate) pomija wartości dopóki spełniają warunek.

Odpowiedź w 2 minuty: Operatory z rodziny skip* służą do pomijania wartości na początku strumienia i rozpoczynania emisji dopiero po spełnieniu określonych warunków. Są one lustrzanym odbiciem operatorów take*.

skip(count) pomija określoną liczbę pierwszych wartości ze źródłowego Observable, a następnie przepuszcza wszystkie kolejne wartości bez zmian. Jest przydatny, gdy chcemy zignorować początkowe wartości strumienia, na przykład początkowe ładowanie danych czy pierwsze zdarzenia inicjalizacyjne. W przeciwieństwie do take(), który kończy strumień po pobraniu wartości, skip() kontynuuje emisję do naturalnego zakończenia źródła.

skipUntil(notifier$) pomija wszystkie wartości ze źródła aż do momentu, gdy Observable notifier$ wyemituje swoją pierwszą wartość. Od tego momentu wszystkie wartości są przepuszczane. Jest to przydatne w scenariuszach, gdzie chcemy rozpocząć obserwację dopiero po wystąpieniu określonego zdarzenia, na przykład rozpoczęcie śledzenia akcji użytkownika dopiero po załadowaniu danych.

skipWhile(predicate) pomija wartości dopóki spełniają one warunek określony w predykacie. Gdy tylko pojawi się pierwsza wartość nie spełniająca warunku, operator zaczyna przepuszczać wszystkie kolejne wartości (nawet jeśli znów spełniałyby warunek). To odróżnia go od filter(), który sprawdza warunek dla każdej wartości osobno.

Przykład kodu:

import { interval, fromEvent, Subject } from 'rxjs';
import { skip, skipUntil, skipWhile } from 'rxjs';

// 1. skip() - pomiń pierwsze 3 wartości
const numbers$ = interval(1000);

numbers$.pipe(
  skip(3)
).subscribe(num => {
  console.log('skip:', num); // Wynik: 3, 4, 5, 6...
});

// 2. skipUntil() - pomiń wartości aż do pierwszego kliknięcia
const clicks$ = fromEvent(document, 'click');

numbers$.pipe(
  skipUntil(clicks$)
).subscribe(num => {
  console.log('skipUntil po kliknięciu:', num);
});
// Liczby zaczynają być emitowane dopiero po pierwszym kliknięciu

// 3. skipWhile() - pomiń wartości dopóki są mniejsze niż 5
numbers$.pipe(
  skipWhile(num => num < 5)
).subscribe(num => {
  console.log('skipWhile:', num);
  // Wynik: 5, 6, 7, 8... (0,1,2,3,4 zostały pominięte)
});

// 4. Różnica między skipWhile() a filter()
numbers$.pipe(
  skipWhile(num => num < 5) // Pomija TYLKO początkowe < 5
).subscribe(num => console.log('skipWhile:', num));
// Wynik: 5, 6, 7, 8, 9, 10...

numbers$.pipe(
  filter(num => num >= 5) // Filtruje WSZYSTKIE < 5
).subscribe(num => console.log('filter:', num));
// Wynik: 5, 6, 7, 8, 9, 10... (identyczny w tym przypadku)

// 5. Praktyczny przykład - ignoruj pierwsze ładowanie
import { BehaviorSubject } from 'rxjs';

const userData$ = new BehaviorSubject({ name: '', email: '' });

// Pomiń początkową pustą wartość
userData$.pipe(
  skip(1) // Pomija pierwszy emit z pustymi danymi
).subscribe(user => {
  console.log('Zaktualizowany użytkownik:', user);
});

userData$.next({ name: '', email: '' }); // Pominięte
userData$.next({ name: 'Jan', email: 'jan@example.com' }); // Wyemitowane

// 6. skipUntil - rozpocznij śledzenie po załadowaniu danych
const dataLoaded$ = new Subject<void>();
const userActions$ = fromEvent(document, 'click');

userActions$.pipe(
  skipUntil(dataLoaded$)
).subscribe(() => {
  console.log('Akcja użytkownika (dane już załadowane)');
});

// Symulacja ładowania danych
setTimeout(() => {
  console.log('Dane załadowane!');
  dataLoaded$.next();
}, 3000);
// Kliknięcia przed 3 sekundami są ignorowane

// 7. Łączenie skip i take - okno wartości
import { take } from 'rxjs';

numbers$.pipe(
  skip(5),  // Pomiń pierwsze 5
  take(3)   // Weź następne 3
).subscribe({
  next: num => console.log('Okno wartości:', num),
  complete: () => console.log('Zakończono')
});
// Wynik: 5, 6, 7, potem complete

// 8. skipWhile z obiektami
import { of } from 'rxjs';

interface Temperature {
  value: number;
  unit: string;
}

const temperatures$ = of<Temperature>(
  { value: -5, unit: 'C' },
  { value: -2, unit: 'C' },
  { value: 1, unit: 'C' },
  { value: 5, unit: 'C' },
  { value: -1, unit: 'C' }, // Ta wartość też będzie wyemitowana!
  { value: 8, unit: 'C' }
);

temperatures$.pipe(
  skipWhile(temp => temp.value < 0) // Pomija tylko początkowe ujemne
).subscribe(temp => {
  console.log('Temperatura:', temp.value);
  // Wynik: 1, 5, -1, 8 (tylko pierwsze dwie ujemne zostały pominięte)
});

Materiały:

21. 21. Jaka jest różnica między debounceTime() a throttleTime()?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
22. 22. Jak działa operator distinctUntilChanged() i kiedy go używać?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
23. 23. Czym jest operator first() i last()?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi

Operatory Łączące w RxJS

24. Jaka jest różnica między merge() a concat()?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
25. Jak działa operator combineLatest() i kiedy go używać?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
26. Czym jest forkJoin() i jak obsługuje równoległe żądania?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
27. Jak używać operatora zip() do synchronizacji strumieni?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
28. Czym jest operator withLatestFrom()?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
29. Diagramy Mermaid - Przepływ danych w operatorach łączących

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
30. Tabela porównawcza operatorów łączących

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi

Subjects i Multicasting - RxJS

31. Czym jest Subject i czym różni się od zwykłego Observable?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
32. Jaka jest różnica między Subject, BehaviorSubject, ReplaySubject i AsyncSubject?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
33. Jak działa operator share() i shareReplay()?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
34. Czym jest multicasting i dlaczego jest ważny?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi

Obsługa Błędów i Retry w RxJS

35. Jak obsługiwać błędy w strumieniach RxJS za pomocą catchError()?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
36. Jak używać operatorów retry() i retryWhen() do ponawiania operacji?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
37. Czym jest operator finalize() i kiedy go używać?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
38. Jak działa operator throwError() do tworzenia strumieni błędów?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi

Najlepsze Praktyki i Wzorce RxJS

39. Jak unikać memory leaks przy pracy z RxJS w Angular?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
40. Jakie są najlepsze praktyki organizacji kodu z RxJS w dużych aplikacjach?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi

Zaawansowane wzorce i techniki

41. Jak tworzyc wlasne operatory wielokrotnego uzytku w RxJS?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
42. Czym są Schedulery w RxJS i jak wpływają na wykonanie Observable?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
43. Jak testować kod RxJS za pomocą marble testing i TestScheduler?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
44. Jak implementować wzorzec exponential backoff dla retry w RxJS?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
45. Jak obsługiwać backpressure w RxJS gdy producent jest szybszy niż konsument?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
46. Jak używać WebSocket z RxJS do komunikacji w czasie rzeczywistym?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi
47. Jak zarządzać stanem aplikacji za pomocą RxJS bez zewnętrznych bibliotek?

Odpowiedź dostępna w pełnej wersji. Odblokuj 27 pozostałych odpowiedzi i ucz się z pełnego zestawu.

Odblokuj odpowiedzi

Chcesz poznać wszystkie odpowiedzi?

Uzyskaj pełny dostęp do 47 pytań rekrutacyjnych z RxJS oraz pozostałych technologii.

Zobacz plany cenowe