wandres.dev
OBSERVABLES · RxJS y streams

Operadores: transformar y combinar streams

El álgebra de los observables: map y filter para moldear un stream, switchMap para aplanar streams de streams con cancelación, debounceTime para domar el tiempo y combineLatest para fusionar fuentes en un único flujo derivado.

⏱ 18 min

Si el observable es el sustantivo de la programación reactiva, los operadores son sus verbos. Un operador toma uno o más streams y devuelve un stream nuevo, sin mutar los originales: es una función pura sobre el tiempo. Con un puñado bien entendido —map, filter, switchMap, debounceTime, combineLatest— puedes expresar coreografías asíncronas que en código imperativo serían un nido de banderas, temporizadores y condiciones de carrera. Aprender operadores es aprender a pensar en transformaciones de flujo en lugar de en pasos.

🎯 Al terminar esta lección sabrás
  • Moldear un stream con los operadores de transformación map y filter.
  • Aplanar streams de streams y cancelar con switchMap.
  • Domar la dimensión temporal con debounceTime.
  • Fusionar varias fuentes en un flujo derivado con combineLatest.

pipe y los operadores puros

Desde RxJS 5.5 los operadores son funciones sueltas que se encadenan con pipe, no métodos del observable. Esto habilita el tree-shaking —solo empaquetas lo que usas— y deja la composición legible de arriba abajo. Técnicamente, cada operador es una función de orden superior: recibe configuración y devuelve una OperatorFunction que va de un observable a otro. Cada eslabón recibe el stream de salida del anterior:

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

const pares$ = numeros$.pipe(
  filter((n) => n % 2 === 0),       // deja pasar solo los pares
  tap((n) => console.log('paso', n)), // efecto lateral, no transforma
  map((n) => n * 10),               // multiplica cada uno por diez
);

filter es un guardián: propaga al siguiente next solo los valores que cumplen el predicado. map es un transformador: aplica una función a cada valor y emite el resultado. tap es la excepción útil: no cambia el stream, solo permite un efecto lateral (un log, una métrica) sin romper la cadena. Ninguno toca numeros$; todos devuelven un stream nuevo. Esta pureza es lo que permite razonar cada línea del pipe de forma aislada, como se razona una tubería de Unix.

El centenar de operadores se organiza en pocas familias que conviene tener en la cabeza: de creación (of, from, interval), de transformación (map, scan), de filtrado (filter, take, distinctUntilChanged), de combinación (combineLatest, merge), de manejo de error (catchError, retry) y de utilidad (tap, finalize). Aprender la categoría antes que el nombre concreto acelera enormemente encontrar el operador que buscas.

switchMap: streams de streams

El operador que separa a quien entiende RxJS de quien lo copia es switchMap. Aparece cuando cada valor de un stream debe disparar otro stream —típicamente una petición de red— y quieres que solo cuente el más reciente. switchMap proyecta cada valor a un observable interno y cancela el interno anterior en cuanto llega uno nuevo:

flowchart LR
A[texto: ne] --> P[peticion A]
B[texto: neo] --> Q[peticion B cancela A]
C[texto: neov] --> R[peticion C cancela B]
R --> S[unico resultado emitido]
style S fill:#89b4fa,color:#11111b

Esa cancelación resuelve gratis la condición de carrera más común de las interfaces: el usuario teclea rápido, se lanzan varias búsquedas y la lenta responde después de la rápida, pisando el resultado correcto con uno viejo. Con switchMap la petición anterior se aborta y jamás emite. En código, el anidamiento desaparece:

import { switchMap, map } from 'rxjs';

const resultados$ = consulta$.pipe(
  switchMap((texto) => buscar(texto)),   // buscar devuelve un Observable
  map((res) => res.items),
);

Sus tres hermanos completan el cuadro: mergeMap no cancela y deja correr todo en paralelo; concatMap encola y respeta el orden; exhaustMap ignora nuevos valores mientras el interno sigue vivo. Elegir entre los cuatro es elegir una política de concurrencia.

💡
La regla mnemotécnica de los cuatro map de orden superior

switchMap para búsquedas y navegación: solo importa lo último. concatMap para escrituras que deben ordenarse: guardar cambios en secuencia. mergeMap para acciones independientes que pueden solaparse: subir varios ficheros a la vez. exhaustMap para acciones que no deben repetirse mientras corren: el botón de login que no admite doble clic. Escoger mal aquí es la causa raíz de la mayoría de bugs sutiles en código RxJS de producción.

debounceTime: domar el tiempo

Los operadores temporales son el terreno donde RxJS no tiene rival. debounceTime(300) retiene cada valor y solo lo emite si pasan 300 ms sin que llegue otro; ideal para no lanzar una búsqueda en cada tecla. Sus vecinos afinan el matiz: throttleTime limita la frecuencia dejando pasar uno cada intervalo, auditTime emite el último de cada ventana, y distinctUntilChanged descarta valores repetidos consecutivos. Juntos convierten un torrente ruidoso en una señal útil:

