¿Qué es la programación reactiva con RxJS? RxJS (Extensiones reactivas para JavaScript) es una biblioteca para programación reactiva que utiliza Observables para gestionar secuencias asincrónicas de eventos. En Angular está integrado de forma nativa: cada llamada a HttpClient, cada FormControl.valueChanges, cada evento de Router es un Observable. El verdadero poder de RxJS reside en sus operadores: funciones puras que transforman, filtran, combinan y gestionan flujos de datos. Conocerlos bien marca la diferencia entre una aplicación reactiva bien diseñada y una maraña de devoluciones de llamadas anidadas. 1. map: transformación de datos Analogía: como Array.map(), pero en un flujo continuo a lo largo del tiempo. Caso de uso: Normalizar la respuesta de una API // Escenario: La API devuelve { data: Usuario[], total: número } // pero el componente solo espera Usuario[] this.http.get>('/api/users') .pipe( map(response => respuesta.data), map(users => users.map(u => ({ ...u, fullName: `${u.firstName} ${u.lastName}`, avatarUrl: u.avatar || '/assets/default-avatar.png' }))) ) .subscribe(usuarios => this.users = usuarios); Regla general: utilice map siempre que desee transformar la forma de los datos sin cambiar la tasa de salida. 2. filtrar: filtrar salidas. Solo pasa valores que satisfagan un predicado booleano. Caso de uso: procesar solo eventos de teclado relevantes de Event(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(); Regla general: use el filtro lo más temprano posible en la cadena para reducir el trabajo de los operadores posteriores. 3. toque: efectos secundarios sin cambiar, toque realiza una acción en el valor actual pero deja que el valor pase sin cambios. Ideal para registrar y actualizar el estado de la interfaz de usuario.
pipe( tap(() => this.loading = true), map(productos => productos.filter(p => p.inStock)), tap(products => console.log(`Productos disponibles: ${products.length}`)), tap(() => this.loading = false) ) .subscribe(products => this.products = productos); Regla general: no utilice tap para la lógica empresarial, solo para efectos secundarios (registro, análisis, control giratorio). 4. debounceTime + distintivoUntilChanged: optimización de la búsqueda. Estos dos operadores a menudo trabajan en pares para evitar inundar el backend con solicitudes innecesarias a medida que el usuario escribe. debounceTime(ms): espera a que el usuario deje de escribir durante N ms antes de emitir distintivoUntilChanged(): no genera resultados si el valor es idéntico al anterior. Caso de uso: cuadro de búsqueda con API llama a this.searchControl.valueChanges.pipe( debounceTime(400), distintivoUntilChanged(), filter(term => term.length >= 2), switchMap(term => this.productService.search(term)), takeUntil(this.destroy$) ).subscribe(resultados => this.results = resultados); Sin estos operadores, cada pulsación de tecla generaría una solicitud HTTP. Con ellos el número de llamadas se reduce en más del 90%. 5. switchMap: borrar solicitudes anteriores switchMap es el más utilizado para solicitudes HTTP en respuesta a eventos de usuario. Cuando llega un nuevo número, elimina automáticamente el Observable anterior y se suscribe al nuevo. Caso de uso: navegación dinámica con parámetros de ruta this.route.paramMap.pipe( map(params => params.get('id')), filter(id => !!id), switchMap(id => this.articleService.getById(id!)) ).subscribe(article => this.article = artículo); Caso de uso: autocompletar con resultados actualizados // Si el usuario escribe "ang", luego "angu" y luego "angul", // solo se ejecutan las solicitudes de "angul" searchTerm$.pipe( debounceTime(300), switchMap(term => this.searchService.getSuggestions(term)) ).subscribe(suggestions => this.suggestions = sugerencias); Regla: use switchMap cuando solo la última respuesta sea relevante. 6.
mergeMap: ejecución paralela A diferencia de switchMap, mergeMap NO elimina solicitudes anteriores: las ejecuta todas en paralelo y genera los resultados a medida que llegan. Caso de uso: carga de varios archivos desde(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(resultado)); Precaución: utilice mergeMap(fn, maxConcurrency) para limitar las solicitudes paralelas simultáneas. 7. concatMap: ejecución secuencial garantizada concatMap ejecuta operaciones una a la vez en orden, esperando a que se complete la anterior antes de comenzar la siguiente. Caso de uso: guardar registros secuencialmente desde([record1, record2, record3]).pipe( concatMap(record => this.db.save(record)) ).subscribe({ next: result => console.log('Saved:', result), complete: () => console.log('¡Todos los registros se guardaron en orden!') }); Regla de oro de los operadores de orden superior: solo cuenta el último → switchMap Todos en paralelo → mergeMap Uno a la vez en orden → concatMap 8. forkJoin: espera múltiples solicitudes paralelas forkJoin es la contraparte reactiva de Promise.all(): espera a que se completen todos los Observables y emite un objeto con los últimos valores de cada uno. Caso de uso: cargar datos desde múltiples puntos finales para un panel ngOnInit() { forkJoin({ usuario: this.userService.getCurrent(), estadísticas: this.statsService.getDashboard(), notificaciones: this.notifService.getUnread(), proyectos: this.projectService.getAll() }).pipe( takeUntil(this.destroy$) ).subscribe({ next: ({ usuario, estadísticas, notificaciones, proyectos }) => { this.user = usuario; this.stats = stats; this.notificaciones = this.isReady;
error = 'Error al cargar el panel' }); } Limitación: si una de las fuentes genera un error, forkJoin propaga el error y elimina las demás. Utilice catchError en cada Observable para tolerar fallas parciales. 9. combineLatest: combinación de transmisiones en vivo A diferencia de forkJoin, combineLatest no espera a que finalice: recalcula la salida cada vez que una de las fuentes emite un nuevo valor. Caso de uso: tabla filtrada con múltiples filtros independientes combineLatest([ this.searchControl.valueChanges.pipe(startWith('')), this.categoryControl.valueChanges.pipe(startWith('all')), this.statusControl.valueChanges.pipe(startWith('active')) ]).pipe( debounceTime(200), switchMap(([búsqueda, categoría, estado]) => this.productService.getFiltered({ búsqueda, categoría, estado }) ), takeUntil(this.destroy$) ).subscribe(products => this.filteredProducts = productos); Nota: startWith() es necesario porque combineLatest no emite hasta que todas las fuentes hayan emitido al menos una vez. 10. takeUntil: evita pérdidas de memoria Cada suscripción no administrada es una posible pérdida de memoria. takeUntil es el patrón estándar en Angular para completar automáticamente las suscripciones cuando se destruye el componente. Caso de uso: patrón de limpieza estándar @Component({ ... }) clase de exportación MyComponent implementa OnInit, OnDestroy { private destroy$ = new Subject(); ngOnInit() { this.dataService.stream$.pipe( takeUntil(this.destroy$) ).subscribe(datos => this.data = datos); intervalo(5000).pipe( takeUntil(this.destroy$), switchMap(() => this.statsService.refresh()) ).subscribe(stats => this.stats = stats); } ngOnDestroy() { this.destroy$.siguiente(); this.destroy$.complete(); } } Alternativa moderna (Angular 16+): use takeUntilDestroyed() de @angular/core/rxjs-interop sin manipular manualmente el Asunto. 11. catchError: manejo de errores. catchError intercepta errores en la transmisión y le permite recuperarse con un respaldo observable en lugar de interrumpir la suscripción.
Caso de uso: respaldo con datos almacenados en caché + notificación de usuario this.http.get('/api/products').pipe( catchError(err => { console.error('[ProductService] Error de API:', err); this.toastService.showError('Error al cargar los productos. Usando datos de caché.'); return of(this.cacheService.getProducts() ?? []); }) ).subscribe(products => this.products = productos); Caso de uso: reintento automático para errores de red this.http.get('/api/critical-data').pipe( retry({ count: 3, delay: 1000 }), catchError(err => { this.errorService.report(err); return throwError(() => new Error('Servicio no disponible')); }) ).subscribe(data => this.data = data); 12. shareReplay: almacenamiento en caché de transmisiones y uso compartido shareReplay(1) transforma un Observable frío en un Observable caliente compartido, con almacenamiento en caché del último valor. La llamada HTTP se ejecuta solo una vez independientemente del número de suscriptores. Caso de uso: configuración global cargada solo una vez @Injectable({ provideIn: 'root' }) export class ConfigService { // La llamada HTTP se ejecuta SÓLO UNA VEZ, // todos los componentes reciben la misma respuesta en caché de solo lectura config$ = this.http.get('/api/config').pipe( shareReplay(1) ); constructor(private http: HttpClient) {} } // En cualquier componente, cero llamadas duplicadas @Component({ ... }) export class NavbarComponent { config$ = inject(ConfigService).config$; // usar caché } Resumen: Cuándo usar cada operador Operador Usar cuando...
Mapa de ejemplo típico Transformar la forma de los datos Normalizar la respuesta Filtro API Ignorar valores irrelevantes Filtrar eventos de pulsación del teclado Efecto secundario sin modificar el flujo Registro, debounceTime Esperar pausa en la entrada Cuadro de búsqueda distintivoUntilChanged Evitar emisiones duplicadas Controles de formulario switchMap Solo cuenta la última solicitud Autocompletar, parámetros de ruta mergeMap Solicitudes paralelas sin orden Carga múltiple concatMap Orden de ejecución garantizada Operaciones secuenciales forkJoin Esperar N solicitudes, todas terminadas Dashboard init combineLatest Multiple transmisiones en vivo combinadas Múltiples filtros tomarHasta Suscripción de limpieza automática Patrón ngOnDestroy catchError Manejar y recuperarse de errores Reintentar + compartir alternativaReproducir Compartir y almacenar en caché transmisiones Configuración global Dominar estos operadores significa escribir código Angular que sea más legible, más eficiente y prácticamente libre de errores asincrónicos. El secreto está en partir siempre de la pregunta: "¿Qué debería pasar cuando llegue una nueva emisión mientras la anterior todavía está en vuelo?" — la respuesta casi siempre indica el operador correcto.