Was ist reaktive Programmierung mit RxJS? RxJS (Reactive Extensions for JavaScript) ist eine Bibliothek für reaktive Programmierung, die Observables verwendet, um asynchrone Ereignissequenzen zu verwalten. In Angular ist es nativ integriert: Jeder HttpClient-Aufruf, jeder FormControl.valueChanges, jedes Router-Ereignis ist ein Observable. Die wahre Stärke von RxJS liegt in seinen Operatoren: reine Funktionen, die Datenströme transformieren, filtern, kombinieren und verwalten. Sie gut zu kennen, macht den Unterschied zwischen einer gut gestalteten reaktiven Anwendung und einem Gewirr verschachtelter Rückrufe aus. 1. map – Datentransformation Analogie: Wie Array.map() , aber in einem kontinuierlichen Strom über die Zeit. Anwendungsfall: Die Antwort einer API normalisieren // Szenario: Die API gibt { data: User[], total: number } // zurück, aber die Komponente erwartet nur User[] this.http.get>('/api/users') .pipe( map(response => Response.data), map(users => users.map(u => ({ ...u, fullName: `${u.firstName} ${u.lastName}`, avatarUrl: u.avatar ||. '/assets/default-avatar.png' }))) ) .subscribe(users => this.users = Benutzer); Faustregel: Verwenden Sie Map immer dann, wenn Sie die Form der Daten ändern möchten, ohne die Ausgaberate zu ändern. 2. Filter – Filterausgaben Übergeben Sie nur Werte, die ein boolesches Prädikat erfüllen. Anwendungsfall: Nur relevante Tastaturereignisse verarbeiten fromEvent(document, 'keydown') .pipe( filter(event => event.key === 'Enter' || event.key === 'Escape'), map(event => event.key) ) .subscribe(key => { if (key === 'Enter') this.confirmAction(); if (key === 'Escape') this.closeModal(); Regel Faustregel: Verwenden Sie den Filter so früh wie möglich in der Kette, um die Arbeit nachfolgender Operatoren zu reduzieren. 3. Tap führt eine Aktion für den aktuellen Wert durch, lässt den Wert jedoch unverändert. Ideal für die Protokollierung und das Laden des Spinners this.http.get('/api/products').
pipe( tap(() => this.loading = true), map(products => products.filter(p => p.inStock)), tap(products => console.log(`Verfügbare Produkte: ${products.length}`)), tap(() => this.loading = false) ) .subscribe(products => this.products = products); Faustregel: Verwenden Sie Tap nicht für Geschäftslogik, sondern nur für Nebeneffekte (Protokoll, Analyse, Spinner). 4. debounceTime + simplyUntilChanged – Suche optimieren Diese beiden Operatoren arbeiten oft paarweise, um zu vermeiden, dass das Backend während der Eingabe durch den Benutzer mit unnötigen Anfragen überflutet wird. debounceTime(ms): wartet N ms darauf, dass der Benutzer mit der Eingabe aufhört, bevor eindeutigUntilChanged() ausgegeben wird: wird nicht ausgegeben, wenn der Wert mit dem vorherigen identisch ist Anwendungsfall: Suchfeld mit API-Aufruf this.searchControl.valueChanges.pipe( debounceTime(400), simplyUntilChanged(), filter(term => term.length >= 2), switchMap(term => this.productService.search(term)), takeUntil(this.destroy$) ).subscribe(results => this.results = results); Ohne diese Operatoren würde jeder Tastendruck eine HTTP-Anfrage generieren. Mit ihnen wird die Anzahl der Anrufe um über 90 % reduziert. 5. switchMap – Vorherige Anfragen löschen switchMap wird am häufigsten für HTTP-Anfragen als Reaktion auf Benutzerereignisse verwendet. Wenn eine neue Ausgabe eintrifft, löscht es automatisch das vorherige Observable und abonniert das neue. Anwendungsfall: dynamische Navigation mit Routenparametern this.route.paramMap.pipe( map(params => params.get('id')), filter(id => !!id), switchMap(id => this.articleService.getById(id!)) ).subscribe(article => this.article = Article); Anwendungsfall: Automatische Vervollständigung mit aktualisierten Ergebnissen // Wenn der Benutzer „ang“, dann „angu“ und dann „angul“ eingibt, // werden nur Anfragen für „angul“ ausgeführt searchTerm$.pipe( debounceTime(300), switchMap(term => this.searchService.getSuggestions(term)) ).subscribe(suggestions => this.suggestions =Suggestions); Regel: Verwenden Sie switchMap, wenn nur die letzte Antwort relevant ist. 6.
mergeMap – Parallele Ausführung Im Gegensatz zu switchMap löscht mergeMap NICHT vorherige Anfragen: Es führt sie alle parallel aus und gibt Ergebnisse aus, sobald sie eintreffen. Anwendungsfall: Hochladen mehrerer Dateien von(selectedFiles).pipe( mergeMap(file => this.uploadService.upload(file).pipe( map(result => ({ file: file.name, url: result.url, status: 'ok' })), CatchError(err => of({ file: file.name, status: 'error', error: err.message })) )) ).subscribe(result => this.uploadResults.push(result)); Achtung: Verwenden Sie mergeMap(fn, maxConcurrency), um gleichzeitige parallele Anforderungen zu begrenzen. 7. concatMap – Garantierte sequentielle Ausführung. concatMap führt Vorgänge nacheinander aus und wartet auf den Abschluss des vorherigen Vorgangs, bevor der nächste gestartet wird. Anwendungsfall: Datensätze nacheinander speichern von([record1, record2, record3]).pipe( concatMap(record => this.db.save(record)) ).subscribe({ next: result => console.log('Saved:', result), complete: () => console.log('Alle Datensätze in der richtigen Reihenfolge gespeichert!') }); Goldene Regel für Operatoren höherer Ordnung: Nur der letzte zählt → switchMap Alle parallel → mergeMap Einer nach dem anderen der Reihe nach → concatMap 8. forkJoin – Auf mehrere parallele Anforderungen warten forkJoin ist das reaktive Gegenstück zu Promise.all(): Es wartet darauf, dass alle Observables abgeschlossen sind, und gibt jeweils ein Objekt mit den neuesten Werten aus. Anwendungsfall: Daten von mehreren Endpunkten für ein Dashboard laden ngOnInit() { forkJoin({ Benutzer: this.userService.getCurrent(), Statistiken: this.statsService.getDashboard(), Benachrichtigungen: this.notifService.getUnread(), Projekte: this.projectService.getAll() }).pipe( takeUntil(this.destroy$) ).subscribe({ next: ({ Benutzer, Statistiken, Benachrichtigungen, Projekte }) => { this.user = user; this.stats = stats;
error = 'Fehler beim Laden des Dashboards' }); } Einschränkung: Wenn eine der Quellen einen Fehler generiert, gibt forkJoin den Fehler weiter und löscht die anderen. Verwenden Sie „catchError“ für jedes Observable, um Teilfehler zu tolerieren. 9. CombineLatest – Live-Streams kombinieren Im Gegensatz zu forkJoin wartet CombineLatest nicht auf den Abschluss: Es berechnet die Ausgabe jedes Mal neu, wenn eine der Quellen einen neuen Wert ausgibt. Anwendungsfall: gefilterte Tabelle mit mehreren unabhängigen Filtern combinLatest([ this.searchControl.valueChanges.pipe(startWith('')), this.categoryControl.valueChanges.pipe(startWith('all')), this.statusControl.valueChanges.pipe(startWith('active')) ]).pipe( debounceTime(200), switchMap(([search, Category, Status]) => this.productService.getFiltered({ Suche, Kategorie, Status }) ), takeUntil(this.destroy$) ).subscribe(products => this.filteredProducts = products); Hinweis: startWith() ist erforderlich, da combinLatest erst dann emittiert, wenn alle Quellen mindestens einmal emittiert haben. 10. takeUntil – Speicherlecks verhindern Jedes nicht verwaltete Abonnement ist ein potenzielles Speicherleck. takeUntil ist das Standardmuster in Angular, um Abonnements automatisch abzuschließen, wenn die Komponente zerstört wird. Anwendungsfall: Standardbereinigungsmuster @Component({ ... }) Exportklasse MyComponent implementiert OnInit, OnDestroy { private destroy$ = new Subject(); ngOnInit() { this.dataService.stream$.pipe( takeUntil(this.destroy$) ).subscribe(data => this.data = data); Interval(5000).pipe( takeUntil(this.destroy$), switchMap(() => this.statsService.refresh()) ).subscribe(stats => this.stats = stats); } ngOnDestroy() { this.destroy$.next(); this.destroy$.complete(); } } Moderne Alternative (Angular 16+): Verwenden Sie takeUntilDestroyed() von @angular/core/rxjs-interop, ohne das Subjekt manuell zu bearbeiten. 11. CatchError – Fehlerbehandlung CatchError fängt Fehler im Stream ab und ermöglicht Ihnen die Wiederherstellung mit einem Observable-Fallback, anstatt das Abonnement zu unterbrechen.
Anwendungsfall: Fallback mit zwischengespeicherten Daten + Benutzerbenachrichtigung this.http.get('/api/products').pipe( CatchError(err => { console.error('[ProductService] API Error:', err); this.toastService.showError('Failed to Load Products. Using Cache Data.'); return of(this.cacheService.getProducts() ?? []); }) ).subscribe(products => this.products = Produkte); Anwendungsfall: Automatischer Wiederholungsversuch bei Netzwerkfehlern this.http.get('/api/critical-data').pipe( retry({ count: 3, delay: 1000 }), CatchError(err => { this.errorService.report(err); return throwError(() => new Error('Service unavailable')); }) ).subscribe(data => this.data = data); 12. shareReplay – Stream-Caching und Teilen shareReplay(1) wandelt ein Cold Observable in ein gemeinsam genutztes Hot Observable mit Zwischenspeicherung des letzten Werts um. Der HTTP-Aufruf wird unabhängig von der Anzahl der Abonnenten nur einmal ausgeführt. Anwendungsfall: Globale Konfiguration wird nur einmal geladen @Injectable({ bereitgestelltIn: 'root' }) export class ConfigService { // HTTP-Aufruf wird NUR EINMAL ausgeführt, // alle Komponenten erhalten die gleiche zwischengespeicherte Antwort readonly config$ = this.http.get('/api/config').pipe( shareReplay(1) ); Konstruktor(private http: HttpClient) {} } // In jeder Komponente keine doppelten Aufrufe @Component({ ... }) export class NavbarComponent { config$ = inject(ConfigService).config$; // Cache verwenden } Zusammenfassung: Wann die einzelnen Operatoren verwendet werden sollen Operator Verwenden Sie, wenn...
Typische Beispielkarte Transformieren Sie die Form der Daten. Antwort-API-Filter normalisieren. Ignorieren Sie irrelevante Werte. Filtern Sie Tastaturtipp-Ereignisse. Nebeneffekt, ohne den Fluss zu ändern. Protokollierung, Spinner-DebounceTime. Warten Sie auf eine Pause in der Eingabe Mehrere Live-Streams kombiniert, mehrere Filter, takeUntil, automatisches Bereinigungsabonnement, ngOnDestroy-Muster, CatchError, Behandeln und Wiederherstellen von Fehlern, Wiederholen + Fallback, shareReplay, Teilen und Zwischenspeichern von Streams, Global Config. Die Beherrschung dieser Operatoren bedeutet, Angular-Code zu schreiben, der besser lesbar, effizienter und praktisch frei von asynchronen Fehlern ist. Das Geheimnis besteht darin, immer von der Frage auszugehen: „Was soll passieren, wenn eine neue Emission eintrifft, während die vorherige noch im Flug ist?“ – Die Antwort weist Sie fast immer auf den richtigen Betreiber hin.