import { debounceTime, distinctUntilChanged, filter, map } from 'rxjs';

const consulta$ = teclas$.pipe(
  map((e) => (e.target as HTMLInputElement).value),
  debounceTime(300),                      // espera a que el usuario pare
  distinctUntilChanged(),                 // ignora si el texto no cambio
  filter((texto) => texto.length >= 2),   // no busques con una letra
);

Lo notable es que la lógica temporal —“espera a que pare de teclear”— queda declarada como un eslabón más de la tubería, no dispersa en setTimeout, variables de último-id y clearTimeout. El tiempo se vuelve un operador, y como tal se lee, se testea y se recombina.

La distinción entre los temporales es sutil pero decisiva. debounceTime espera al silencio y descarta lo intermedio; throttleTime deja pasar a un ritmo máximo constante. El buscador quiere lo primero; un manejador de scroll o de resize quiere lo segundo, para muestrear sin ahogar el hilo principal:

import { throttleTime, map } from 'rxjs';

const scrollY$ = fromEvent(window, 'scroll').pipe(
  throttleTime(100),                 // como mucho una lectura cada 100 ms
  map(() => window.scrollY),
);

combineLatest: fusionar fuentes

Hasta aquí hemos transformado un stream. combineLatest hace lo contrario: toma varios y produce uno que emite cada vez que cualquiera cambia, entregando el último valor de todos. Es la herramienta del estado derivado a partir de fuentes independientes:

import { combineLatest, map } from 'rxjs';

// filtro y orden son dos controles independientes de la UI
const listaVisible$ = combineLatest([datos$, filtro$, orden$]).pipe(
  map(([datos, filtro, orden]) => aplicar(datos, filtro, orden)),
);

listaVisible$ es una función pura de tres entradas reactivas: cambia el filtro y la lista se recalcula; cambian los datos y se recalcula igual. Nunca queda desincronizada porque no hay un paso manual que actualizar. Sus parientes merge (intercala emisiones sin combinar) y withLatestFrom (emite solo cuando lo hace el principal, tomando de reojo el último del secundario) cubren las demás formas de unir flujos.

Un operador merece mención aparte por su vínculo directo con el estado: scan, el reduce del tiempo. Acumula un valor a lo largo de las emisiones y emite el acumulado en cada paso; es la forma en que un stream de eventos se convierte en estado derivado sin salir del pipe. Rematado con shareReplay, ese estado se comparte entre suscriptores y se recuerda para quien llegue tarde:

import { scan, shareReplay } from 'rxjs';

const total$ = importes$.pipe(
  scan((acc, x) => acc + x, 0),   // acumula: estado a partir de eventos
  shareReplay(1),                 // comparte y recuerda el ultimo valor
);
⚠️
combineLatest no emite hasta que todos han emitido

La trampa clásica: combineLatest espera a tener un valor de cada fuente antes de su primera emisión. Si uno de los streams es perezoso y nunca ha emitido, el combinado permanece mudo y parece roto. La cura es dar valores iniciales —un startWith en las fuentes que puedan tardar, o un BehaviorSubject como origen—. Diagnóstico rápido: si tu combineLatest no emite nunca, busca cuál de sus entradas sigue en silencio.

Los operadores son un álgebra cerrada sobre el tiempo

La idea que hay que llevarse no es la lista de operadores —hay más de cien— sino su naturaleza algebraica: cada operador es una función pura que va de observables a observables, y como el conjunto es cerrado bajo composición, cualquier coreografía asíncrona, por barroca que sea, se expresa encadenando piezas que ya sabes razonar por separado. Esto es exactamente lo que un Array te da en el espacio con map, filter y reduce; RxJS lo eleva a la dimensión del tiempo. La consecuencia práctica es profunda: un problema que en estilo imperativo exige coordinar callbacks, banderas de estado mutable y limpieza manual de temporizadores —el caldo de cultivo de las condiciones de carrera— se convierte en una expresión declarativa donde la cancelación, el orden y la sincronización están garantizados por el operador, no por tu vigilancia. Dominar switchMap frente a sus hermanos, o saber cuándo combineLatest en vez de withLatestFrom, no es memorizar API: es haber interiorizado que estás componiendo transformaciones de flujo. Quien piensa así deja de preguntarse “cuándo llega el dato” y empieza a declarar “de qué es función este dato”. Ese es el salto de operar observables a diseñar con ellos.

⚔️ Construye un buscador reactivo entero
  1. Parte de un fromEvent sobre un input y extrae el texto con map.
  2. Encadena debounceTime, distinctUntilChanged y un filter de longitud mínima.
  3. Usa switchMap para lanzar la petición y demuestra, con red lenta simulada, que las respuestas viejas no pisan a las nuevas.
  4. Cambia switchMap por mergeMap y observa la condición de carrera reaparecer: explica exactamente por qué y en qué caso mergeMap sí sería correcto.