11. Nest avanzado: caché, colas, tiempo real, GraphQL y microservicios
Un CRUD correcto se escribe con lo del capítulo 10. Una aplicación que aguanta tráfico real, que no bloquea la petición del usuario para enviar un correo, que notifica cambios al instante y que sobrevive a tres réplicas detrás de un balanceador necesita otra caja de herramientas. Este capítulo recorre esa caja pieza a pieza: caché, colas de trabajos, tareas programadas, eventos, limitación de peticiones, WebSockets, GraphQL, microservicios, CQRS, archivos y sondas de salud. En cada una interesa lo mismo: qué problema resuelve, qué problema nuevo introduce y cuándo no usarla.
11.1 Qué vas a poder hacer al terminar
- Decidir qué se cachea, dónde se cachea y cuándo se invalida, y razonar por qué la invalidación es la parte difícil y no la lectura.
- Montar caché en memoria y con Redis usando
@nestjs/cache-manager, y explicar por qué elCacheInterceptorautomático es peligroso en endpoints que dependen del usuario. - Sacar de la petición HTTP todo el trabajo que no necesita estar ahí: colas BullMQ con reintentos, retroceso exponencial, deduplicación y consumidores idempotentes.
- Programar tareas con
@nestjs/scheduley resolver el problema clásico del cron que se ejecuta N veces cuando hay N réplicas. - Distinguir un evento interno de un mensaje en una cola real, y elegir con criterio en lugar de por costumbre.
- Limitar peticiones por IP y por usuario con almacenamiento compartido, funcionando correctamente detrás de un proxy inverso.
- Construir un gateway de WebSockets autenticado en el handshake, con salas, difusión selectiva y adaptador de Redis para escalar horizontalmente; y saber cuándo bastan los Server-Sent Events.
- Exponer una API GraphQL code-first sin caer en el problema N+1, con DataLoader y con límites de profundidad y complejidad.
- Argumentar con datos por qué todavía no necesitas microservicios, y si los necesitas, elegir transporte y aplicar saga, outbox, idempotencia y circuit breaker.
- Implementar subida y descarga de archivos con validación real del contenido y URLs prefirmadas.
- Publicar sondas de liveness y readiness con Terminus y entender qué hace Kubernetes con cada una.
11.2 Caché: el arte de no volver a calcular lo mismo
11.2.1 Por qué se cachea, y qué se puede cachear
Cachear es guardar el resultado de una operación costosa para reutilizarlo en peticiones posteriores. El motivo no es solo la latencia: es que el trabajo que no haces no consume CPU, ni conexiones a la base de datos, ni ancho de banda. Una consulta que tarda 300 ms y se repite 500 veces por minuto no es un problema de 300 ms, es un problema de 2,5 minutos de CPU de base de datos por minuto de reloj, es decir, un sistema que ya no da abasto.
| Buen candidato a caché | Mal candidato a caché |
|---|---|
| Datos que se leen mucho más de lo que se escriben (catálogos, configuración, tarifas, países) | Datos que cambian en cada lectura (saldo en un sistema de pagos, stock durante una venta flash) |
| Cálculos deterministas y caros (agregados, informes, renders, conversiones) | Respuestas que dependen de permisos finos y variables por usuario, si no puedes aislar la clave |
| Respuestas de terceros con cuota limitada (geocodificación, tipos de cambio) | Datos personales o secretos, salvo con cifrado y control de acceso explícito |
| Resultados tolerantes a estar «un poco desactualizados» (contadores, rankings, estadísticas) | Operaciones con efectos secundarios: nunca se cachea un POST que crea algo |
La pregunta previa a cualquier caché no es técnica, es de producto: ¿cuántos segundos de desactualización tolera este dato? Si la respuesta es «cero», no hay caché posible y hay que optimizar la consulta (índices, capítulo 17). Si la respuesta es «treinta segundos», ya tienes tu TTL y todo lo demás es implementación.
11.2.2 Los cinco niveles de caché
Antes de escribir una línea de código conviene situarse: hay caché en muchos sitios y la más barata siempre es la que está más cerca del usuario.
NAVEGADOR CDN / EDGE APLICACIÓN (Node) DISTRIBUIDA BASE DE DATOS
┌──────────┐ ┌────────────┐ ┌──────────────────┐ ┌─────────────┐ ┌──────────────┐
│ HTTP │ │ Cloudflare │ │ Map en memoria │ │ Redis │ │ buffer pool │
│ Cache- │──miss────► │ Fastly │─miss──►│ del proceso │─► │ Valkey │─► │ plan cache │
│ Control │ │ Varnish │ │ (LRU + TTL) │ │ Memcached │ │ vistas mat. │
│ ETag │ └────────────┘ └──────────────────┘ └─────────────┘ └──────────────┘
└──────────┘
0 ms ~5-30 ms ~0,01 ms ~0,5-2 ms ~5-300 ms
gratis por petición por proceso compartida última parada
── Coste de un fallo de caché ────────────────────────────────────────────────────────────────────►
◄─ Coste de una invalidación incorrecta (datos obsoletos servidos a más usuarios) ──────────────────
REGLA: cachea lo más arriba que la corrección te permita, no lo más arriba posible.
- Navegador. Cabeceras
Cache-Control,ETagyLast-Modified. Es la caché más eficiente porque ahorra incluso la conexión, pero no la controlas: una vez enviada una respuesta conmax-age=3600no hay forma de retirarla de los navegadores. - CDN o proxy inverso. Ideal para recursos estáticos y para respuestas públicas idénticas para
todos. Se invalida con purgas explícitas. Cuidado: si la respuesta depende del usuario y no envías
Cache-Control: privateoVary: Authorization, la CDN puede servir la respuesta de un usuario a otro. Es un incidente de seguridad clásico. - Aplicación, en memoria del proceso. Rapidísima y trivial de montar. Su límite es estructural: con N réplicas tienes N cachés incoherentes y cada despliegue las vacía. Sirve para datos de configuración inmutables y para reducir «ráfagas» dentro de una misma petición.
- Distribuida (Redis). Una única fuente de verdad para todas las réplicas, con TTL, invalidación centralizada y estructuras de datos útiles (contadores, conjuntos, pub/sub). Es la opción por defecto en producción. Coste: una dependencia más que puede caerse, y latencia de red.
- Base de datos. Tiene su propia caché de páginas y de planes, y además vistas materializadas. Muchas veces «necesito Redis» significa en realidad «me falta un índice». Mide antes de añadir capas.
11.2.3 Invalidación, TTL y estampida: los tres problemas difíciles
Hay una cita atribuida a Phil Karlton que resume tres décadas de dolor: «solo hay dos cosas difíciles en informática: la invalidación de caché y poner nombres a las cosas». Veamos las estrategias reales.
| Estrategia | Cómo funciona | Ventaja | Riesgo |
|---|---|---|---|
| TTL (caducidad) | La entrada expira sola pasado un tiempo | Trivial, sin código de invalidación, autolimpiante | Ventana de datos obsoletos igual al TTL |
| Invalidación explícita | Al escribir, se borra la clave afectada | Coherencia casi inmediata | Hay que acordarse de todas las claves derivadas; un olvido son datos obsoletos eternos |
| Write-through | Al escribir, se actualiza caché y base de datos | La caché nunca está vacía | Escrituras más lentas y necesidad de transaccionalidad entre dos sistemas |
| Clave versionada | La clave incluye una versión o marca de tiempo: tareas:v7:... |
Invalidar es incrementar un contador; no hace falta enumerar claves | Las entradas viejas siguen ocupando memoria hasta que caducan |
| Etiquetas | Cada entrada se asocia a etiquetas y se invalida por etiqueta | Invalidación por grupos lógicos | Requiere estructura adicional; no viene de serie en cache-manager |
Recomendación profesional: combina siempre TTL corto como red de seguridad con invalidación explícita. El TTL garantiza que un olvido tuyo dure minutos y no meses; la invalidación explícita da la frescura que el usuario espera. Nunca dependas solo de la invalidación explícita: es código, y el código tiene huecos.
La estampida de caché (cache stampede o thundering herd)
Una clave muy solicitada caduca. En ese instante llegan 400 peticiones concurrentes, las 400 encuentran la caché vacía y las 400 lanzan la misma consulta costosa a la base de datos. La caché, que existía para proteger la base de datos, acaba de convertirse en el detonante de su caída.
SIN PROTECCIÓN CON BLOQUEO (single-flight)
t=0 clave caduca t=0 clave caduca
┌── req 1 ──► BD (300 ms) ┌── req 1 ──► adquiere lock ──► BD
├── req 2 ──► BD (300 ms) ├── req 2 ──► lock ocupado ──► espera/valor viejo
├── req 3 ──► BD (300 ms) ├── req 3 ──► lock ocupado ──► espera/valor viejo
│ ... 400 consultas idénticas │ ...
└── req 400 ► BD (300 ms) └── req 400 ► lock ocupado ──► espera/valor viejo
t=300ms req 1 escribe caché y libera
RESULTADO: BD saturada, latencia x20 RESULTADO: 1 consulta, latencia normal
Mitigaciones, de la más simple a la más robusta:
- TTL con jitter. En lugar de 300 s exactos para todas las entradas, usa
300 s ± 10 %aleatorio. No resuelve la estampida de una clave única, pero evita que miles de claves caduquen en el mismo segundo (típico tras un despliegue que precalienta la caché). - Single-flight en proceso. Guarda en un
Mapla promesa en curso por clave: la segunda petición concurrente reutiliza la promesa de la primera en lugar de lanzar otra consulta. Elimina la estampida dentro de cada réplica, que suele ser el 90 % del problema. - Bloqueo distribuido. Un
SET clave:lock valor NX PX 5000en Redis: solo quien adquiere el lock recalcula; el resto espera brevemente o devuelve el valor antiguo. - Refresco proactivo (stale-while-revalidate). Guarda dos tiempos: uno «fresco» y otro «tolerable». Si el dato está entre ambos, devuélvelo inmediatamente y lanza el recálculo en segundo plano (una cola, sección 11.3). El usuario nunca espera y la base de datos nunca recibe una avalancha.
- Precalentamiento. Un trabajo programado que rellena las claves críticas antes de que caduquen.
/**
* Colapsa llamadas concurrentes con la misma clave en una sola ejecución.
* Ámbito: el proceso actual. Con varias réplicas necesitas además un lock en Redis.
*/
export class SingleFlight {
private readonly enVuelo = new Map<string, Promise<unknown>>();
async run<T>(clave: string, fn: () => Promise<T>): Promise<T> {
const existente = this.enVuelo.get(clave);
if (existente) return existente as Promise<T>;
// Importante: guardamos la promesa ANTES del primer await,
// si no, dos llamadas simultáneas seguirían pasando de largo.
const promesa = fn().finally(() => this.enVuelo.delete(clave));
this.enVuelo.set(clave, promesa);
return promesa;
}
}
11.2.4 @nestjs/cache-manager: configuración en memoria y con Redis
Nest no implementa su propia caché: envuelve la librería cache-manager, que abstrae el
almacén subyacente. Eso permite empezar en memoria y pasar a Redis cambiando la configuración del módulo, sin
tocar los servicios.
Antes de copiar código de un tutorial, comprueba la versión que tienes instalada
(npm ls @nestjs/cache-manager cache-manager), porque hay tres diferencias que rompen el
comportamiento en silencio:
1. Unidad del TTL. En cache-manager v4 y anteriores el TTL se expresaba en
segundos. Desde v5 se expresa en milisegundos. Un ttl: 60 heredado de un tutorial
antiguo no cachea un minuto: cachea 60 milisegundos, y parece que «la caché no funciona».
2. Almacenes. @nestjs/cache-manager v2 (con cache-manager v5) usaba la opción
store y paquetes como cache-manager-redis-yet.
@nestjs/cache-manager v3 (con cache-manager v6) se apoya en Keyv y usa la opción
stores (en plural, un array que permite varios niveles), con @keyv/redis.
3. Borrado total. El método reset() de cache-manager v5 desapareció en v6; el
equivalente es clear() sobre la instancia. Si tu código llama a reset() tras
actualizar, fallará en tiempo de ejecución y no en compilación si el tipo es laxo.
En este capítulo se muestran las dos formas, etiquetadas. Consulta siempre la documentación de tu versión exacta en lugar de fiarte de la memoria (o de la mía).
import { Module } from '@nestjs/common';
import { CacheModule } from '@nestjs/cache-manager';
@Module({
imports: [
CacheModule.register({
isGlobal: true, // evita importar CacheModule en cada módulo
ttl: 30_000, // ms en cache-manager v5+ (¡segundos en v4!)
max: 500, // número máximo de elementos en memoria: SIEMPRE ponlo
}),
],
})
export class AppModule {}
max es una fuga de memoria con otro nombre
Si cacheas por clave dinámica (por ejemplo, por identificador de usuario o por combinación de filtros) sin
límite de tamaño, el proceso crecerá hasta que el heap de Node reviente con
JavaScript heap out of memory. El almacén en memoria aplica una política LRU cuando le das un
max; sin él, no hay nada que expulse entradas antes de su TTL.
import { CacheModule } from '@nestjs/cache-manager';
import { redisStore } from 'cache-manager-redis-yet';
import { ConfigModule, ConfigService } from '@nestjs/config';
CacheModule.registerAsync({
isGlobal: true,
imports: [ConfigModule],
inject: [ConfigService],
useFactory: async (config: ConfigService) => ({
store: await redisStore({
url: config.getOrThrow<string>('REDIS_URL'),
ttl: 30_000,
}),
}),
});
import { CacheModule } from '@nestjs/cache-manager';
import { createKeyv } from '@keyv/redis';
import { ConfigModule, ConfigService } from '@nestjs/config';
CacheModule.registerAsync({
isGlobal: true,
imports: [ConfigModule],
inject: [ConfigService],
useFactory: (config: ConfigService) => ({
// 'stores' en plural: admite varios niveles en cascada
stores: [createKeyv(config.getOrThrow<string>('REDIS_URL'))],
ttl: 30_000,
}),
});
stores en plural puedes poner primero un almacén en memoria muy pequeño y con
TTL de pocos segundos, y detrás Redis. Las claves más calientes se resuelven sin salir del proceso y Redis
absorbe el resto. Es un patrón excelente para catálogos, pero recuerda que el nivel en memoria no se
puede invalidar desde otra réplica: mantén su TTL en el orden de segundos para que la incoherencia sea
irrelevante.
11.2.5 CacheInterceptor automático y sus limitaciones
Nest incluye un interceptor que cachea respuestas sin escribir código. Es cómodo y, mal usado, es la vía más rápida a un incidente de privacidad.
import { Controller, Get, UseInterceptors } from '@nestjs/common';
import { CacheInterceptor, CacheKey, CacheTTL } from '@nestjs/cache-manager';
@Controller('paises')
@UseInterceptors(CacheInterceptor) // aplicado a todo el controlador
export class PaisesController {
@Get()
@CacheKey('paises:todos') // clave fija en lugar de la URL
@CacheTTL(24 * 60 * 60 * 1000) // 24 h en ms
findAll() {
return this.service.findAll(); // catálogo público e idéntico para todos: caso ideal
}
}
Lo que hace exactamente el interceptor y por qué importa cada detalle:
- Solo cachea peticiones
GETen el transporte HTTP. UnPOSTnunca se cachea, lo cual es correcto y deliberado. - La clave por defecto es la URL de la petición (ruta más query string). No incluye cabeceras, ni cookies, ni el usuario autenticado. Esto es el origen del error más grave de esta sección.
- Cachea la respuesta ya serializada por el controlador, es decir, después de los interceptores anteriores pero sin volver a pasar por la lógica de negocio. Si tu autorización vive en un guard, el guard sí se ejecuta (los guards van antes de los interceptores); pero si el filtrado de datos por permisos ocurre en el servicio, ese filtrado se salta.
- No sabe nada de tu dominio, así que no puede invalidar nada. Solo TTL.
- Para adaptarlo, se extiende y se sobrescribe
trackBy(context), que devuelve la clave oundefinedpara no cachear.
// GET /tareas devuelve SOLO las tareas del usuario autenticado.
// La clave por defecto es la URL: '/tareas' para todos.
@Controller('tareas')
@UseInterceptors(CacheInterceptor)
export class TareasController {
@Get()
@CacheTTL(60_000)
findMine(@CurrentUser() user: JwtUser) {
return this.service.findByOwner(user.id);
}
}
// Ana entra primero → se cachea su lista bajo '/tareas'
// Luis entra después → recibe LAS TAREAS DE ANA.
// Fuga de datos entre usuarios: no es un bug de rendimiento,
// es una brecha de seguridad notificable.
// Opción A: no uses el interceptor automático en datos por usuario.
// Opción B: si lo usas, la identidad SIEMPRE forma parte de la clave.
@Injectable()
export class CachePorUsuarioInterceptor extends CacheInterceptor {
protected trackBy(context: ExecutionContext): string | undefined {
const req = context.switchToHttp().getRequest();
if (req.method !== 'GET') return undefined; // no cachear escrituras
const userId = req.user?.sub;
// Sin usuario identificado, mejor no cachear que cachear mal.
if (!userId) return undefined;
return `u:${userId}:${req.originalUrl}`;
}
}
Cache-Control: private, no-store) para que ninguna
CDN ni ningún proxy intermedio guarde una copia. Los dos fallos —clave sin usuario y cabecera pública— se
cometen por separado y se pagan igual.
11.2.6 Uso manual: el patrón cache-aside
El interceptor sirve para catálogos. Para todo lo demás se inyecta la caché y se controla a mano. El patrón canónico se llama cache-aside (o lazy loading): leer de caché, y si no está, calcular, guardar y devolver.
┌──────────────┐ 1. get(clave) ┌───────────┐
│ Servicio │───────────────────►│ CACHÉ │
│ │◄───────────────────│ (Redis) │
└──────┬───────┘ 2. hit → fin └───────────┘
│ 2'. miss ▲
▼ │ 4. set(clave, valor, ttl)
┌──────────────┐ 3. consulta │
│ Base de datos│──────────────────────────┘
└──────────────┘
import { Inject, Injectable, Logger } from '@nestjs/common';
import { CACHE_MANAGER } from '@nestjs/cache-manager';
import { Cache } from 'cache-manager';
@Injectable()
export class InformesService {
private readonly log = new Logger(InformesService.name);
constructor(
@Inject(CACHE_MANAGER) private readonly cache: Cache,
private readonly em: EntityManager,
) {}
async resumenMensual(proyectoId: string, mes: string): Promise<Resumen> {
const clave = `informe:resumen:${proyectoId}:${mes}`;
const enCache = await this.cache.get<Resumen>(clave);
if (enCache) return enCache; // acierto
const resumen = await this.calcularResumen(proyectoId, mes); // ~800 ms
// TTL con jitter: evita que todas las claves caduquen a la vez.
const ttl = 300_000 + Math.floor(Math.random() * 30_000);
await this.cache.set(clave, resumen, ttl);
return resumen;
}
/** Envoltorio reutilizable: no repitas el try/get/set en veinte servicios. */
async conCache<T>(clave: string, ttlMs: number, fn: () => Promise<T>): Promise<T> {
try {
const hit = await this.cache.get<T>(clave);
if (hit !== undefined && hit !== null) return hit;
} catch (e) {
// Redis caído NO debe tumbar la aplicación: degradar, no fallar.
this.log.warn(`Caché no disponible en lectura (${clave}): ${e}`);
return fn();
}
const valor = await fn();
try {
await this.cache.set(clave, valor, ttlMs);
} catch (e) {
this.log.warn(`Caché no disponible en escritura (${clave}): ${e}`);
}
return valor;
}
}
1. La caché es opcional, la aplicación no. Si Redis se cae y tu get lanza una
excepción sin capturar, has convertido una caché en un punto único de fallo. Envuelve siempre en
try/catch y degrada a la fuente original.
2. Distingue «no está» de «está y vale null». Si un valor legítimo puede ser
null, 0 o cadena vacía, comprobar if (!hit) hace que nunca se
considere un acierto y la caché no sirve para nada. Compara contra undefined o guarda un objeto
envoltorio.
3. Cachear los «no encontrado» evita el cache penetration. Si un atacante pide mil identificadores inexistentes, cada petición atraviesa la caché y llega a la base de datos. Guarda un marcador de ausencia con TTL corto (30 s) para amortiguarlo.
11.2.7 Ejemplo de producción: listado con filtros, usuario e invalidación por evento
Este es el caso realista: un listado paginado, con filtros arbitrarios, específico del usuario, que debe reflejar de inmediato los cambios que ese usuario hace. La solución combina tres ideas: clave canónica (los mismos filtros en distinto orden deben producir la misma clave), versión por usuario (invalidar sin enumerar claves) y eventos de dominio (quien escribe no conoce a quien cachea).
import { Inject, Injectable } from '@nestjs/common';
import { CACHE_MANAGER } from '@nestjs/cache-manager';
import { Cache } from 'cache-manager';
import { OnEvent } from '@nestjs/event-emitter';
import { createHash } from 'node:crypto';
interface FiltroTareas {
estado?: 'pendiente' | 'hecha';
etiquetas?: string[];
page?: number;
limit?: number;
}
@Injectable()
export class TareasCacheService {
private static readonly TTL = 120_000; // 2 min como red de seguridad
constructor(@Inject(CACHE_MANAGER) private readonly cache: Cache) {}
/**
* Clave canónica: ordenamos las entradas y los arrays para que
* ?estado=hecha&page=1 y ?page=1&estado=hecha compartan entrada.
* El hash mantiene la clave corta y sin caracteres problemáticos.
*/
private huella(filtro: FiltroTareas): string {
const normalizado = {
estado: filtro.estado ?? null,
etiquetas: [...(filtro.etiquetas ?? [])].sort(),
page: filtro.page ?? 1,
limit: Math.min(filtro.limit ?? 20, 100),
};
return createHash('sha1').update(JSON.stringify(normalizado)).digest('hex').slice(0, 16);
}
/** Versión del espacio de claves de un usuario: invalidar = incrementarla. */
private async version(userId: string): Promise<number> {
const v = await this.cache.get<number>(`tareas:ver:${userId}`);
if (v === undefined || v === null) {
await this.cache.set(`tareas:ver:${userId}`, 1, 0); // 0 = sin caducidad
return 1;
}
return v;
}
async listar(
userId: string,
filtro: FiltroTareas,
cargar: () => Promise<Pagina<TareaDto>>,
): Promise<Pagina<TareaDto>> {
const v = await this.version(userId);
const clave = `tareas:u:${userId}:v:${v}:f:${this.huella(filtro)}`;
const hit = await this.cache.get<Pagina<TareaDto>>(clave);
if (hit) return hit;
const pagina = await cargar();
await this.cache.set(clave, pagina, TareasCacheService.TTL);
return pagina;
}
/**
* Invalidación por evento: una sola escritura invalida TODAS las
* combinaciones de filtros y páginas de ese usuario, sin recorrer claves.
* (Buscar claves con KEYS * en Redis en producción es un antipatrón:
* bloquea el servidor. Con este truco no hace falta.)
*/
@OnEvent('tarea.creada')
@OnEvent('tarea.actualizada')
@OnEvent('tarea.eliminada')
async invalidar(evento: { propietarioId: string }): Promise<void> {
const actual = await this.version(evento.propietarioId);
await this.cache.set(`tareas:ver:${evento.propietarioId}`, actual + 1, 0);
}
}
No necesita conocer de antemano qué claves existen, cosa imposible con filtros libres. No usa
KEYS ni SCAN, que en un Redis con millones de claves son operaciones peligrosas.
Es correcta en concurrencia: si dos escrituras invalidan a la vez, el peor caso es que la versión avance de
más, lo que solo provoca un fallo de caché adicional. Y mantiene el TTL como red: si un evento se pierde,
la incoherencia dura dos minutos, no para siempre.
El único coste es que las entradas de versiones antiguas siguen en memoria hasta caducar. Con un TTL de dos minutos, ese coste está acotado y es irrelevante.
get seguido de set del ejemplo no es atómico. Para este caso concreto da igual
(perder una invalidación simultánea equivale a invalidar una sola vez, que es el efecto deseado). Si
necesitas un contador realmente atómico —por ejemplo para límites de uso facturables— no uses la
abstracción de caché: usa el cliente de Redis directamente y su comando INCR, que sí es
atómico.
11.3 Colas y trabajos en segundo plano
11.3.1 Por qué una petición HTTP no debe hacerlo todo
Un usuario se registra. Hay que crear la fila en la base de datos, enviar un correo de bienvenida, generar
su avatar, avisar al CRM y apuntar el evento en la analítica. Si todo eso ocurre dentro del
POST /registro, la respuesta tarda tres segundos y, lo que es peor, si el servidor de correo
está caído el registro falla aunque la parte importante haya funcionado. Has acoplado la corrección de tu
caso de uso a la disponibilidad de un tercero.
Una cola invierte el planteamiento: la petición hace solo lo imprescindible y transaccional, encola el resto y responde. Un proceso aparte —el worker— consume esos trabajos con sus propios reintentos y a su propio ritmo.
PETICIÓN HTTP (rápida, transaccional) WORKER (lento, reintentable)
┌───────────────────────────┐ ┌──────────────────────────────┐
│ 1. valida DTO │ │ 5. toma el trabajo (BRPOPLPUSH│
│ 2. INSERT usuario ──┐ │ │ atómico: nadie más lo ve) │
│ 3. queue.add(...) ───┼────┼── PRODUCTOR ──► │ 6. envía el correo │
│ 4. responde 201 │ │ ┌────────┐ │ 7a. ok → completed │
└──────────────────────┼────┘ │ REDIS │◄────│ 7b. err → retry con backoff │
~40 ms │ │ cola │ │ 7c. agotado → failed (DLQ) │
└──────────►└────────┘ └──────────────────────────────┘
▲ ▲
│ └── otros workers: escalado horizontal
└────── panel de monitorización
Estados de un trabajo: waiting ──► active ──► completed
▲ │
└─ delayed ─┴──► failed ──► (reintento) ──► waiting
Las cuatro razones que justifican una cola, por orden de frecuencia real:
- No bloquear la respuesta. Todo lo que el usuario no necesita ver para continuar debe salir del camino crítico: correos, notificaciones push, sincronizaciones, webhooks salientes.
- Reintentos con garantías. Un
try/catchen el controlador no sobrevive a un reinicio del proceso. Un trabajo en Redis sí: se reintentará aunque despliegues en medio. - Absorber picos. Si llegan 10 000 peticiones en un minuto, la cola crece y los workers las procesan a su ritmo sostenible. Sin cola, se saturan las conexiones a la base de datos y caen todas.
- Trabajos largos. Generar un PDF de 400 páginas o procesar un CSV de 2 millones de filas no cabe en el timeout de un balanceador (típicamente 30-60 s) y bloquearía el event loop (capítulo 1).
11.3.2 BullMQ con @nestjs/bullmq
BullMQ es la evolución de Bull: una cola de trabajos sobre Redis, escrita en TypeScript, que usa scripts Lua para garantizar atomicidad en las transiciones de estado. Es el estándar de facto en el ecosistema Node.
@nestjs/bull y @nestjs/bullmq no son intercambiables
Son dos paquetes distintos para dos librerías distintas (bull y bullmq). El
código de uno no compila en el otro. Si sigues un tutorial y las importaciones «no existen», es esto:
| Concepto | @nestjs/bull (Bull, legado) | @nestjs/bullmq (BullMQ, actual) |
|---|---|---|
| Conexión | redis: { host, port } | connection: { host, port } |
| Consumidor | Clase con métodos @Process() | Clase que extiende WorkerHost e implementa process() |
| Varios tipos de trabajo | @Process('nombre') por método | Un único process(); se discrimina con switch (job.name) |
| Eventos | @OnQueueCompleted(), @OnQueueFailed() | @OnWorkerEvent('completed'), @OnWorkerEvent('failed') |
| Progreso | job.progress(50) | await job.updateProgress(50) |
| Nombre del trabajo al encolar | Opcional | Obligatorio: queue.add('nombre', datos) |
| Requisitos | Redis 2.8+ | Redis 6.2 o superior recomendado |
Además, dentro de BullMQ la API de trabajos repetibles cambió: las versiones recientes de la serie
5 introdujeron los job schedulers (upsertJobScheduler) y marcaron como obsoleta la
antigua repeat con removeRepeatableByKey. Verifica en la documentación de tu
versión cuál está disponible antes de escribir la migración.
import { BullModule } from '@nestjs/bullmq';
@Module({
imports: [
BullModule.forRootAsync({
imports: [ConfigModule],
inject: [ConfigService],
useFactory: (c: ConfigService) => ({
connection: {
host: c.getOrThrow('REDIS_HOST'),
port: c.get<number>('REDIS_PORT', 6379),
// Obligatorio si pasas tu propia instancia de ioredis:
// los workers usan comandos bloqueantes y no admiten
// el reintento por comando de ioredis.
maxRetriesPerRequest: null,
},
defaultJobOptions: {
attempts: 3,
backoff: { type: 'exponential', delay: 2000 },
removeOnComplete: { age: 3600, count: 1000 },
removeOnFail: { age: 24 * 3600 },
},
}),
}),
BullModule.registerQueue({ name: 'correos' }, { name: 'informes' }),
],
})
export class AppModule {}
import { InjectQueue } from '@nestjs/bullmq';
import { Queue } from 'bullmq';
@Injectable()
export class CorreosProducer {
constructor(@InjectQueue('correos') private readonly cola: Queue) {}
async bienvenida(usuarioId: string, email: string) {
await this.cola.add(
'bienvenida', // nombre: obligatorio en BullMQ
{ usuarioId, email }, // datos: deben ser serializables a JSON
{
// jobId estable = deduplicación: si ya existe un trabajo
// con este id, add() no crea otro.
jobId: `bienvenida:${usuarioId}`,
attempts: 5,
backoff: { type: 'exponential', delay: 5000 },
priority: 1, // 1 = máxima prioridad
},
);
}
}
add() pierde sus métodos, sus referencias y su vínculo con el
EntityManager; y un Date vuelve como cadena. Encola identificadores, no
objetos: el worker recarga lo que necesite. Además, evita meter datos personales o secretos en el trabajo,
porque quedan visibles en Redis y en el panel de monitorización durante días.
11.3.3 El consumidor: @Processor y WorkerHost
import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq';
import { Job, UnrecoverableError } from 'bullmq';
import { Logger } from '@nestjs/common';
@Processor('correos', {
concurrency: 5, // 5 trabajos en paralelo POR réplica
limiter: { max: 100, duration: 60_000 }, // tope de 100/min: respeta la cuota del proveedor
})
export class CorreosProcessor extends WorkerHost {
private readonly log = new Logger(CorreosProcessor.name);
constructor(
private readonly smtp: SmtpService,
private readonly usuarios: UsuariosService,
) {
super();
}
/** Punto de entrada único: un solo process() para todos los nombres. */
async process(job: Job<DatosCorreo, void, string>): Promise<void> {
switch (job.name) {
case 'bienvenida':
return this.enviarBienvenida(job);
case 'recuperar-clave':
return this.enviarRecuperacion(job);
default:
// Un nombre desconocido nunca se arreglará reintentando.
throw new UnrecoverableError(`Trabajo no soportado: ${job.name}`);
}
}
private async enviarBienvenida(job: Job<DatosCorreo>): Promise<void> {
const usuario = await this.usuarios.findById(job.data.usuarioId);
// El usuario se borró entre el encolado y el consumo: no es un error del sistema.
if (!usuario) throw new UnrecoverableError('Usuario inexistente');
// Idempotencia: si ya se envió, no lo repitas (ver 11.3.5).
if (usuario.bienvenidaEnviadaEn) {
this.log.log(`Bienvenida ya enviada a ${usuario.id}, se omite`);
return;
}
await job.updateProgress(30);
await this.smtp.enviar({ to: usuario.email, plantilla: 'bienvenida' });
await this.usuarios.marcarBienvenidaEnviada(usuario.id, new Date());
await job.updateProgress(100);
}
@OnWorkerEvent('completed')
onCompleted(job: Job) {
this.log.log(`OK ${job.name}#${job.id} en ${Date.now() - job.timestamp} ms`);
}
@OnWorkerEvent('failed')
onFailed(job: Job | undefined, err: Error) {
// attemptsMade incluye el intento actual: compáralo con opts.attempts
const agotado = job ? job.attemptsMade >= (job.opts.attempts ?? 1) : true;
const nivel = agotado ? 'error' : 'warn';
this.log[nivel](`FALLO ${job?.name}#${job?.id} (${job?.attemptsMade}): ${err.message}`);
if (agotado) this.alertas.notificar(job, err); // ahora sí, avisa a una persona
}
@OnWorkerEvent('stalled')
onStalled(jobId: string) {
// Un trabajo "atascado" es un worker que dejó de renovar su lock:
// casi siempre significa que bloqueaste el event loop o que te mató el OOM killer.
this.log.warn(`Trabajo atascado ${jobId}: revisa CPU y memoria del worker`);
}
}
concurrency es por instancia del worker, así que la concurrencia real es
concurrency × réplicas. El límite lo marca casi siempre el recurso más escaso, y rara vez es la
CPU de Node: es el pool de conexiones de la base de datos. Si tu pool tiene 10 conexiones y arrancas
4 réplicas con concurrencia 10, tienes 40 trabajos peleando por 10 conexiones y verás timeouts de
adquisición que parecen errores de la base de datos.
Para trabajo intensivo de CPU, la concurrencia no ayuda en absoluto (un solo hilo) y conviene usar
useWorkerThreads o procesos separados. Para trabajo de E/S (correos, HTTP, consultas), valores
entre 5 y 25 son razonables.
11.3.4 Opciones de trabajo: el vocabulario completo
| Opción | Qué hace | Cuándo y cómo usarla |
|---|---|---|
attempts | Número total de intentos, incluido el primero | 3-5 para errores transitorios. 1 si el trabajo no es idempotente y no puedes arreglarlo |
backoff | { type: 'fixed' | 'exponential', delay } | Exponencial siempre que el fallo pueda ser saturación del destino; fijo para reintentos rápidos previsibles |
delay | Retrasa la primera ejecución (ms) | «Enviar recordatorio a las 24 h»; también para esperar a que una transacción confirme |
priority | Entero; menor es más prioritario | Correos transaccionales por delante de campañas. Ojo: usar prioridades tiene coste de rendimiento en colas enormes |
jobId | Identificador propio; si ya existe, no se crea otro | Deduplicación natural. Clave para el patrón «un trabajo por entidad y día» |
removeOnComplete | true, o { age, count } | Ponlo siempre. Sin esto Redis crece sin límite hasta llenar la memoria |
removeOnFail | Igual, para fallidos | Conserva los fallidos más tiempo que los completados: son tu material de diagnóstico |
lifo | Último en entrar, primero en salir | Casos interactivos donde lo reciente vale más que lo antiguo |
repeat | { pattern, tz, limit } | Trabajos periódicos gestionados por la cola. Revisa si tu versión ya prefiere los job schedulers |
jobId tiene una letra pequeña importante
Un jobId repetido evita crear un trabajo duplicado mientras ese trabajo siga existiendo en
Redis. Si ya se completó y fue eliminado por removeOnComplete, el mismo jobId
vuelve a ser aceptado. Por tanto jobId protege de duplicados cercanos en el tiempo, no es
una garantía de «exactamente una vez» para siempre. Las versiones recientes de BullMQ añaden además una
opción explícita de deduplication con su propia ventana de tiempo; compruébala en tu versión.
La garantía real de no duplicar efectos la da la idempotencia del consumidor, nunca el productor.
11.3.5 Idempotencia del consumidor: lo más importante de esta sección
Toda cola distribuida ofrece entrega «al menos una vez», no «exactamente una vez». Esto no es un defecto de BullMQ: es una consecuencia de que la red y los procesos fallan. Un worker puede enviar el correo y morir antes de marcar el trabajo como completado; el trabajo volverá a la cola y el correo se enviará dos veces. La única defensa es que ejecutar el trabajo dos veces produzca el mismo resultado que ejecutarlo una.
async process(job: Job<{ pedidoId: string }>) {
const pedido = await this.pedidos.find(job.data.pedidoId);
// Si el proceso muere DESPUÉS del cargo y ANTES de
// completar el trabajo, se reintenta y se cobra otra vez.
await this.pasarela.cobrar(pedido.total, pedido.tarjeta);
pedido.estado = 'pagado';
await this.em.flush();
// Efecto secundario adicional, también duplicable:
await this.smtp.enviar({ to: pedido.email, plantilla: 'recibo' });
}
async process(job: Job<{ pedidoId: string }>) {
const pedido = await this.pedidos.find(job.data.pedidoId);
// 1. Comprobación de estado: barrera barata para el caso normal.
if (pedido.estado === 'pagado') return;
// 2. Clave de idempotencia enviada a la pasarela: si repites la
// misma clave, la pasarela devuelve el cargo original en lugar
// de crear otro. Es el mecanismo que ofrecen Stripe, Adyen, etc.
const clave = `pedido:${pedido.id}:cobro`;
const cargo = await this.pasarela.cobrar(pedido.total, pedido.tarjeta, {
idempotencyKey: clave,
});
// 3. Escritura condicional: solo pasa de pendiente a pagado.
// Si otro worker ganó la carrera, afectadas = 0 y salimos.
const afectadas = await this.pedidos.marcarPagadoSiPendiente(pedido.id, cargo.id);
if (afectadas === 0) return;
// 4. El correo se encola aparte, con su propio jobId estable.
await this.correos.add('recibo', { pedidoId: pedido.id }, {
jobId: `recibo:${pedido.id}`,
});
}
Las cuatro técnicas de idempotencia, por orden de robustez:
- Operaciones naturalmente idempotentes.
UPDATE ... SET estado = 'pagado'es idempotente;UPDATE ... SET intentos = intentos + 1no lo es. Diseña hacia la primera forma siempre que puedas. - Escritura condicional.
UPDATE ... WHERE id = ? AND estado = 'pendiente'y comprueba las filas afectadas. Es una transición de estado atómica en la propia base de datos, sin locks explícitos. - Tabla de trabajos procesados. Una tabla con clave primaria
(tipo, clave); el consumidor inserta al empezar y una violación de unicidad significa «ya hecho». Funciona con cualquier efecto secundario, incluso los que no controlas. - Clave de idempotencia en el tercero. Cuando el efecto secundario está fuera de tu sistema (cobros, correos, mensajes), es la única forma de deduplicar de verdad. Casi todas las APIs de pago la ofrecen; úsala.
11.3.6 Eventos, fallos y monitorización
- Eventos locales (
@OnWorkerEvent): los emite tu worker sobre sus trabajos. Perfecto para métricas y trazas del consumidor. - Eventos globales (clase
QueueEventsde BullMQ): llegan a cualquier proceso suscrito, incluido el que solo produce. Es lo que necesitas si la API quiere reaccionar cuando termina un trabajo que encoló (por ejemplo, para avisar por WebSocket, sección 11.6). - Cola de fallidos. En BullMQ los trabajos agotados quedan en el conjunto
failed, con sufailedReasony su stack trace. Desde ahí se pueden reintentar en masa (job.retry()) una vez arreglada la causa. Es tu dead letter queue de facto. - Métricas mínimas que debes vigilar: tamaño de
waiting(¿crece sin parar?), edad del trabajo más antiguo en espera, ratio defailedy trabajosstalled. Una cola que crece indefinidamente es un incidente aunque nada dé error. - Panel.
bull-board(con su adaptador para BullMQ y su integración para Nest) ofrece una interfaz web para inspeccionar, reintentar y limpiar. Protégelo siempre con autenticación: expone los datos de todos los trabajos y permite borrar la cola.
@Injectable()
export class QueueMetricsService {
constructor(@InjectQueue('correos') private readonly cola: Queue) {}
/** Expón esto en /metrics o en el health check de readiness. */
async estado() {
const [waiting, active, delayed, failed] = await Promise.all([
this.cola.getWaitingCount(),
this.cola.getActiveCount(),
this.cola.getDelayedCount(),
this.cola.getFailedCount(),
]);
return { waiting, active, delayed, failed, saludable: waiting < 5_000 };
}
}
11.3.7 Ejemplo real: informe pesado con seguimiento de progreso
Caso de uso completo: el usuario pide un informe que tarda dos minutos. La API responde de inmediato con un identificador, el worker lo genera actualizando el progreso, lo sube a un almacén de objetos y avisa al usuario. Es el patrón 202 Accepted con consulta de estado.
@Controller('informes')
export class InformesController {
constructor(
@InjectQueue('informes') private readonly cola: Queue,
private readonly usuarios: UsuariosService,
) {}
@Post()
@HttpCode(202) // 202 Accepted: aceptado, aún no terminado
async solicitar(@Body() dto: CrearInformeDto, @CurrentUser() user: JwtUser) {
const jobId = `informe:${user.id}:${dto.mes}`;
// Deduplicación: si ya está generándose el mismo informe, no dupliques trabajo.
const existente = await this.cola.getJob(jobId);
if (existente) {
return { id: jobId, estado: await existente.getState(), duplicado: true };
}
await this.cola.add('mensual', { userId: user.id, mes: dto.mes }, {
jobId,
attempts: 3,
backoff: { type: 'exponential', delay: 10_000 },
removeOnComplete: { age: 7 * 24 * 3600 }, // conserva una semana para consultar estado
});
return { id: jobId, estado: 'waiting' };
}
@Get(':id/estado')
async estado(@Param('id') id: string, @CurrentUser() user: JwtUser) {
const job = await this.cola.getJob(id);
if (!job) throw new NotFoundException();
// AUTORIZACIÓN: el id contiene el usuario, pero nunca te fíes del formato.
if (job.data.userId !== user.id) throw new ForbiddenException();
return {
estado: await job.getState(),
progreso: job.progress,
url: job.returnvalue?.url ?? null,
error: job.failedReason ?? null,
};
}
}
@Processor('informes', { concurrency: 2 }) // pesados: poca concurrencia
export class InformesProcessor extends WorkerHost {
constructor(
private readonly em: EntityManager,
private readonly almacen: AlmacenService,
private readonly eventos: EventEmitter2,
) {
super();
}
async process(job: Job<{ userId: string; mes: string }>): Promise<{ url: string }> {
const { userId, mes } = job.data;
// El worker vive fuera de una petición HTTP: NO hay contexto de petición.
// Con MikroORM hay que crear un contexto propio (ver capítulo 15):
// em.fork() da un EntityManager con su propia Identity Map y Unit of Work.
const em = this.em.fork();
const filas: FilaInforme[] = [];
const total = await em.count(Movimiento, { mes });
const LOTE = 1_000;
// Paginar en lotes: cargar 2 millones de filas de golpe reventaría la memoria.
for (let offset = 0; offset < total; offset += LOTE) {
const lote = await em.find(Movimiento, { mes }, { limit: LOTE, offset, orderBy: { id: 'ASC' } });
filas.push(...lote.map(transformar));
// Liberamos la Identity Map: si no, el fork acumula todas las entidades.
em.clear();
await job.updateProgress(Math.round((offset / total) * 90));
}
const pdf = await generarPdf(filas);
const url = await this.almacen.subir(`informes/${userId}/${mes}.pdf`, pdf);
await job.updateProgress(100);
// Avisamos al resto del sistema (correo, WebSocket) sin acoplarnos a ellos.
this.eventos.emit('informe.listo', { userId, mes, url });
// Lo devuelto queda en job.returnvalue y lo lee el endpoint de estado.
return { url };
}
}
1. Sin contexto de petición no hay EntityManager por defecto. El
RequestContext de MikroORM lo crea un middleware HTTP que en el worker no existe. Usa
em.fork() por trabajo (o el helper de contexto correspondiente) para que dos trabajos
concurrentes no compartan Identity Map ni Unit of Work.
2. La memoria crece si no limpias. Un bucle que carga lotes sin em.clear() mantiene
vivas todas las entidades por la Identity Map. Es la causa número uno de workers que mueren por OOM
procesando ficheros grandes.
11.3.8 Tareas programadas con @nestjs/schedule
Cuando el trabajo no lo dispara un usuario sino el reloj, no hace falta una cola: basta un planificador en
el proceso. @nestjs/schedule envuelve la librería cron y añade tres decoradores.
import { Cron, CronExpression, Interval, Timeout, SchedulerRegistry } from '@nestjs/schedule';
@Injectable()
export class MantenimientoService {
private readonly log = new Logger(MantenimientoService.name);
// Expresión explícita con zona horaria: imprescindible en producción,
// porque el contenedor casi siempre corre en UTC.
@Cron('0 30 3 * * *', { name: 'purga-nocturna', timeZone: 'Europe/Madrid' })
async purgarSesiones() {
const borradas = await this.sesiones.borrarExpiradas();
this.log.log(`Sesiones purgadas: ${borradas}`);
}
// Constantes legibles para los casos habituales.
@Cron(CronExpression.EVERY_10_MINUTES)
async sincronizarTarifas() { /* ... */ }
// Cada 30 s desde el arranque. No es cron: es un setInterval gestionado.
@Interval('latido', 30_000)
latido() { this.metricas.gauge('app.vivo', 1); }
// Una sola vez, 10 s después del arranque: ideal para precalentar caché.
@Timeout(10_000)
async precalentar() { await this.catalogo.precargar(); }
/** Programación dinámica: añadir o quitar tareas en tiempo de ejecución. */
constructor(private readonly registry: SchedulerRegistry) {}
pausarPurga() {
this.registry.getCronJob('purga-nocturna').stop();
}
}
| Campo | 1 | 2 | 3 | 4 | 5 | 6 |
|---|---|---|---|---|---|---|
| Significado | segundo opcional | minuto | hora | día del mes | mes | día de la semana |
| Rango | 0-59 | 0-59 | 0-23 | 1-31 | 1-12 | 0-7 (0 y 7 = domingo) |
Ejemplos leídos en voz alta, que es la única forma fiable de revisar un cron:
0 30 3 * * *→ «segundo 0, minuto 30, hora 3, cualquier día»: todos los días a las 3:30.0 */15 * * * *→ cada 15 minutos, en el segundo 0 (*/nsignifica «cada n»).0 0 9 * * 1-5→ de lunes a viernes a las 9:00.0 0 0 1 * *→ el día 1 de cada mes a medianoche.0 0 2 * * 0→ los domingos a las 2:00. Cuidado: si además pones día del mes, la mayoría de implementaciones interpretan los dos campos como un O lógico, no como un Y.
Escalas la API a tres pods. El decorador @Cron se registra en cada proceso, así que a
las 3:30 se ejecuta tres veces. Con una purga idempotente no pasa nada; con «generar y enviar las
facturas del mes» acabas de cometer un error con consecuencias legales y contables.
Soluciones, de peor a mejor:
a) Variable de entorno. Solo una réplica tiene SCHEDULER=true y el módulo de tareas se
registra condicionalmente. Simple, pero si esa réplica cae, no hay tareas y nadie se enterará.
b) Despliegue dedicado. Un proceso aparte, con una sola réplica, que solo planifica y encola.
Es la opción más limpia arquitectónicamente y la habitual en Kubernetes junto con CronJob.
c) Bloqueo distribuido. Todas las réplicas despiertan, pero solo una consigue el cerrojo en Redis y ejecuta. Tolerante a fallos y sin configuración especial por réplica.
d) Trabajos repetibles de BullMQ. La propia cola garantiza que el trabajo periódico se encola una sola vez, aunque veinte réplicas lo declaren, porque la programación vive en Redis y no en el proceso. Es la respuesta correcta cuando ya tienes BullMQ, y por eso conviene usar el planificador solo para encolar, nunca para ejecutar.
@Cron('0 0 3 1 * *')
async facturarMes() {
// Con 3 réplicas: 3 ejecuciones simultáneas.
// Además, todo el trabajo pesado ocurre DENTRO del cron:
// si falla a mitad, no hay reintento ni rastro de por dónde iba.
const clientes = await this.clientes.activos();
for (const c of clientes) {
const factura = await this.crearFactura(c);
await this.smtp.enviar(factura);
}
}
@Cron('0 0 3 1 * *', { timeZone: 'Europe/Madrid' })
async planificarFacturacion() {
// 1. Cerrojo distribuido: SET NX PX es atómico en Redis.
// Solo una réplica sigue; el TTL evita cerrojos huérfanos.
const periodo = new Date().toISOString().slice(0, 7); // '2026-07'
const conseguido = await this.redis.set(
`lock:facturacion:${periodo}`, process.env.HOSTNAME ?? '1',
'PX', 300_000, 'NX',
);
if (!conseguido) return;
// 2. El cron solo ENCOLA; el trabajo pesado, con reintentos, va a la cola.
const clientes = await this.clientes.activos({ fields: ['id'] });
await this.cola.addBulk(clientes.map((c) => ({
name: 'facturar-cliente',
data: { clienteId: c.id, periodo },
// jobId por cliente y periodo: la garantía real de "una factura por mes".
opts: { jobId: `factura:${c.id}:${periodo}`, attempts: 5 },
})));
}
11.4 Eventos internos: desacoplar módulos sin salir del proceso
Cuando un caso de uso termina, a menudo hay que avisar a otros módulos: se ha creado una tarea, hay que invalidar la caché, notificar por WebSocket, apuntar la métrica y quizá encolar un correo. Si el servicio de tareas llama a esos cuatro servicios, el módulo de tareas acaba importando media aplicación y cada nueva reacción obliga a modificar el código que la origina, violando el principio de abierto/cerrado.
@nestjs/event-emitter invierte esa dependencia: el emisor publica un hecho y no sabe quién
escucha. Por debajo es eventemitter2, es decir, memoria del proceso actual.
EventEmitterModule.forRoot({
wildcard: true, // permite escuchar 'tarea.*'
delimiter: '.',
maxListeners: 20,
verboseMemoryLeak: true, // avisa si un evento acumula listeners
});
// Evento de dominio: una clase, no un objeto suelto.
// Nombre en pasado: describe algo que YA ocurrió.
export class TareaCreadaEvent {
constructor(
readonly tareaId: string,
readonly propietarioId: string,
readonly proyectoId: string,
) {}
}
@Injectable()
export class TareasService {
constructor(private readonly eventos: EventEmitter2) {}
async crear(dto: CrearTareaDto, userId: string) {
const tarea = this.em.create(Tarea, { ...dto, propietario: userId });
await this.em.flush(); // primero persistimos...
// ...y solo después publicamos: nunca anuncies algo que aún puede fallar.
this.eventos.emit('tarea.creada', new TareaCreadaEvent(tarea.id, userId, dto.proyectoId));
return tarea;
}
}
@Injectable()
export class TareasListener {
private readonly log = new Logger(TareasListener.name);
// async: true ejecuta el listener sin bloquear al emisor.
@OnEvent('tarea.creada', { async: true })
async alCrear(e: TareaCreadaEvent) {
// Un listener NUNCA debe romper al emisor: captura tus errores.
try {
await this.push.notificar(e.propietarioId, 'Tarea creada');
} catch (err) {
this.log.error(`Fallo notificando ${e.tareaId}`, err as Error);
}
}
// Comodín: útil para auditoría y métricas transversales.
@OnEvent('tarea.*')
auditar(payload: unknown, ...args: unknown[]) {
this.auditoria.registrar(payload);
}
}
No hay persistencia. Si el proceso muere entre el emit y el listener, el evento se
pierde para siempre y nadie se enterará: no hay reintento, ni cola de fallidos, ni traza.
No sale del proceso. Con tres réplicas, un evento emitido en la réplica A no llega a los listeners de B ni de C. Esto convierte en incorrecto el patrón «invalido la caché en memoria con un evento».
Por defecto es sincrónico. emit() ejecuta los listeners en la misma pila; si uno tarda
200 ms, la petición del emisor tarda 200 ms más. Con { async: true } se desacopla, pero
entonces las excepciones del listener quedan fuera del try/catch del emisor y pueden
convertirse en rechazos de promesa no gestionados.
Se pierde la trazabilidad. Leyendo emit('tarea.creada') no se sabe qué ocurre después.
Documenta los eventos y sus consumidores, o el sistema se vuelve imposible de razonar («programación
vudú»).
| Criterio | Llamada directa | Evento interno | Cola real (BullMQ) |
|---|---|---|---|
| Acoplamiento | Alto: el emisor importa al receptor | Bajo: solo comparten el nombre del evento | Bajo, además entre procesos |
| Persistencia | N/A | Ninguna | Sí, en Redis |
| Reintentos | Manuales | Ninguno | Automáticos con backoff |
| Cruza réplicas | No aplica | No | Sí |
| Trazabilidad | Máxima: se sigue leyendo el código | Baja | Media: hay panel e histórico |
| Coste | Cero | Cero | Redis + workers + operación |
| Úsalo para… | Lógica que forma parte del caso de uso y debe fallar con él | Reacciones internas, opcionales e idempotentes: métricas, caché, auditoría | Efectos que no pueden perderse: correos, pagos, integraciones |
Regla de decisión. Pregúntate: «si esta reacción no ocurre, ¿alguien se queja?». Si la respuesta es sí, necesitas una cola. Si es «se pierde una métrica», un evento interno es perfecto. Y si la reacción es parte de la corrección del caso de uso —descontar el stock al confirmar el pedido— entonces no es un evento: es una llamada directa dentro de la misma transacción.
11.5 Rate limiting y protección de la API
Sin límite de peticiones, cualquiera puede probar diez mil contraseñas por minuto en tu /login,
agotar tu cuota de un servicio de pago o tumbar la base de datos con un bucle. @nestjs/throttler
implementa una ventana temporal por identificador: N peticiones cada T milisegundos.
@nestjs/throttler v4 y anteriores, ttl se expresaba en segundos y la
configuración era un objeto plano. Desde v5 se declara un array throttlers, cada uno con
nombre opcional, y el ttl va en milisegundos. Un ttl: 60 copiado de una guía
antigua produce una ventana de 60 ms, es decir, ningún límite efectivo. El decorador también cambió de
@Throttle(limit, ttl) a @Throttle({ nombre: { limit, ttl } }).
import { ThrottlerModule, ThrottlerGuard } from '@nestjs/throttler';
import { ThrottlerStorageRedisService } from 'nestjs-throttler-storage-redis';
@Module({
imports: [
ThrottlerModule.forRootAsync({
inject: [ConfigService],
useFactory: (c: ConfigService) => ({
// Varias ventanas simultáneas: todas deben cumplirse.
// Corta las ráfagas sin castigar el uso sostenido normal.
throttlers: [
{ name: 'corta', ttl: 1_000, limit: 10 }, // 10 por segundo
{ name: 'media', ttl: 60_000, limit: 120 }, // 120 por minuto
{ name: 'larga', ttl: 3_600_000, limit: 2_000 },
],
// Sin almacenamiento compartido, cada réplica cuenta por su cuenta:
// con 4 réplicas el límite real es 4 veces el configurado.
storage: new ThrottlerStorageRedisService(c.getOrThrow('REDIS_URL')),
// No limitar las sondas de salud ni el scrapeo de métricas.
skipIf: (ctx) => ['/health', '/metrics'].includes(
ctx.switchToHttp().getRequest().path,
),
}),
}),
],
providers: [{ provide: APP_GUARD, useClass: ThrottlerGuard }], // global
})
export class AppModule {}
// El login hereda el límite global de 120/min: suficiente
// para 172 800 intentos de contraseña al día por IP.
@Post('login')
login(@Body() dto: LoginDto) {
return this.auth.login(dto);
}
// Y en el arranque, detrás de Nginx o de un balanceador,
// req.ip es la IP DEL PROXY: todos los usuarios del mundo
// comparten un único contador. Un solo usuario activo
// bloquea la API para todos los demás.
// Límite específico y agresivo en el punto sensible.
@Throttle({ corta: { limit: 3, ttl: 60_000 } })
@Post('login')
login(@Body() dto: LoginDto) {
return this.auth.login(dto);
}
@SkipThrottle() // webhook interno de confianza
@Post('webhooks/stripe')
webhook(@Body() body: unknown) { /* ... */ }
// main.ts: decirle a Express de cuántos proxies fiarse.
// El número es la cantidad de saltos de confianza; NUNCA
// uses 'true' en internet abierta, porque entonces cualquiera
// puede falsificar X-Forwarded-For y saltarse el límite.
const app = await NestFactory.create(AppModule);
app.set('trust proxy', 1); // Fastify: { trustProxy: 1 }
Para limitar por usuario en lugar de por IP —lo correcto en endpoints autenticados, porque una oficina
entera comparte IP— se extiende el guard y se sobrescribe getTracker:
@Injectable()
export class ThrottlerUsuarioGuard extends ThrottlerGuard {
protected async getTracker(req: Record<string, any>): Promise<string> {
// Usuario autenticado → contador propio; anónimo → por IP.
// Prefijos distintos para que un id numérico no colisione con una IP.
return req.user?.sub ? `u:${req.user.sub}` : `ip:${req.ips?.[0] ?? req.ip}`;
}
}
helmet para cabeceras de seguridad, límite de tamaño del cuerpo
(express.json({ limit: '100kb' })) y timeouts de servidor. En el capítulo 12 se detalla
la protección del flujo de autenticación.
11.6 WebSockets y tiempo real
11.6.1 El gateway: decoradores y ciclo de vida
HTTP es petición-respuesta: el servidor no puede hablar primero. Para notificar cambios en el momento en que
ocurren hace falta una conexión persistente y bidireccional. Un gateway de Nest es el equivalente a un
controlador, pero para mensajes de WebSocket; por defecto usa Socket.IO, y con WsAdapter puede
usar ws nativo.
import {
WebSocketGateway, WebSocketServer, SubscribeMessage, MessageBody,
ConnectedSocket, OnGatewayInit, OnGatewayConnection, OnGatewayDisconnect,
WsException,
} from '@nestjs/websockets';
import { Server, Socket } from 'socket.io';
@WebSocketGateway({
namespace: '/tareas', // aísla este canal de otros
cors: { origin: ['https://app.ejemplo.com'], credentials: true },
transports: ['websocket'], // evita el fallback a long-polling
})
export class TareasGateway implements OnGatewayInit, OnGatewayConnection, OnGatewayDisconnect {
private readonly log = new Logger(TareasGateway.name);
@WebSocketServer() private readonly server!: Server; // instancia de Socket.IO
constructor(private readonly jwt: JwtService) {}
/** 1. Una vez, al inicializar. Sitio idóneo para middleware de conexión. */
afterInit(server: Server): void {
// El middleware corre ANTES de handleConnection: rechaza aquí lo inválido
// para no gastar recursos en sockets que no deberían existir.
server.use(async (socket: Socket, next) => {
try {
const token = socket.handshake.auth?.token
?? (socket.handshake.headers.authorization ?? '').replace('Bearer ', '');
if (!token) throw new Error('sin token');
const payload = await this.jwt.verifyAsync(token);
socket.data.user = { id: payload.sub, roles: payload.roles };
next();
} catch {
// El cliente recibe un connect_error, no una conexión a medias.
next(new Error('No autorizado'));
}
});
}
/** 2. Por cada cliente que se conecta (ya autenticado por el middleware). */
async handleConnection(client: Socket): Promise<void> {
const userId = client.data.user.id as string;
// Sala personal: permite enviar a "este usuario" en todos sus dispositivos.
await client.join(`user:${userId}`);
this.log.log(`Conectado ${client.id} (usuario ${userId})`);
}
/** 3. Al desconectar. Socket.IO abandona las salas solo; libera TU estado. */
handleDisconnect(client: Socket): void {
this.presencia.marcarOffline(client.data.user?.id);
}
/** Mensaje entrante: el cliente pide seguir un proyecto. */
@SubscribeMessage('proyecto:seguir')
async seguir(
@ConnectedSocket() client: Socket,
@MessageBody() body: { proyectoId: string },
): Promise<{ ok: true }> {
// AUTORIZACIÓN: quien se conecta no puede unirse a la sala que quiera.
// Olvidar esta comprobación es el equivalente en WebSocket a un IDOR.
const permitido = await this.acl.puedeVerProyecto(client.data.user.id, body.proyectoId);
if (!permitido) throw new WsException('Sin acceso a ese proyecto');
await client.join(`proyecto:${body.proyectoId}`);
return { ok: true }; // se entrega como acknowledgement al cliente
}
/** API interna: la usan los servicios y los listeners de eventos. */
notificarCambio(proyectoId: string, tarea: TareaDto): void {
this.server.to(`proyecto:${proyectoId}`).emit('tarea:actualizada', tarea);
}
notificarUsuario(userId: string, payload: unknown): void {
this.server.to(`user:${userId}`).emit('notificacion', payload);
}
}
| Elemento | Equivalente HTTP | Nota práctica |
|---|---|---|
@WebSocketGateway(opts) | @Controller() | Con namespace separas dominios; con path cambias la ruta del handshake |
@SubscribeMessage('x') | @Post('x') | Lo devuelto se envía como ack si el cliente pasó un callback |
@MessageBody() | @Body() | Admite pipes: @UsePipes(new ValidationPipe()) valida igual que en HTTP |
@ConnectedSocket() | @Req() | Guarda datos de sesión en socket.data, no en variables del gateway |
@WebSocketServer() | — | Servidor completo: difusión a salas y a todos |
OnGatewayInit | onModuleInit | afterInit(server): registra middleware y adaptadores |
OnGatewayConnection | Middleware de entrada | handleConnection(client): unir a salas, marcar presencia |
OnGatewayDisconnect | — | handleDisconnect(client): limpieza; se llama también al perder la red |
Un patrón habitual y erróneo es aceptar la conexión y esperar un mensaje auth con el token.
Mientras llega, tienes un socket anónimo consumiendo memoria y capaz de emitir mensajes; y si olvidas
comprobar el estado en algún @SubscribeMessage, queda abierto. Autentica en el
handshake y rechaza antes de establecer la conexión.
Sobre dónde va el token: en handshake.auth (Socket.IO lo envía en el propio protocolo)
es preferible a handshake.query, porque los parámetros de consulta acaban en los registros de
acceso del proxy y en el historial de errores. La cabecera Authorization solo está disponible en
el handshake HTTP inicial y no en el navegador con la API WebSocket nativa: por eso
Socket.IO ofrece auth.
Los guards de Nest funcionan en gateways, pero recuerda que se ejecutan por mensaje, no por
conexión, y que allí se accede al cliente con context.switchToWs().getClient(). Un guard por
mensaje es útil para autorización fina; la autenticación conviene resolverla una vez en la conexión.
11.6.2 Escalado horizontal: el adaptador de Redis no es opcional
Un socket está atado al proceso que lo aceptó. Con una sola instancia todo funciona; en cuanto hay dos, la mitad de tus notificaciones desaparecen y el error es desconcertante porque «en local funciona».
SIN ADAPTADOR (roto) CON ADAPTADOR DE REDIS (correcto)
Ana ──ws──► [ Réplica A ] Ana ──ws──► [ Réplica A ]──┐
Luis ─ws──► [ Réplica B ] Luis ─ws──► [ Réplica B ]──┤
│ pub/sub
POST /tareas/7 ──► Réplica A POST /tareas/7 ──► Réplica A│
server.to('proyecto:1').emit(...) emit(...) ──► publica ──►├──► REDIS
└─► llega SOLO a los sockets de A │
Luis NO recibe nada Réplicas A y B reciben ◄────┘
y entregan a SUS sockets
Síntoma: "a veces llega y a veces no", Resultado: entrega a todos los
proporcional al número de réplicas sockets de la sala, esté donde esté
import { IoAdapter } from '@nestjs/platform-socket.io';
import { createAdapter } from '@socket.io/redis-adapter';
import { createClient } from 'redis';
import { ServerOptions } from 'socket.io';
export class RedisIoAdapter extends IoAdapter {
private adapterConstructor!: ReturnType<typeof createAdapter>;
async connectToRedis(url: string): Promise<void> {
// Dos clientes: uno publica y otro se suscribe. Un cliente en modo
// suscripción no puede ejecutar otros comandos, de ahí el duplicado.
const pubClient = createClient({ url });
const subClient = pubClient.duplicate();
await Promise.all([pubClient.connect(), subClient.connect()]);
this.adapterConstructor = createAdapter(pubClient, subClient);
}
createIOServer(port: number, options?: ServerOptions): unknown {
const server = super.createIOServer(port, options);
server.adapter(this.adapterConstructor);
return server;
}
}
// main.ts
const app = await NestFactory.create(AppModule);
const redisAdapter = new RedisIoAdapter(app);
await redisAdapter.connectToRedis(process.env.REDIS_URL!);
app.useWebSocketAdapter(redisAdapter); // ANTES de app.listen()
transports: ['websocket'] como en el ejemplo, que evita el problema de raíz a cambio de perder
el fallback en redes que bloquean WebSocket.
11.6.3 Server-Sent Events: cuando media conexión basta
Muchas veces solo necesitas que el servidor empuje datos y el cliente no tiene nada que decir: progreso de un informe, contador de notificaciones, precios en directo. Para eso hay una solución mucho más simple que WebSocket, integrada en HTTP: Server-Sent Events. Es texto sobre una respuesta que no se cierra, con reconexión automática por parte del navegador.
import { Sse, MessageEvent } from '@nestjs/common';
import { interval, map, switchMap, takeWhile, from } from 'rxjs';
@Controller('informes')
export class InformesSseController {
/** El método devuelve un Observable<MessageEvent>; Nest mantiene la conexión abierta. */
@Sse(':id/progreso')
@UseGuards(JwtAuthGuard) // los guards funcionan igual que en HTTP normal
progreso(@Param('id') id: string): Observable<MessageEvent> {
return interval(1000).pipe(
switchMap(() => from(this.cola.getJob(id))),
map((job) => ({
// 'data' es obligatorio; 'type' se corresponde con el nombre del evento
// y 'id' permite al navegador reanudar con Last-Event-ID.
type: 'progreso',
data: { progreso: job?.progress ?? 0, estado: job ? 'activo' : 'desconocido' },
})),
takeWhile((e) => (e.data as { progreso: number }).progreso < 100, true),
);
}
}
| Técnica | Dirección | Latencia | Coste servidor | Complejidad | Cuándo elegirla |
|---|---|---|---|---|---|
Polling (setInterval + GET) |
Cliente pregunta | Media del intervalo | Alto: peticiones aunque no haya cambios | Mínima | Prototipos, datos que cambian cada minutos, clientes con red hostil |
| Long-polling | Cliente pregunta y el servidor retiene | Baja | Alto: una conexión y un hilo lógico por cliente | Media | Solo como fallback cuando WebSocket está bloqueado |
| SSE | Solo servidor → cliente | Muy baja | Bajo: una conexión HTTP por cliente | Baja | Notificaciones, progreso, métricas en vivo, tokens de un LLM. Pasa por proxies y CDNs sin configuración especial |
| WebSocket | Bidireccional | Mínima | Bajo por mensaje, alto en memoria por conexión | Alta: adaptador, reconexión, autenticación propia | Chat, colaboración simultánea, juegos, edición concurrente, cualquier flujo con mensajes del cliente frecuentes |
11.6.4 Ejemplo completo: cambios de tarea en tiempo real
Unimos las piezas: el servicio de tareas emite un evento de dominio, un listener lo traduce a difusión por WebSocket y el cliente Angular lo consume como una señal. Nótese que el servicio de tareas no conoce el gateway: si mañana se sustituye Socket.IO, el dominio no cambia.
@Injectable()
export class RealtimeListener {
constructor(private readonly gateway: TareasGateway) {}
@OnEvent('tarea.actualizada', { async: true })
async difundir(e: TareaActualizadaEvent) {
// Enviamos un DTO, nunca la entidad: evita filtrar campos internos
// y las referencias circulares del ORM al serializar.
this.gateway.notificarCambio(e.proyectoId, {
id: e.tareaId,
estado: e.estado,
actualizadaEn: e.fecha.toISOString(),
// Marca de versión: el cliente descarta mensajes desordenados.
version: e.version,
});
}
@OnEvent('informe.listo', { async: true })
async informe(e: { userId: string; url: string }) {
this.gateway.notificarUsuario(e.userId, { tipo: 'informe', url: e.url });
}
}
import { Injectable, signal, inject, DestroyRef } from '@angular/core';
import { io, Socket } from 'socket.io-client';
@Injectable({ providedIn: 'root' })
export class RealtimeService {
private socket?: Socket;
private readonly destroyRef = inject(DestroyRef);
readonly conectado = signal(false);
readonly ultimoCambio = signal<TareaDto | null>(null);
conectar(token: string): void {
this.socket = io('https://api.ejemplo.com/tareas', {
transports: ['websocket'],
auth: { token }, // llega a handshake.auth en el servidor
reconnectionDelay: 1000,
reconnectionDelayMax: 10_000, // retroceso: no martillees un servidor caído
});
this.socket.on('connect', () => this.conectado.set(true));
this.socket.on('disconnect', () => this.conectado.set(false));
this.socket.on('tarea:actualizada', (t: TareaDto) => this.ultimoCambio.set(t));
// Token caducado o revocado: el servidor rechazó el handshake.
this.socket.on('connect_error', (err) => {
if (err.message === 'No autorizado') this.auth.renovarYReconectar();
});
// Liberar SIEMPRE: un socket huérfano sigue recibiendo y filtra memoria.
this.destroyRef.onDestroy(() => this.socket?.disconnect());
}
seguirProyecto(proyectoId: string): Promise<void> {
// El tercer argumento es el ack: la promesa se resuelve con lo que
// devuelva el @SubscribeMessage del servidor.
return new Promise((resolve, reject) => {
this.socket?.emit('proyecto:seguir', { proyectoId }, (res: { ok?: boolean }) =>
res?.ok ? resolve() : reject(new Error('No autorizado')),
);
});
}
}
11.7 GraphQL
11.7.1 Qué problema resuelve y qué problemas trae
REST expone recursos con una forma fija decidida por el servidor. Eso genera tres fricciones cuando el cliente es una aplicación rica: over-fetching (la pantalla necesita el nombre y llegan cuarenta campos), under-fetching (necesita el nombre del autor, que no viene incluido) y múltiples viajes (para pintar una lista con autores y etiquetas hacen falta tres peticiones en cascada, con la latencia sumada).
GraphQL invierte el control: el cliente declara exactamente qué campos quiere, en una sola petición, y el
servidor responde con esa forma. A cambio introduce cuatro problemas nuevos que hay que aceptar con los ojos
abiertos: la caché HTTP deja de funcionar (todo es un POST a /graphql), el
N+1 aparece de forma natural en los resolvers, la seguridad requiere limitar consultas maliciosas
y la complejidad operativa sube (esquema, resolvers, herramientas, monitorización por operación).
11.7.2 Code-first con @nestjs/graphql
Hay dos enfoques: schema-first (escribes SDL y se generan los tipos) y code-first (escribes clases con decoradores y se genera el SDL). En un proyecto Nest, code-first es casi siempre mejor: una única fuente de verdad en TypeScript, sin desincronización entre esquema y código.
import { ObjectType, Field, ID, Int, registerEnumType } from '@nestjs/graphql';
export enum EstadoTarea { PENDIENTE = 'pendiente', HECHA = 'hecha' }
registerEnumType(EstadoTarea, { name: 'EstadoTarea' });
@ObjectType({ description: 'Tarea de un proyecto' })
export class Tarea {
@Field(() => ID) id!: string;
@Field() titulo!: string;
// nullable explícito: en GraphQL todo es opcional salvo que digas lo contrario
@Field({ nullable: true }) descripcion?: string;
@Field(() => EstadoTarea) estado!: EstadoTarea;
@Field(() => Int) comentariosTotal!: number;
// El tipo se resuelve en un @ResolveField: no viaja en la consulta base.
@Field(() => Usuario) propietario!: Usuario;
@Field(() => [Etiqueta]) etiquetas!: Etiqueta[];
}
@InputType()
export class CrearTareaInput {
@Field() @IsString() @Length(3, 120) titulo!: string;
@Field(() => ID) @IsUUID() proyectoId!: string;
}
@Resolver(() => Tarea)
@UseGuards(GqlAuthGuard)
export class TareasResolver {
constructor(private readonly servicio: TareasService) {}
@Query(() => [Tarea], { name: 'tareas' })
listar(@Args('estado', { nullable: true }) estado?: EstadoTarea) {
return this.servicio.listar({ estado });
}
@Mutation(() => Tarea)
crearTarea(
@Args('input') input: CrearTareaInput,
@CurrentUser() user: JwtUser,
) {
return this.servicio.crear(input, user.id);
}
/** Campo calculado: solo se ejecuta si el cliente lo pide. */
@ResolveField(() => Usuario)
propietario(@Parent() tarea: Tarea, @Context() ctx: GqlContext) {
// OJO: aquí nace el N+1 (ver 11.7.3). Con DataLoader, no.
return ctx.loaders.usuarios.load(tarea.propietarioId);
}
@Subscription(() => Tarea, {
// Filtro en el servidor: no envíes a quien no le corresponde.
filter: (payload, variables) =>
payload.tareaActualizada.proyectoId === variables.proyectoId,
})
tareaActualizada(@Args('proyectoId') proyectoId: string) {
return this.pubSub.asyncIterableIterator('tareaActualizada');
}
}
PubSub de graphql-subscriptions pasó de
asyncIterator() a asyncIterableIterator() en la versión 3; con la 2 sigue siendo el
primero. En Apollo Server 4 (el que usa el driver actual) la opción playground: true quedó
obsoleta y el entorno de pruebas se activa con un plugin de página de aterrizaje. Y el
PubSub de serie es en memoria: con varias réplicas necesitas
graphql-redis-subscriptions, por el mismo motivo que los WebSockets necesitan adaptador.
11.7.3 El problema N+1 y DataLoader
Un resolver de campo se ejecuta una vez por elemento del resultado padre. Si la consulta devuelve 50 tareas y cada una resuelve su propietario con una consulta, se ejecutan 51 consultas: una para la lista y 50 para los autores. Con dos niveles de anidamiento, cientos. La API parece correcta en desarrollo con tres filas y se derrumba en producción.
SIN DATALOADER (N+1) CON DATALOADER (2 consultas)
query { tareas { titulo propietario { nombre } } }
1 SELECT * FROM tarea LIMIT 50 1 SELECT * FROM tarea LIMIT 50
2 SELECT * FROM usuario WHERE id = 'a' │ ┌─ load('a') ─┐
3 SELECT * FROM usuario WHERE id = 'b' │ ├─ load('b') ─┤ mismo tick del
4 SELECT * FROM usuario WHERE id = 'a' ◄── │ ├─ load('a') ─┤ event loop: se
… … 50 consultas, muchas repetidas dup │ └─ load('c') ─┘ agrupan y deduplican
51 SELECT * FROM usuario WHERE id = 'z' 2 SELECT * FROM usuario WHERE id IN ('a','b','c')
Latencia: 51 × RTT Latencia: 2 × RTT
DataLoader resuelve dos cosas a la vez: agrupa (batching) las claves solicitadas en el mismo turno del
event loop en una única consulta con IN, y deduplica con una caché por petición. Es
imprescindible que los loaders se creen por petición: si son singletons, la caché mezcla datos entre
usuarios y entre momentos, exactamente el problema de la sección 11.2.
import DataLoader from 'dataloader';
export function crearLoaders(em: EntityManager) {
return {
usuarios: new DataLoader<string, Usuario>(async (ids) => {
const filas = await em.find(Usuario, { id: { $in: [...ids] } });
const porId = new Map(filas.map((u) => [u.id, u]));
// CONTRATO: el array devuelto debe tener el MISMO tamaño y orden
// que las claves recibidas. Si no, DataLoader devuelve valores cruzados.
return ids.map((id) => porId.get(id) ?? new Error(`Usuario ${id} no existe`));
}),
// Relación uno-a-muchos: agrupa hijos por clave del padre.
etiquetasPorTarea: new DataLoader<string, Etiqueta[]>(async (tareaIds) => {
const filas = await em.find(TareaEtiqueta, { tarea: { $in: [...tareaIds] } },
{ populate: ['etiqueta'] });
const grupos = new Map<string, Etiqueta[]>();
for (const f of filas) {
const lista = grupos.get(f.tarea.id) ?? [];
lista.push(f.etiqueta);
grupos.set(f.tarea.id, lista);
}
// Sin resultados debe devolver [], no undefined.
return tareaIds.map((id) => grupos.get(id) ?? []);
}),
};
}
export type GqlContext = { req: Request; loaders: ReturnType<typeof crearLoaders> };
GraphQLModule.forRoot<ApolloDriverConfig>({
driver: ApolloDriver,
autoSchemaFile: 'schema.gql',
sortSchema: true,
// Una fábrica de contexto POR PETICIÓN: aquí nacen y mueren los loaders.
context: ({ req }: { req: Request }) => ({ req, loaders: crearLoaders(orm.em.fork()) }),
introspection: process.env.NODE_ENV !== 'production',
validationRules: [depthLimit(7)],
});
11.7.4 Seguridad: profundidad, complejidad e introspección
- Profundidad. Un esquema con relaciones cíclicas permite
tarea { proyecto { tareas { proyecto { … } } } }hasta agotar el servidor. Limita condepthLimit(n)envalidationRules; entre 7 y 10 suele bastar. - Complejidad. La profundidad no captura el coste de
primeros: 10000. Congraphql-query-complexityse asigna un coste a cada campo y se rechaza la consulta si el total excede un máximo, antes de ejecutar nada. - Paginación obligatoria. Ningún campo de lista debe poder devolver todo. Impón un
limitmáximo en el propio resolver, no confíes en el argumento del cliente. - Introspección desactivada en producción. Es la que permite descargar el esquema completo y facilita la búsqueda de campos olvidados. Desactiva también el explorador gráfico.
- Consultas persistidas (allowlist). El nivel máximo de protección: solo se aceptan consultas previamente registradas, identificadas por su hash. Elimina de golpe todo el vector de consultas arbitrarias.
- Autorización por campo. Un guard a nivel de resolver no protege un
@ResolveFieldal que se llega por otro camino del grafo. Autoriza en la capa de servicio o con directivas, no solo en la entrada.
| Criterio | REST | GraphQL | tRPC |
|---|---|---|---|
| Contrato | OpenAPI (generado o escrito) | Esquema tipado y explorable | Los tipos de TypeScript, sin generación |
| Clientes distintos | Endpoints a medida o campos de más | Cada cliente pide lo que necesita | Pensado para un cliente TS propio |
| Caché HTTP | Nativa y gratuita (CDN, ETag) | Hay que renunciar a ella o usar consultas persistidas con GET | Igual que REST si expones GET |
| Riesgo N+1 | Bajo: tú escribes la consulta | Alto: exige DataLoader | Bajo |
| Consumidores externos | Ideal: universal | Viable | No: acopla al cliente TypeScript |
| Coste de entrada | Mínimo | Alto | Muy bajo en un monorepo TS |
| Elígelo cuando… | API pública, integraciones, CRUD, ficheros, caché importante | Muchos clientes heterogéneos, grafos de datos profundos, pantallas muy variables | Frontend y backend TypeScript del mismo equipo y repositorio |
HttpClient y DTOs compartidos (sección 11.11) obtiene el tipado extremo a extremo sin
GraphQL. Adoptarlo solo por «pedir menos campos» rara vez compensa; hazlo cuando tengas varios consumidores
con necesidades divergentes o un grafo de datos genuinamente profundo. Y nada impide combinarlo: GraphQL para
la aplicación interna, REST para las integraciones de terceros.
11.8 Microservicios
11.8.1 Cuándo NO usarlos, que es casi siempre
Los microservicios resuelven un problema organizativo: permitir que muchos equipos despliegen sin coordinarse. No resuelven problemas de rendimiento —añadir latencia de red nunca hace nada más rápido— ni de calidad de código: un monolito mal estructurado se convierte en varios servicios mal estructurados, ahora con red en medio.
Consistencia eventual. Desaparecen las transacciones ACID entre módulos. Lo que era un
flush() se convierte en una saga con compensaciones. Todo el código que asumía «o todo o nada»
debe reescribirse asumiendo estados intermedios visibles.
Depuración distribuida. Un error deja de ser un stack trace y pasa a ser una investigación en cinco servicios. Sin trazas distribuidas (OpenTelemetry, capítulo 13) e identificadores de correlación es literalmente inviable.
Coste operativo real. Cada servicio necesita su despliegue, su CI, sus alertas, sus secretos, su base de datos, su versionado de contratos y su plan de compatibilidad hacia atrás. Multiplica por N.
Fallos parciales. Con cinco servicios al 99,9 % de disponibilidad en serie, el conjunto baja al 99,5 %. Aparecen modos de fallo que no existían: timeouts, reintentos duplicados, cascadas.
Respuesta correcta por defecto: monolito modular. Un solo despliegue, módulos de Nest con fronteras explícitas, cada uno con su propio esquema de base de datos y comunicación entre módulos únicamente por interfaces públicas o eventos. Cuando un módulo necesite escalar o cambiar de ritmo, ya está listo para extraerse, y probablemente para entonces nunca haga falta.
11.8.2 Transportes de @nestjs/microservices
Nest abstrae el transporte: el mismo @MessagePattern funciona sobre TCP o sobre Kafka. La
elección, sin embargo, no es un detalle: determina las garantías de entrega.
| Transporte | Modelo | Persistencia | Caso de uso típico | Contraindicación |
|---|---|---|---|---|
| TCP | Petición-respuesta directa | Ninguna | Pruebas, dos servicios internos, monorepo | Producción con descubrimiento dinámico o muchos consumidores |
| Redis | Pub/sub | Ninguna: si nadie escucha, el mensaje se pierde | Notificaciones efímeras cuando ya tienes Redis | Cualquier cosa que no se pueda perder |
| NATS | Pub/sub y petición-respuesta nativa; colas de trabajo | Opcional con JetStream | Comunicación interna de muy baja latencia y operación sencilla | Necesidad de reproducir el histórico sin JetStream |
| MQTT | Pub/sub con niveles de calidad de servicio | Según QoS y broker | IoT, dispositivos con red intermitente, mensajes pequeños | Mensajería de negocio compleja |
| RabbitMQ | Colas con enrutado, confirmación y DLQ | Sí, con colas durable | Trabajo distribuido con reintentos y prioridades: el «todoterreno» | Cientos de miles de mensajes por segundo |
| Kafka | Log particionado con desplazamientos | Sí, y reproducible | Analítica, auditoría, event sourcing, alto caudal, varios consumidores del mismo flujo | Petición-respuesta y equipos pequeños: operarlo es un trabajo en sí mismo |
| gRPC | RPC con contrato .proto sobre HTTP/2 | Ninguna | Llamadas internas tipadas, alto rendimiento, streaming bidireccional, políglota | Navegadores sin pasarela y contratos que cambian a diario |
11.8.3 @MessagePattern, @EventPattern y ClientProxy
@Controller()
export class FacturasController {
/** Petición-respuesta: el emisor ESPERA un valor. */
@MessagePattern({ cmd: 'facturas.obtener' })
async obtener(@Payload() data: { id: string }): Promise<FacturaDto> {
const f = await this.servicio.buscar(data.id);
// Los errores deben viajar como RpcException para que el
// cliente reciba un error estructurado y no un objeto vacío.
if (!f) throw new RpcException({ code: 'NOT_FOUND', message: 'No existe' });
return f;
}
/** Evento: notificación sin respuesta. Puede llegar más de una vez. */
@EventPattern('pedido.confirmado')
async alConfirmar(@Payload() e: PedidoConfirmado, @Ctx() ctx: RmqContext) {
// Con noAck: false, TÚ decides cuándo se confirma el mensaje.
const canal = ctx.getChannelRef();
const mensaje = ctx.getMessage();
try {
await this.servicio.emitirFactura(e.pedidoId); // idempotente
canal.ack(mensaje);
} catch (err) {
// requeue en false: va a la dead letter queue en lugar de
// volver a la cola para siempre (bucle infinito de veneno).
canal.nack(mensaje, false, false);
}
}
}
@Module({
imports: [ClientsModule.register([{
name: 'FACTURAS',
transport: Transport.RMQ,
options: {
urls: [process.env.RMQ_URL!],
queue: 'facturas',
queueOptions: { durable: true },
noAck: false,
},
}])],
})
export class PedidosModule {}
@Injectable()
export class PedidosService {
constructor(@Inject('FACTURAS') private readonly facturas: ClientProxy) {}
async detalle(id: string) {
// send() devuelve un Observable FRÍO: sin subscribe no se envía nada.
// Un timeout SIEMPRE: sin él, un servicio lento cuelga esta petición
// hasta agotar el pool de conexiones del gateway.
return firstValueFrom(
this.facturas.send({ cmd: 'facturas.obtener' }, { id }).pipe(
timeout(3_000),
retry({ count: 2, delay: 300 }),
catchError((err) => {
if (err instanceof TimeoutError) {
// Degradar es mejor que fallar: devuelve lo que tengas.
return of({ id, estado: 'desconocido', degradado: true });
}
return throwError(() => new ServiceUnavailableException('Facturas no disponible'));
}),
),
);
}
notificar(e: PedidoConfirmado) {
// emit() es fuego y olvido: no esperes confirmación de proceso.
this.facturas.emit('pedido.confirmado', e);
}
}
1. Serialización. Todo lo que cruza el transporte es JSON (o Protobuf en gRPC): las fechas llegan
como cadenas, los Map, Set y Buffer se degradan y las clases pierden
sus métodos. Define DTOs de transporte explícitos y valida a la entrada del microservicio igual que en HTTP,
porque el emisor es tan poco de fiar como un navegador.
2. Timeouts en cascada. Si el gateway espera 30 s y el servicio A espera 30 s por el B, un fallo en B mantiene ocupada la petición del usuario un minuto. Presupuesta el tiempo: el gateway 3 s, A 1,5 s, B 700 ms. Cada salto debe tener menos presupuesto que quien le llama.
3. Kafka y la petición-respuesta. Sobre Kafka, el patrón petición-respuesta necesita un tema de
respuesta y suscribirse explícitamente en onModuleInit con
subscribeToResponseOf(). Si no lo haces, la llamada nunca se resuelve y el síntoma es un
timeout sin más pistas.
4. El transporte Redis no persiste. Está construido sobre pub/sub: si el consumidor está reiniciándose, el mensaje se pierde sin error en el emisor. Para eventos de negocio usa RabbitMQ, NATS con JetStream o Kafka; Redis solo para lo verdaderamente efímero.
11.8.4 Aplicaciones híbridas
Un mismo proceso puede atender HTTP y escuchar mensajes. Es la forma habitual de migrar sin big bang: el servicio conserva su API REST y empieza a consumir eventos.
const app = await NestFactory.create(AppModule);
app.useGlobalPipes(new ValidationPipe({ whitelist: true, transform: true }));
// inheritAppConfig hace que pipes, filtros e interceptores globales
// se apliquen TAMBIÉN a los mensajes; sin él, los mensajes entran sin validar.
app.connectMicroservice<MicroserviceOptions>(
{ transport: Transport.RMQ, options: { urls: [process.env.RMQ_URL!], queue: 'pedidos' } },
{ inheritAppConfig: true },
);
await app.startAllMicroservices();
await app.listen(3000); // el mismo proceso: mismo pool de BD, mismo event loop
main) para aislar los recursos.
11.8.5 Patrones distribuidos imprescindibles
- API Gateway. Un único punto de entrada que autentica, limita, enruta y compone respuestas. Evita que el cliente conozca la topología interna. Riesgo: convertirse en un monolito con lógica de negocio; debe quedarse en orquestación y traducción.
- Outbox transaccional. El problema: no puedes escribir en la base de datos y publicar en el broker
atómicamente; si publicas antes del commit, anuncias algo que puede no ocurrir, y si publicas después,
puedes perder el mensaje. Solución: en la misma transacción inserta la fila y un registro en la tabla
outbox; un proceso aparte lee esa tabla y publica, marcando lo enviado. Así el mensaje se publica «al menos una vez» y siempre coherente con el estado. Es el patrón que hace innecesario el commit en dos fases. - Idempotencia. Igual que en las colas (11.3.5): identificador de mensaje, tabla de procesados y operaciones idempotentes. En un sistema distribuido es obligatoria, no recomendable.
- Reintentos con retroceso exponencial y jitter. Reintentar de inmediato y a la vez es la receta de la manada atronadora: el servicio que se está recuperando vuelve a caer.
- Circuit breaker. Si un servicio falla el 50 % de las últimas 20 llamadas, el interruptor «abre» y
las siguientes fallan de inmediato sin llamar. Pasado un tiempo entra en «semiabierto» y prueba con una
llamada. Evita gastar hilos y timeouts en algo que se sabe roto y da margen a la recuperación. Suele
implementarse con
opossumenvolviendo el cliente. - Dead letter queue. Los mensajes que fallan repetidamente salen del flujo a una cola aparte. Sin DLQ, un solo mensaje «venenoso» bloquea el consumo indefinidamente. Con DLQ hay que vigilarla: una DLQ que nadie mira es una papelera.
- Saga. La sustituta de la transacción distribuida: una secuencia de transacciones locales, cada una con su compensación. No hay rollback global; hay operaciones inversas.
SAGA COREOGRAFIADA (cada servicio reacciona a eventos)
PEDIDOS PAGOS ALMACÉN ENVÍOS
│ pedido.creado │ │ │
├─────────────────►│ cobra │ │
│ ├──pago.ok─────────►│ reserva stock │
│ │ ├──stock.ok───────►│ crea envío
│ │ │ ├──envio.ok──┐
│◄─────────────────┴───────────────────┴──────────────────┴────────────┘
│ pedido.completado (sin coordinador central)
COMPENSACIÓN (el almacén no tiene stock)
│ │ │ stock.fallido │
│ │◄──────────────────┤ │
│ ├─ reembolsa (compensación de "cobra") │
│◄─────────────────┤ pago.reembolsado │
├─ marca pedido como cancelado │
SAGA ORQUESTADA (un coordinador dirige) ┌─────────────────┐
│ ORQUESTADOR │
1. cobrar ──► 2. reservar ──► 3. enviar │ (máquina de │
▲ │ │ estados) │
└──── compensar si falla ◄─────────────────┤ guarda el │
│ estado de cada │
Coreografía: menos acoplamiento, flujo │ saga en curso │
difícil de ver. Orquestación: flujo └─────────────────┘
explícito y depurable, un componente más.
La elección entre coreografía y orquestación no es estética: decide dónde vive la complejidad. En la coreografía cada servicio conoce solo los eventos que consume y los que publica, el acoplamiento es mínimo y, en cambio, el flujo completo no está escrito en ningún sitio: para saber qué ocurre al confirmar un pedido hay que leer cuatro repositorios distintos y confiar en que la documentación esté al día. En la orquestación hay una máquina de estados explícita, consultable y depurable, al precio de un componente más que desplegar, escalar y hacer tolerante a fallos. Regla práctica: hasta tres pasos, coreografía; a partir de cuatro pasos con compensaciones, plazos y estados intermedios visibles para el usuario, orquestación.
Un microservicio no es un módulo con red en medio: es un límite transaccional nuevo. En el momento
en que dos datos que antes se escribían en un flush() pasan a vivir en dos bases de datos, has
cambiado un problema de código por un problema de protocolos: idempotencia, outbox, saga, compensaciones,
presupuestos de tiempo y observabilidad distribuida. Todo eso es trabajo permanente, no una migración
puntual.
Por eso el orden correcto de las decisiones es: primero módulos con fronteras reales dentro del mismo despliegue, después extracción de un solo módulo cuando exista una razón medible (un equipo que se bloquea, un componente que necesita escalar diez veces más que el resto, un requisito legal de aislamiento), y solo entonces el transporte. Si empiezas eligiendo Kafka, has decidido la respuesta antes de conocer la pregunta.
11.9 Errores comunes y cómo solucionarlos
Los fallos de este capítulo comparten una firma característica: casi ninguno se manifiesta en desarrollo. Con una réplica, un usuario y tres filas en la base de datos, la caché es coherente, el cron se ejecuta una vez, todos los sockets están en el mismo proceso y no hay N+1 apreciable. Aparecen el día del despliegue a producción, y por eso conviene conocerlos antes de sufrirlos.
| Síntoma o error | Causa real | Solución |
|---|---|---|
| Caché (11.2) | ||
| Un usuario ve el listado de otro usuario tras recargar la página | La clave del CacheInterceptor es la URL y no incluye la identidad: /tareas es la misma cadena para todos |
Incluir el identificador del usuario en la clave sobrescribiendo trackBy(), o no cachear ese endpoint. Devolver además Cache-Control: private, no-store para que ninguna CDN guarde copia |
| La misma petición devuelve datos distintos según el momento, sin patrón aparente | Caché en memoria del proceso con varias réplicas: hay N cachés independientes y el balanceador reparte | Mover la caché a Redis, que es una única fuente de verdad. Si se mantiene un nivel en memoria (L1), bajar su TTL al orden de segundos para que la incoherencia sea irrelevante |
| El usuario edita algo, ve el cambio en el formulario y el listado sigue mostrando el valor antiguo | Invalidación olvidada: la ruta de escritura no borra ni versiona las claves derivadas (listados con filtros, contadores, agregados) | Versión por usuario o por entidad en la clave e incremento al escribir (11.2.7), disparado por evento de dominio, y TTL corto como red de seguridad para que un olvido dure minutos y no meses |
| Un fallo puntual de la base de datos se sirve durante media hora a todos los usuarios | Se cacheó la respuesta de error o un objeto vacío con el TTL normal | Cachear únicamente resultados correctos. Para los «no encontrado» legítimos, un marcador de ausencia con TTL de 15-30 s que amortigua el cache penetration sin fijar el error |
| «La caché no funciona»: los aciertos son casi cero pese a repetir la misma petición | ttl: 60 copiado de una guía antigua: desde cache-manager v5 el TTL está en milisegundos, así que caduca en 60 ms |
Comprobar la versión instalada y expresar los TTL en milisegundos con separadores legibles (60_000). Instrumentar la tasa de aciertos: sin métrica no se detecta |
| La base de datos se satura de golpe cada vez que caduca la clave más consultada | Estampida de caché: cientos de peticiones concurrentes encuentran el hueco y lanzan la misma consulta costosa | Single-flight por clave dentro del proceso, cerrojo distribuido SET NX PX entre réplicas, TTL con jitter y refresco proactivo stale-while-revalidate |
El proceso muere con JavaScript heap out of memory tras unas horas |
Caché en memoria con claves dinámicas y sin max: nada expulsa entradas antes de su TTL |
Fijar siempre max para que actúe la política LRU, acotar el tamaño de los valores cacheados y vigilar el uso de heap del proceso |
| Reiniciar Redis tumba toda la API con errores 500 | Las llamadas a get y set no están protegidas: una caché opcional se convirtió en punto único de fallo |
Envolver en try/catch y degradar a la fuente original registrando un aviso. La aplicación debe funcionar, más lenta, con la caché caída |
| Redis se queda sin memoria y empieza a expulsar claves de negocio | Entradas guardadas sin TTL (ttl: 0) mezcladas con colas y sesiones en la misma instancia |
TTL en todo lo que sea caché, política maxmemory-policy adecuada y separación de instancias o bases entre caché (volátil) y colas (persistente) |
Un valor legítimo 0, false o cadena vacía nunca se considera acierto |
Comprobación if (!hit), que confunde «no está en caché» con «está y vale falsy» |
Comparar explícitamente contra undefined y null, o envolver el valor en un objeto ({ v: valor }) para distinguir ausencia de contenido |
| Colas y tareas programadas (11.3) | ||
| El cliente recibe dos correos de recibo, o peor, se le cobra dos veces | El consumidor no es idempotente y la cola garantiza entrega «al menos una vez»: un worker que muere tras el efecto y antes de confirmar provoca el reintento | Comprobación de estado al entrar, escritura condicional (UPDATE … WHERE estado = 'pendiente') comprobando filas afectadas y clave de idempotencia en el tercero para los efectos externos |
| Trabajos que llevan semanas fallando y nadie se había dado cuenta | El conjunto failed hace de cola de fallidos, pero no hay panel, ni alerta, ni responsable |
Alerta sobre el número de fallidos y sobre la antigüedad del más viejo, panel protegido con autenticación y un procedimiento escrito de reintento tras corregir la causa |
| Al desplegar se pierden trabajos que estaban a medio procesar | El contenedor recibe SIGTERM y el proceso muere sin cerrar el worker ni renovar el lock |
Apagado ordenado: enableShutdownHooks(), cierre explícito del worker para que termine el trabajo en curso y terminationGracePeriodSeconds mayor que el trabajo más largo. Los trabajos con el lock vencido se recuperan como stalled, pero solo si son idempotentes |
| Un servicio externo con una incidencia pasa de degradado a completamente caído en cuanto reintentamos | Reintentos inmediatos y sin retroceso: todos los trabajos vuelven a la vez y multiplican la carga sobre algo que ya estaba mal | backoff: { type: 'exponential', delay } con jitter, limiter por cola para respetar la cuota del proveedor y UnrecoverableError para los errores que jamás se arreglarán reintentando |
| El worker procesa el trabajo y la entidad «no existe» aunque acaba de crearse | Se encoló dentro de la transacción, antes del commit: el worker fue más rápido que la base de datos | Encolar después del commit, o usar el patrón outbox si el mensaje no puede perderse. Un delay de unos segundos disimula el problema pero no lo resuelve |
Avisos constantes de trabajos stalled y trabajos ejecutados dos veces |
El worker bloquea el event loop con trabajo intensivo de CPU y no renueva su cerrojo, así que la cola lo considera muerto | Sacar el cálculo a worker_threads o a un proceso aparte, reducir la concurrencia y ajustar lockDuration. Y, en cualquier caso, mantener el consumidor idempotente |
| La memoria de Redis crece sin parar aunque las colas estén vacías de trabajo pendiente | Falta removeOnComplete: el histórico de trabajos terminados se conserva indefinidamente |
removeOnComplete: { age, count } y removeOnFail: { age } en las opciones por defecto, conservando los fallidos más tiempo que los completados |
| La facturación mensual se emite tres veces y hay que anular facturas a mano | @Cron se registra en cada proceso: con tres réplicas hay tres planificadores idénticos |
Cerrojo distribuido en Redis, despliegue dedicado con una sola réplica o trabajos repetibles de BullMQ, cuya programación vive en Redis. El planificador solo debe encolar, nunca ejecutar el trabajo pesado |
| Errores de timeout al adquirir conexión que parecen fallos de la base de datos | Concurrencia real igual a concurrency × réplicas, muy por encima del tamaño del pool de conexiones |
Dimensionar la concurrencia contra el recurso más escaso (casi siempre el pool), no contra la CPU, y separar el despliegue del worker del de la API para que no compitan |
| WebSockets y tiempo real (11.6) | ||
| Las notificaciones llegan «a veces»: aproximadamente a uno de cada N usuarios | Cada socket vive atado al proceso que lo aceptó y server.to(...).emit(...) solo alcanza a los sockets de esa réplica |
Adaptador de Redis (@socket.io/redis-adapter) registrado con useWebSocketAdapter() antes de listen(). Con GraphQL, el equivalente es sustituir el PubSub en memoria por graphql-redis-subscriptions |
| Conexiones que se caen tras un rato largo abiertas y no se recuperan | El token se validó solo en el handshake y ha caducado; la reconexión reutiliza el mismo token muerto | Renovar el token en el cliente y reconectar con el nuevo al recibir connect_error. Si la sesión debe poder revocarse en caliente, revalidar periódicamente y cerrar el socket cuando deje de ser válida |
| Tras una desconexión breve, la interfaz muestra datos incompletos o incoherentes | La vista se construyó solo con los mensajes recibidos: lo emitido mientras no había conexión se perdió | Cargar el estado por HTTP, aplicar encima los mensajes en vivo y recargar al reconectar. Añadir número de versión por entidad y descartar mensajes con versión anterior a la ya aplicada |
| Errores de sesión desconocida en el handshake detrás del balanceador | El transporte de long-polling reparte varias peticiones HTTP de la misma sesión entre réplicas distintas | Afinidad de sesión por cookie en el balanceador, o forzar transports: ['websocket'] para eliminar el problema de raíz |
| Un usuario recibe eventos de un proyecto al que no tiene acceso | El manejador de @SubscribeMessage une al cliente a la sala que este pide sin comprobar permisos: es un IDOR sobre salas |
Autorizar en el momento de unirse a la sala y filtrar también en el emisor. En GraphQL, usar el filter de @Subscription para no enviar a quien no corresponde |
| GraphQL (11.7) | ||
| Una consulta que devuelve 50 elementos genera 51 o más consultas SQL | N+1: los @ResolveField se ejecutan una vez por elemento del resultado padre |
DataLoader por relación, creado por petición en la fábrica de contexto, agrupando con IN y respetando el contrato de devolver un array del mismo tamaño y orden que las claves |
| Con DataLoader activo, un usuario recibe ocasionalmente datos de otro | Los loaders se registraron como singleton: su caché interna vive todo el proceso y mezcla peticiones | Crear los loaders en context: ({ req }) => …, uno por petición, sobre un EntityManager propio (em.fork()), y no reutilizarlos nunca entre peticiones |
| Una sola consulta consume toda la CPU y tumba el servidor | Anidamiento cíclico o listas sin límite: el esquema permite pedir un grafo de coste exponencial | depthLimit entre 7 y 10 en validationRules, análisis de complejidad con graphql-query-complexity, paginación obligatoria con tope impuesto por el servidor e introspección desactivada en producción. Para el máximo nivel, consultas persistidas |
| La monitorización no ve ningún error y los usuarios sí los ven | GraphQL responde 200 OK con los fallos dentro del array errors del cuerpo |
Instrumentar los errores del propio GraphQL con un plugin de Apollo y alertar sobre ellos; en el cliente, comprobar siempre errors además del estado HTTP. Normalizar los códigos con formatError sin filtrar detalles internos |
| Microservicios (11.8) | ||
| Una petición del usuario se queda colgada hasta que el navegador desiste | send() devuelve un observable sin límite de tiempo: si el destinatario no contesta, nadie corta |
timeout() en toda llamada entre servicios, con presupuesto decreciente por salto (gateway 3 s, servicio A 1,5 s, servicio B 700 ms) y respuesta degradada en catchError |
| Un servicio secundario lento deja fuera de servicio a toda la aplicación | Fallo en cascada: las peticiones esperando agotan conexiones y memoria del que llama | Circuit breaker (opossum) para dejar de llamar a lo que se sabe roto, aislamiento de recursos por dependencia y degradación explícita del caso de uso cuando esa parte no es imprescindible |
| El mismo evento se procesa dos veces y se duplican filas o notificaciones | Los brokers entregan «al menos una vez»; un nack con reencolado o un reinicio reproducen el mensaje |
Identificador único de mensaje y tabla de procesados con clave primaria (tipo, clave): la violación de unicidad significa «ya hecho». Publicar con outbox transaccional para no perder ni inventar mensajes |
| El estado final de una entidad es incorrecto aunque todos los mensajes se procesaron | El orden de entrega no está garantizado entre particiones ni entre consumidores concurrentes: llegó primero la actualización y después la creación | Número de versión o marca temporal en el mensaje y descarte de los atrasados; clave de partición estable (por ejemplo el identificador de la entidad) en Kafka para que todo lo de una entidad vaya en orden; y consumo con concurrencia 1 por clave cuando el orden sea imprescindible |
Con Kafka, send() nunca se resuelve y solo se ve un timeout |
Falta suscribirse al tema de respuesta con subscribeToResponseOf() en onModuleInit |
Declarar las respuestas esperadas al inicializar el cliente y, en general, preferir eventos sobre petición-respuesta cuando el transporte es un log como Kafka |
| Mensajes que desaparecen sin ningún error en el emisor | Transporte Redis, que es pub/sub puro: si el consumidor está reiniciándose, nadie recibe el mensaje y nadie se queja | RabbitMQ con colas durable, NATS con JetStream o Kafka para todo lo que sea negocio. Reservar Redis para lo verdaderamente efímero, como una presencia o un indicador de escritura |
| Un mensaje defectuoso bloquea el consumo de la cola indefinidamente | Mensaje «venenoso»: falla siempre y vuelve a la cola con nack(msg, false, true), formando un bucle infinito |
Límite de reintentos y salida a una dead letter queue (nack(msg, false, false)), con alerta sobre su tamaño. Una DLQ que nadie vigila es una papelera con datos de negocio dentro |
11.10 Buenas y malas prácticas
Haz esto
- Mide antes de añadir una capa. Muchos «necesito Redis» son en realidad «me falta un índice» o «traigo cien columnas para mostrar tres». Optimizar la consulta no introduce estado duplicado; cachear, sí.
- Combina TTL corto con invalidación explícita. La invalidación da la frescura que el usuario espera y el TTL garantiza que un olvido tuyo se corrija solo en minutos. Depender únicamente de la invalidación es depender de que el código no tenga huecos.
- Mete la identidad y todos los parámetros en la clave de caché. Si la respuesta depende de quién pregunta o de qué filtros aplica, eso forma parte de la clave; normaliza el orden de los filtros para que peticiones equivalentes compartan entrada.
- Trata la caché como opcional y la fuente como obligatoria. Captura los errores del almacén, degrada a la base de datos y registra el aviso: una caída de Redis debe traducirse en latencia, nunca en errores 500.
- Saca de la petición todo lo que el usuario no necesita ver. Correos, webhooks salientes, sincronizaciones, informes y generación de ficheros van a una cola; la petición hace solo lo transaccional y responde
202con un identificador de seguimiento. - Escribe cada consumidor como si fuera a ejecutarse dos veces, porque lo hará. Comprobación de estado, escritura condicional con filas afectadas y clave de idempotencia en los terceros: es la única garantía real de «exactamente una vez» observable.
- Encola identificadores, no objetos, y siempre después del commit. Los datos viajan como JSON por Redis, donde quedan visibles durante días, así que el worker debe recargar lo que necesite en lugar de recibir entidades y datos personales.
- Usa el planificador solo para encolar. Un
@Cronprotegido con cerrojo distribuido que añade trabajos a una cola te da reintentos, trazabilidad y una sola ejecución aunque haya veinte réplicas. - Autentica el WebSocket en el handshake y autoriza cada sala. Rechazar antes de establecer la conexión evita sockets anónimos consumiendo memoria, y comprobar el permiso al unirse a una sala evita el equivalente en tiempo real de un IDOR.
- Empieza por Server-Sent Events y sube a WebSocket cuando haga falta. Si el flujo es del servidor al cliente, SSE se resuelve con un controlador normal, guards normales y sin adaptador ni sesiones adheridas.
- Crea los DataLoader por petición y limita profundidad y complejidad desde el primer día. Un esquema sin límites es una denegación de servicio a la espera de que alguien la descubra, y un loader compartido es una fuga de datos entre usuarios.
- Pon un timeout con presupuesto decreciente en cada llamada entre servicios. Cada salto debe disponer de menos tiempo que quien le llama, y el caso de uso debe saber qué devolver cuando ese tiempo se agota.
Evita esto
- Activar el
CacheInterceptora nivel de controlador «por si acaso». En cuanto un endpoint devuelva datos que dependan del usuario, la clave basada en la URL servirá la respuesta de Ana a Luis, y eso es un incidente de privacidad, no un fallo de rendimiento. - Cachear en memoria del proceso y creer que se puede invalidar con un evento interno. Los eventos de
eventemitter2no salen del proceso: con tres réplicas invalidas una de cada tres cachés y el resultado es incoherencia aleatoria. - Recorrer claves con
KEYS *para invalidar. En un Redis con millones de claves bloquea el servidor para todos los clientes. La alternativa es una clave versionada que se invalida incrementando un contador. - Poner TTL largos «para que vaya más rápido» sin preguntar cuánta desactualización tolera el dato. La decisión es de producto: si nadie sabe responder cuántos segundos son aceptables, todavía no se puede cachear.
- Usar una cola para algo que debe fallar con el caso de uso. Descontar el stock al confirmar un pedido pertenece a la transacción, no a un trabajo asíncrono; si se encola, el pedido se confirma y el stock puede no bajar nunca.
- Dejar
attemptsalto en un consumidor no idempotente. Cada reintento repite el efecto secundario: cinco intentos de un cobro sin clave de idempotencia son cinco cargos posibles en la tarjeta del cliente. - Reintentar de inmediato, sin retroceso y todos a la vez. Es la manada atronadora: el servicio que estaba recuperándose vuelve a caer justo cuando empezaba a responder.
- Confiar el «una sola vez» al
jobIddel productor. Solo deduplica mientras el trabajo siga existiendo en Redis; en cuantoremoveOnCompletelo borra, el mismo identificador vuelve a aceptarse. - Construir la vista únicamente con los mensajes del socket. Quien se conecta tarde, pierde la red un instante o deja la pestaña en segundo plano se queda con datos incompletos y sin saberlo.
- Meter el token en la query string del handshake. Acaba en los registros de acceso del proxy y en los informes de error; el sitio correcto es
handshake.auth. - Adoptar GraphQL por «pedir menos campos» y descubrir después que has perdido la caché HTTP. Con un único cliente TypeScript propio, DTOs compartidos sobre REST dan el mismo tipado sin renunciar a
ETag, CDN niGET. - Partir el monolito para «escalar» o para «tener código más limpio». Añadir red nunca hace nada más rápido, y un diseño confuso repartido en seis servicios sigue siendo confuso, ahora con consistencia eventual y depuración distribuida.
11.11 Preguntas frecuentes
¿Cuándo debo cachear y cuándo no?
¿Qué diferencia hay entre la caché de aplicación y la caché HTTP con ETag?
ETag o Last-Modified evita transmitir: el cliente manda If-None-Match con la huella que ya tiene y el servidor responde 304 Not Modified sin cuerpo, aunque para calcular ese ETag normalmente haya tenido que hacer el trabajo. El escenario ideal es apilarlas: Redis te da la respuesta en un milisegundo, calculas su huella y devuelves un 304 de trescientos bytes. Dos precauciones importantes: en respuestas que dependen del usuario hay que enviar Cache-Control: private (y Vary sobre la cabecera relevante) para que ninguna CDN ni proxy intermedio guarde una copia compartida; y recuerda que una vez enviado un max-age largo al navegador no hay forma de retirarlo, mientras que una entrada de Redis puedes borrarla cuando quieras.¿Cómo se invalida una caché correctamente?
KEYS * para localizar qué borrar: en producción bloquea Redis para todos los clientes. Y ten presente que la incoherencia entre réplicas solo se elimina de verdad centralizando la caché; un nivel en memoria no se puede invalidar desde otro proceso, así que su TTL debe medirse en segundos.¿Cuándo uso una cola en lugar de una llamada síncrona?
¿Qué significa exactamente que un trabajo sea idempotente y cómo se consigue?
estado = 'pagado') lo es; incrementar un contador no. Segunda, escritura condicional: UPDATE … WHERE id = ? AND estado = 'pendiente' y comprobar las filas afectadas, que es una transición de estado atómica sin cerrojos explícitos. Tercera, tabla de procesados con clave primaria (tipo, clave), donde una violación de unicidad significa «ya hecho»; sirve para cualquier efecto. Cuarta, clave de idempotencia en el tercero, imprescindible cuando el efecto está fuera de tu sistema: las pasarelas de pago la ofrecen precisamente para esto y devuelven el cargo original en lugar de crear otro. Y una advertencia: el jobId del productor no da idempotencia, solo evita duplicados mientras ese trabajo siga existiendo en Redis.¿Llamada directa, evento interno o cola? Dame un criterio que pueda aplicar sin pensarlo mucho.
@nestjs/event-emitter es perfecto y gratis. Y si la reacción forma parte de la corrección del caso de uso —debe fallar todo junto o no fallar nada—, no es un evento ni un trabajo: es una llamada directa dentro de la misma transacción. Recuerda las dos limitaciones del evento interno, que se olvidan constantemente: no persiste (si el proceso muere entre el emit y el listener, no queda ni rastro) y no cruza réplicas, porque es memoria del proceso actual.¿Cómo se escalan los WebSockets a varias instancias?
server.to('sala').emit(...) solo alcanza a los sockets de esa réplica. Con dos instancias, la mitad de las notificaciones desaparecen; el síntoma clásico es «a veces llega y a veces no», proporcional al número de réplicas. La solución es un adaptador que propague las emisiones entre procesos: con Socket.IO, @socket.io/redis-adapter registrado mediante app.useWebSocketAdapter() antes de listen(), con dos clientes de Redis (uno publica y otro se suscribe, porque un cliente en modo suscripción no puede ejecutar otros comandos). Con las suscripciones de GraphQL el problema es idéntico y la solución equivalente: sustituir el PubSub en memoria por graphql-redis-subscriptions. Falta un segundo detalle: si mantienes activo el transporte de long-polling, necesitas afinidad de sesión en el balanceador, porque varias peticiones HTTP de la misma sesión deben caer en la misma réplica; forzar transports: ['websocket'] elimina ese requisito a cambio de perder el mecanismo de reserva en redes que bloquean WebSocket. Y ten en cuenta que el adaptador resuelve el reparto, no la memoria: cada conexión abierta cuesta, así que las decenas de miles de sockets exigen dimensionar réplicas y ajustar límites de descriptores de fichero.¿Server-Sent Events, WebSockets o sondeo?
setInterval es defendible: cuesta cero en complejidad y funciona en cualquier red, a costa de peticiones inútiles cuando no hay cambios. Si el flujo va únicamente del servidor al cliente —progreso de un informe, contador de notificaciones, precios, la salida de un modelo de lenguaje token a token—, usa SSE: es un controlador normal con @Sse que devuelve un Observable, hereda los guards de HTTP, reconecta solo en el navegador, atraviesa proxies sin configuración especial y no necesita adaptador ni sesiones adheridas. Reserva WebSocket para lo genuinamente bidireccional y frecuente: chat, colaboración simultánea, edición concurrente, juegos, indicadores de escritura. La recomendación de arquitecto es empezar por SSE, porque la mayoría de los «necesitamos WebSockets» son en realidad «el servidor debe avisar al cliente». Limitaciones de SSE que conviene conocer: solo texto UTF-8, sobre HTTP/1.1 el navegador limita a unas seis conexiones por dominio (con HTTP/2 el problema desaparece) y algunos proxies requieren desactivar explícitamente el buffering.¿REST o GraphQL? ¿Cuáles son los compromisos reales?
POST a /graphql y ni la CDN ni el ETag te sirven (se recupera parcialmente con consultas persistidas servidas por GET); el N+1 aparece de forma natural en los resolvers y obliga a DataLoader en todas las relaciones; la seguridad exige límites de profundidad y complejidad porque el cliente construye la consulta; la monitorización tiene que mirar dentro del cuerpo, ya que los errores llegan con estado 200; y la autorización debe vivir en la capa de servicio, porque a un @ResolveField se puede llegar por caminos que tu guard de entrada no vigila. En este stack concreto, Angular con HttpClient y DTOs compartidos ya te da tipado extremo a extremo, así que adoptar GraphQL solo para «pedir menos campos» rara vez compensa. Nada impide combinarlos: GraphQL para la aplicación interna y REST para las integraciones de terceros.¿Cómo se evita el N+1 en GraphQL?
IN y deduplica las repetidas con una caché interna. Primera condición: los loaders deben crearse por petición, en la fábrica de contexto (context: ({ req }) => ({ req, loaders: crearLoaders(orm.em.fork()) })); si son singleton, su caché vive todo el proceso y acabará devolviendo datos de un usuario a otro. Segunda condición: la función de lote debe devolver un array del mismo tamaño y en el mismo orden que las claves recibidas, rellenando los huecos con null, [] o un Error; si no se respeta, DataLoader entrega valores cruzados y el fallo es silencioso y desconcertante. Complementos útiles: contar las consultas en los tests para que una regresión salte sola, activar el registro SQL en desarrollo y usar la carga anticipada del ORM (populate) cuando el campo siempre se pide, porque una sola consulta con join es aún mejor que dos.¿Cuándo debo pasar a microservicios y por qué casi nunca es la respuesta?
flush() pasa a ser una saga con compensaciones), depuración distribuida (un error deja de ser un stack trace y se convierte en una investigación en cinco servicios), coste operativo multiplicado por N (despliegue, CI, alertas, secretos, base de datos, versionado de contratos) y modos de fallo nuevos: timeouts, duplicados, cascadas. Los indicadores legítimos para extraer un servicio son medibles: un equipo bloqueado de forma sistemática por los despliegues de otro, un componente que necesita escalar diez veces más que el resto o con un perfil de recursos radicalmente distinto, o un requisito legal de aislamiento de datos. La respuesta por defecto es el monolito modular: un despliegue, módulos con fronteras explícitas, esquema de base de datos propio por módulo y comunicación solo por interfaces públicas o eventos. Así, el día que haga falta extraer uno, el trabajo será mover código y no rediseñar el sistema; y lo más probable es que ese día no llegue.¿Petición-respuesta o comunicación basada en eventos? ¿Y cómo se garantiza la entrega?
send() con @MessagePattern) cuando quien llama necesita el resultado ahora para continuar: es fácil de razonar y de depurar, pero acopla temporalmente los dos servicios, porque si el destinatario está caído la operación falla. Eventos (emit() con @EventPattern) cuando publicas un hecho consumado y no te importa quién reacciona: el emisor no espera, se pueden añadir consumidores sin tocarlo y cada uno falla y reintenta por su cuenta; el precio es consistencia eventual y trazabilidad más difícil. Sobre las garantías, hay tres niveles y conviene no engañarse: «como máximo una vez» es lo que ofrece el pub/sub de Redis, donde un mensaje emitido mientras el consumidor se reinicia se pierde sin error; «al menos una vez» es lo que ofrecen RabbitMQ con colas durable, NATS con JetStream y Kafka, y es lo que debes asumir siempre; y «exactamente una vez» no existe a nivel de transporte, se construye con confirmación manual (noAck: false y ack después de procesar), identificador de mensaje y consumidor idempotente. Añade a eso el patrón outbox para que publicar y escribir en la base de datos sean atómicos, una dead letter queue vigilada para los mensajes que fallan siempre, y presupuestos de timeout decrecientes por salto en el lado de petición-respuesta.¿Qué es una saga y por qué aparece en cuanto distribuyes?
ROLLBACK global; el commit en dos fases existe pero bloquea recursos, no escala y casi ningún broker ni ORM lo hace práctico. Ejemplo: cobrar, reservar stock y crear el envío. Si el stock falla, no se «deshace» el cobro con un rollback, se emite un reembolso, que es una operación de negocio nueva, visible en los extractos y probablemente con implicaciones contables. Hay dos formas de coordinarla: coreografiada, donde cada servicio reacciona a los eventos de los demás (mínimo acoplamiento, pero el flujo completo no está escrito en ningún sitio), y orquestada, con una máquina de estados explícita que dirige y guarda el progreso de cada saga en curso (depurable y consultable, a cambio de un componente más). Tres consecuencias que hay que aceptar antes de empezar: los estados intermedios son visibles para el usuario y hay que diseñarlos en la interfaz; toda compensación debe ser idempotente porque puede reintentarse; y algunas operaciones no se pueden compensar (un correo enviado no se puede desenviar), lo que obliga a reordenar los pasos para dejar lo irreversible al final.¿Necesito Redis para todo esto? ¿Puedo usar la misma instancia para caché, colas y sockets?
maxmemory-policy con expulsión de claves poco usadas, mientras que las colas y el pub/sub no toleran que Redis borre sus claves para hacer sitio; mezclarlas con una política agresiva significa perder trabajos. La segunda es el aislamiento de fallos y de carga: una avalancha de trabajos o un KEYS descuidado afectan a todo lo demás, incluido el limitador que protege tu login. En producción con volumen, la recomendación habitual es una instancia (o al menos una base de datos) para caché volátil y otra para colas persistentes, cada una con su configuración de persistencia. Y no olvides que BullMQ requiere maxRetriesPerRequest: null si le pasas tu propia conexión de ioredis, porque sus workers usan comandos bloqueantes.¿Cómo se prueban la caché, las colas y los WebSockets sin volverse loco?
CacheModule.register() en memoria basta en los tests, y en integración levanta Redis con contenedores efímeros. Un test que casi nadie escribe y que evita el peor fallo de la sección: comprobar que dos usuarios distintos con la misma URL no comparten entrada. Para las colas, prueba el consumidor como lo que es, una clase con un método: constrúyela con dependencias falsas e invoca process() con un objeto de trabajo simulado; después añade el test de idempotencia llamando dos veces al mismo trabajo y verificando que el efecto ocurrió una sola vez. Para los WebSockets, arranca la aplicación en un puerto efímero y conéctate con un cliente real de socket.io-client: es la única forma de validar el handshake, el rechazo por token inválido, el permiso al unirse a una sala y la entrega a la sala correcta. En GraphQL, además de los tests de resolvers, cuenta las consultas SQL emitidas para que una regresión de N+1 rompa la construcción en lugar de aparecer en producción. El capítulo 13 desarrolla la estrategia completa de testing y observabilidad.11.12 Ejercicios
Los ejercicios están pensados para hacerse sobre el proyecto de tareas que se viene construyendo desde el
capítulo 9. Necesitarás un Redis local; la forma más rápida es
docker run -p 6379:6379 redis:7-alpine. Cuando un ejercicio pida «medir», mide de verdad: casi
todas las decisiones de este capítulo se justifican con números y se estropean con intuiciones.
11.1 Implementa el patrón cache-aside en un servicio que calcule un agregado costoso (por ejemplo, tareas por estado y proyecto de los últimos treinta días). Añade TTL con jitter y registra los aciertos y los fallos. Mide la latencia del primer acceso y de los siguientes, y anota qué porcentaje de aciertos obtienes en un uso normal de la aplicación.
11.2 Reproduce a propósito la fuga de datos del CacheInterceptor: cachea GET /tareas con la clave por defecto, entra con dos usuarios distintos y comprueba que el segundo ve las tareas del primero. Después arréglalo sobrescribiendo trackBy() y escribe un test de integración que falle si alguien vuelve a introducir el problema.
11.3 Monta una cola correos con BullMQ: productor, consumidor que extiende WorkerHost, attempts, retroceso exponencial, removeOnComplete y jobId estable. Haz que el envío falle de forma aleatoria una vez de cada tres y observa en los registros la secuencia de reintentos con sus retardos reales.
11.4 Expón el progreso de un trabajo largo con Server-Sent Events y consúmelo desde Angular. Comprueba qué ocurre exactamente cuando el servidor se reinicia a mitad del flujo: ¿reconecta el navegador solo? ¿Cuánto tarda? ¿Se pierde algún valor intermedio?
11.5 Implementa la invalidación por clave versionada de la sección 11.2.7 y verifica dos propiedades con tests: que ?estado=hecha&page=1 y ?page=1&estado=hecha comparten entrada de caché, y que crear una tarea invalida en un solo paso todas las combinaciones de filtros y páginas de ese usuario, sin recorrer claves de Redis.
11.6 Escribe un consumidor de cobros idempotente y demuéstralo: un test que ejecute el mismo trabajo dos veces y compruebe que solo hay un cargo, un correo y un cambio de estado. Añade un segundo test con dos ejecuciones concurrentes y explica cuál de las dos barreras (comprobación de estado o escritura condicional) es la que realmente te salva en cada caso.
11.7 Construye el gateway de WebSockets autenticado en el handshake, con sala personal por usuario y salas por proyecto autorizadas al unirse. Arranca dos instancias en puertos distintos, comprueba que las notificaciones se pierden, añade el adaptador de Redis y comprueba que dejan de perderse. Documenta el síntoma antes y después.
11.8 Expón tareas y su campo propietario por GraphQL code-first. Activa el registro de SQL, ejecuta una consulta que devuelva cincuenta tareas con su propietario y cuenta las sentencias emitidas. Introduce DataLoader y vuelve a contar. Escribe después un test que falle si el número de consultas supera un umbral, para que una regresión de N+1 rompa la construcción.
11.9 Convierte un @Cron que hoy hace todo el trabajo en un planificador que únicamente encola, protegido con un cerrojo distribuido SET NX PX. Arranca tres instancias a la vez y comprueba que el trabajo se encola una sola vez y que los identificadores de trabajo impiden duplicados incluso si el cerrojo caduca antes de tiempo.
11.10 Implementa una caché resistente a la estampida que combine single-flight en el proceso, cerrojo distribuido entre réplicas y stale-while-revalidate. Provoca la estampida con una herramienta de carga (autocannon -c 200 -d 20 …) contra una clave que caduca durante la prueba y compara consultas a la base de datos y latencia del percentil 99 con y sin las protecciones.
11.11 Implementa el patrón outbox transaccional: tabla outbox escrita en la misma transacción que el cambio de estado y un publicador aparte que la vacía. Demuestra la garantía provocando un fallo del broker justo después del commit y comprobando que el mensaje se publica igualmente cuando el broker vuelve. Añade el consumidor idempotente que deduplica por identificador de mensaje.
11.12 Modela la saga «confirmar pedido» (cobrar, reservar stock, crear envío) primero de forma coreografiada con eventos y después con un orquestador que guarde el estado de cada saga. Implementa las compensaciones, provoca el fallo del stock y verifica el reembolso. Escribe media página comparando ambas versiones en trazabilidad, acoplamiento y facilidad de depuración.
11.13 Añade un circuit breaker con opossum a una llamada entre servicios y un presupuesto de timeout decreciente por salto. Simula un servicio que responde en cinco segundos y comprueba tres cosas: que el interruptor abre, que las llamadas siguientes fallan de inmediato sin consumir conexiones y que el caso de uso devuelve una respuesta degradada útil en lugar de un error 500.
Solución comentada · 11.6 · Consumidor idempotente y test de doble ejecución
// src/pagos/pagos.processor.ts
@Processor('pagos', { concurrency: 3 })
export class PagosProcessor extends WorkerHost {
constructor(
private readonly em: EntityManager,
private readonly pasarela: PasarelaService,
@InjectQueue('correos') private readonly correos: Queue,
) { super(); }
async process(job: Job<{ pedidoId: string }>): Promise<void> {
// El worker vive fuera de una petición: sin fork(), dos trabajos
// concurrentes compartirían Identity Map y Unit of Work.
const em = this.em.fork();
// 1 · Barrera barata. Cubre el caso normal (el reintento llega segundos
// después) sin molestar a la pasarela. NO es suficiente por sí sola:
// dos workers concurrentes pueden leer 'pendiente' los dos.
const pedido = await em.findOneOrFail(Pedido, { id: job.data.pedidoId });
if (pedido.estado === 'pagado') return;
// 2 · Clave de idempotencia ESTABLE: la misma en todos los reintentos del
// mismo pedido. Nunca uses job.attemptsMade ni Date.now() aquí: si la
// clave cambia, la pasarela crea un cargo nuevo y has cobrado dos veces.
const cargo = await this.pasarela.cobrar(pedido.total, pedido.tarjeta, {
idempotencyKey: `pedido:${pedido.id}:cobro`,
});
// 3 · Transición atómica en la base de datos: la carrera la resuelve el
// motor, no nuestro código. Solo una de las ejecuciones ve 1 fila.
const afectadas = await em.nativeUpdate(
Pedido,
{ id: pedido.id, estado: 'pendiente' },
{ estado: 'pagado', cargoId: cargo.id, pagadoEn: new Date() },
);
if (afectadas === 0) return; // otro ganó: no dupliques los efectos siguientes
// 4 · El recibo es OTRO trabajo, con su propio jobId estable: si el envío
// falla, se reintenta el correo y no el cobro.
await this.correos.add('recibo', { pedidoId: pedido.id }, {
jobId: `recibo:${pedido.id}`,
attempts: 5,
backoff: { type: 'exponential', delay: 5_000 },
});
}
}
// test/pagos.processor.spec.ts
describe('PagosProcessor · idempotencia', () => {
it('ejecutar el MISMO trabajo dos veces cobra una sola vez', async () => {
const pedido = await crearPedido({ total: 4990, estado: 'pendiente' });
const job = { name: 'cobrar', data: { pedidoId: pedido.id } } as Job<DatosCobro>;
await processor.process(job);
await processor.process(job); // reentrega: es EXACTAMENTE lo que hará la cola
expect(pasarela.cobrar).toHaveBeenCalledTimes(1); // barrera de estado
expect(await contarTrabajos('correos', `recibo:${pedido.id}`)).toBe(1);
expect((await recargar(pedido)).estado).toBe('pagado');
});
it('dos ejecuciones CONCURRENTES no duplican efectos', async () => {
const pedido = await crearPedido({ total: 4990, estado: 'pendiente' });
const job = { name: 'cobrar', data: { pedidoId: pedido.id } } as Job<DatosCobro>;
// Aquí la comprobación de estado NO sirve: las dos leen 'pendiente'.
// Quien salva la situación es la escritura condicional del paso 3, y la
// clave de idempotencia, que hace que la pasarela devuelva el MISMO cargo.
await Promise.all([processor.process(job), processor.process(job)]);
expect(await contarTrabajos('correos', `recibo:${pedido.id}`)).toBe(1);
expect(cargosCreadosEn(pasarela)).toBe(1);
});
});
Qué demuestra cada test y por qué hacen falta los dos. El primero simula la reentrega secuencial, que es el caso habitual: el worker murió tras cobrar y antes de confirmar el trabajo. Ahí basta la comprobación de estado. El segundo simula dos workers que toman el trabajo casi a la vez —posible tras un lock vencido, un trabajo stalled o un reintento adelantado— y en ese escenario la comprobación de estado es inútil, porque las dos ejecuciones leen pendiente antes de que ninguna escriba. Solo la escritura condicional y la clave de idempotencia del tercero evitan el doble efecto.
Errores frecuentes al resolverlo. Usar em.flush() con la entidad cargada en lugar de nativeUpdate con condición: eso escribe siempre y pierde la carrera silenciosamente. Comprobar el estado y cobrar en pasos separados sin ninguna condición en la escritura, que es el mismo bug con más código. Y encolar el recibo sin jobId, que devuelve el duplicado por la puerta de atrás.
Solución comentada · 11.10 · Caché resistente a la estampida con refresco en segundo plano
// src/common/cache/cache-robusta.service.ts
interface Envoltorio<T> { valor: T; frescoHasta: number; }
@Injectable()
export class CacheRobustaService {
private readonly log = new Logger(CacheRobustaService.name);
/** Promesas en curso por clave: elimina la estampida DENTRO de este proceso. */
private readonly enVuelo = new Map<string, Promise<unknown>>();
constructor(
@Inject(CACHE_MANAGER) private readonly cache: Cache,
@Inject('REDIS') private readonly redis: Redis,
) {}
/**
* frescoMs → durante este tiempo se sirve sin más.
* toleradoMs → pasado el frescor, se sigue sirviendo el valor viejo
* mientras se recalcula por detrás (stale-while-revalidate).
*/
async obtener<T>(
clave: string, frescoMs: number, toleradoMs: number, calcular: () => Promise<T>,
): Promise<T> {
const guardado = await this.leer<Envoltorio<T>>(clave);
// 1 · Fresco: camino rápido, sin sorpresas.
if (guardado && Date.now() < guardado.frescoHasta) return guardado.valor;
// 2 · Rancio pero tolerable: respondemos YA y refrescamos en segundo plano.
// El usuario nunca espera y la base de datos nunca recibe la avalancha.
if (guardado) {
void this.refrescar(clave, frescoMs, toleradoMs, calcular)
.catch((e) => this.log.warn(`Refresco en segundo plano falló: ${e}`));
return guardado.valor;
}
// 3 · Caché vacía de verdad (arranque en frío): colapsamos las llamadas
// concurrentes de este proceso en una sola ejecución.
return this.singleFlight(clave, () =>
this.refrescar(clave, frescoMs, toleradoMs, calcular));
}
private singleFlight<T>(clave: string, fn: () => Promise<T>): Promise<T> {
const existente = this.enVuelo.get(clave) as Promise<T> | undefined;
if (existente) return existente;
// Se registra ANTES del primer await: si no, dos llamadas simultáneas
// pasarían de largo y volveríamos a tener dos consultas.
const promesa = fn().finally(() => this.enVuelo.delete(clave));
this.enVuelo.set(clave, promesa);
return promesa;
}
private async refrescar<T>(
clave: string, frescoMs: number, toleradoMs: number, calcular: () => Promise<T>,
): Promise<T> {
// Cerrojo distribuido: SET con NX y PX es una única operación atómica.
// El PX evita cerrojos huérfanos si el proceso muere sin liberarlo.
const conseguido = await this.redis.set(`lock:${clave}`, '1', 'PX', 10_000, 'NX');
if (!conseguido) {
// Otra réplica está recalculando: devolvemos lo que haya antes de insistir.
const actual = await this.leer<Envoltorio<T>>(clave);
if (actual) return actual.valor;
}
try {
const valor = await calcular();
// Jitter sobre el frescor: miles de claves no deben caducar en el
// mismo segundo (típico tras un despliegue que precalienta la caché).
const jitter = Math.floor(Math.random() * frescoMs * 0.1);
const envoltorio: Envoltorio<T> = { valor, frescoHasta: Date.now() + frescoMs + jitter };
// TTL real = frescor + tolerancia: la entrada sobrevive rancia a propósito.
await this.cache.set(clave, envoltorio, frescoMs + jitter + toleradoMs);
return valor;
} finally {
if (conseguido) await this.redis.del(`lock:${clave}`);
}
}
/** La caché es opcional; la aplicación no. Un Redis caído degrada, no rompe. */
private async leer<T>(clave: string): Promise<T | null> {
try {
return (await this.cache.get<T>(clave)) ?? null;
} catch (e) {
this.log.warn(`Caché no disponible en lectura (${clave}): ${e}`);
return null;
}
}
}
// Uso: 60 s de frescor y 5 min de tolerancia. En la práctica, la clave
// caliente no vuelve a provocar una consulta síncrona nunca más.
const resumen = await this.cacheRobusta.obtener(
`informe:${proyectoId}:${mes}`, 60_000, 300_000,
() => this.calcularResumen(proyectoId, mes),
);
Cómo medirlo. Instrumenta un contador de invocaciones de calcular y lanza
autocannon -c 200 -d 20 http://localhost:3000/informes/… con un frescor de dos segundos para que la
clave caduque varias veces durante la prueba. Sin protecciones verás cientos de invocaciones y un percentil 99
del orden de la consulta costosa; con las tres capas activas verás una invocación por ventana de frescor y un
percentil 99 prácticamente igual a la mediana, porque nadie espera nunca al recálculo.
Limitaciones honestas de esta solución. El get y el set no son una operación
atómica, de modo que en el caso extremo dos réplicas pueden recalcular a la vez justo cuando expira el
cerrojo; el peor efecto es una consulta de más, no un dato incorrecto. Si necesitas exclusión mutua estricta
entre réplicas, el paso siguiente es Redlock, con la advertencia de que su corrección bajo particiones de red
está discutida. Y una nota de diseño: guardar el frescor dentro del valor es lo que permite servir datos
rancios; si te limitas al TTL del almacén, la entrada desaparece y no hay nada viejo que devolver.
Solución comentada · 11.11 · Outbox transaccional y consumidor deduplicado
// src/outbox/outbox-message.entity.ts
@Entity({ tableName: 'outbox' })
@Index({ properties: ['publicadoEn', 'creadoEn'] }) // el índice del publicador
export class OutboxMessage {
@PrimaryKey() id: string = randomUUID();
@Property() tipo!: string; // 'pedido.confirmado'
@Property({ type: 'json' }) payload!: Record<string, unknown>;
@Property() creadoEn: Date = new Date();
@Property({ nullable: true }) publicadoEn?: Date; // null = pendiente
@Property({ default: 0 }) intentos!: number;
}
// src/pedidos/pedidos.service.ts
async confirmar(pedidoId: string): Promise<void> {
await this.em.transactional(async (em) => {
const pedido = await em.findOneOrFail(Pedido, { id: pedidoId },
{ lockMode: LockMode.PESSIMISTIC_WRITE });
if (pedido.estado !== 'pendiente') return; // idempotencia del caso de uso
pedido.estado = 'confirmado';
// LA CLAVE DEL PATRÓN: el mensaje se guarda en la MISMA transacción que
// el cambio de estado. O se confirman los dos, o ninguno. Desaparece la
// ventana en la que anunciamos algo que aún puede no ocurrir (publicar
// antes del commit) o cambiamos el estado sin avisar a nadie (publicar
// después y morir en medio).
em.persist(em.create(OutboxMessage, {
tipo: 'pedido.confirmado',
payload: { pedidoId: pedido.id, total: pedido.total, version: pedido.version },
}));
});
}
// src/outbox/outbox.publisher.ts
@Injectable()
export class OutboxPublisher {
private readonly log = new Logger(OutboxPublisher.name);
constructor(
private readonly em: EntityManager,
@Inject('REDIS') private readonly redis: Redis,
@Inject('BROKER') private readonly broker: ClientProxy,
) {}
@Cron(CronExpression.EVERY_5_SECONDS)
async publicar(): Promise<void> {
// Sin cerrojo, las N réplicas leerían las mismas filas y publicarían N veces.
// No es catastrófico (el consumidor deduplica), pero es ruido evitable.
const turno = await this.redis.set('lock:outbox', '1', 'PX', 10_000, 'NX');
if (!turno) return;
const em = this.em.fork();
const pendientes = await em.find(OutboxMessage, { publicadoEn: null },
{ orderBy: { creadoEn: 'ASC' }, limit: 100 });
for (const msg of pendientes) {
try {
// El id del mensaje viaja al consumidor: es su clave de deduplicación.
this.broker.emit(msg.tipo, { ...msg.payload, mensajeId: msg.id });
msg.publicadoEn = new Date();
} catch (e) {
// No marcamos como publicado: la próxima pasada lo reintenta.
msg.intentos += 1;
this.log.warn(`Outbox ${msg.id} falló (intento ${msg.intentos}): ${e}`);
if (msg.intentos > 20) this.alertas.avisar(msg); // algo va mal de verdad
}
await em.flush();
}
}
/** La tabla outbox crece: purga lo publicado hace días o se convierte en un lastre. */
@Cron(CronExpression.EVERY_DAY_AT_4AM)
async purgar(): Promise<void> {
const limite = new Date(Date.now() - 7 * 24 * 3_600_000);
await this.em.nativeDelete(OutboxMessage, { publicadoEn: { $lt: limite } });
}
}
// src/facturas/facturas.consumer.ts · el otro extremo del contrato
@EventPattern('pedido.confirmado')
async alConfirmar(@Payload() e: PedidoConfirmado & { mensajeId: string }) {
try {
// Insertar primero, procesar después: la clave primaria (tipo, mensajeId)
// convierte el duplicado en una violación de unicidad, no en un efecto doble.
await this.em.insert(MensajeProcesado, {
tipo: 'pedido.confirmado', mensajeId: e.mensajeId, procesadoEn: new Date(),
});
} catch (err) {
if (esViolacionDeUnicidad(err)) return; // ya lo hicimos: nada que hacer
throw err;
}
await this.facturas.emitir(e.pedidoId);
}
Qué garantiza exactamente y qué no. Garantiza que si el estado cambió, el mensaje se publicará «al menos una vez», y que si el mensaje existe, el estado cambió. No garantiza «exactamente una vez»: el publicador puede emitir y morir antes de marcar la fila, y entonces repetirá. Por eso el patrón exige un consumidor idempotente; la tabla de procesados del ejemplo es la forma más general de conseguirlo, porque funciona con cualquier efecto secundario. Tampoco garantiza el orden entre mensajes de entidades distintas: si el orden importa, usa una clave de partición estable por entidad y un número de versión en el payload para descartar los atrasados.
Cómo demostrarlo en el ejercicio. Detén el broker, confirma un pedido y comprueba tres cosas:
que la petición HTTP termina correctamente, que la fila del pedido está confirmada y que hay una fila
pendiente en outbox. Arranca el broker y verifica que en la siguiente pasada el mensaje sale
y la fila queda marcada. Después publica dos veces el mismo mensaje a mano y comprueba que solo se emite una
factura. Como refinamiento, sustituye el sondeo cada cinco segundos por captura de cambios del registro de
transacciones (change data capture) con Debezium: elimina la latencia del sondeo, a cambio de una pieza
de infraestructura considerablemente más compleja de operar.
11.13 Resumen del capítulo
- La caché no es una optimización gratuita: es estado duplicado. La pregunta previa no es técnica sino de producto —cuántos segundos de desactualización tolera el dato— y la regla de colocación es cachear lo más arriba que la corrección permita, no lo más arriba posible. Antes de añadir Redis, comprueba que no te falta un índice.
- Si la respuesta depende de quién pregunta, la identidad forma parte de la clave. El
CacheInterceptorautomático usa la URL, así que en un endpoint por usuario sirve la lista de Ana a Luis: no es un fallo de rendimiento, es una fuga de datos. La segunda barrera son las cabecerasCache-Control: private, no-storepara que ninguna CDN guarde copia. - Invalidar bien es combinar tres cosas: TTL corto como red de seguridad, invalidación explícita disparada por eventos de dominio y claves versionadas cuando los filtros libres hacen imposible enumerar entradas. Y frente a la estampida, cuatro capas: jitter, single-flight, cerrojo distribuido y refresco en segundo plano.
- La caché es opcional y la aplicación no. Envuelve las lecturas y escrituras en
try/catchpara degradar a la fuente original, pon siempremaxen el almacén en memoria y verifica las unidades del TTL: pasaron de segundos a milisegundos, y unttl: 60heredado cachea sesenta milisegundos. - Una cola existe para que la petición haga solo lo imprescindible. Aporta respuesta rápida, reintentos que sobreviven a un reinicio, absorción de picos y trabajos que no caben en el timeout del balanceador. Encola identificadores y no entidades, siempre después del commit, y no salgas a producción sin
attempts, retroceso exponencial yremoveOnComplete. - Toda cola entrega «al menos una vez», así que el consumidor debe ser idempotente. Cuatro técnicas por orden de robustez: operaciones naturalmente idempotentes, escritura condicional comprobando filas afectadas, tabla de procesados con clave única y clave de idempotencia en el tercero. El
jobIddel productor no es una garantía: caduca con el trabajo. - El reloj y los eventos internos tienen letra pequeña. Un
@Cronse registra en cada réplica y se ejecuta N veces, por lo que el planificador solo debe encolar, con cerrojo distribuido o trabajos repetibles. Y un evento deeventemitter2no persiste ni cruza procesos: sirve para métricas, auditoría y caché con TTL, nunca para efectos que no puedan perderse. - En tiempo real, empieza por Server-Sent Events y sube a WebSocket cuando el cliente tenga que hablar. Si usas WebSocket: autentica en el handshake, autoriza cada sala a la que alguien se une, instala el adaptador de Redis en cuanto haya más de una réplica y recuerda que el socket complementa la carga inicial por HTTP y obliga a recargar al reconectar.
- GraphQL traslada al cliente el control de la forma de la respuesta y te pasa la factura en otro sitio: adiós a la caché HTTP, N+1 en cada resolver de campo —resuelto con DataLoader creado por petición—, límites obligatorios de profundidad y complejidad, y errores que llegan con estado 200 y hay que monitorizar por dentro.
- Los microservicios resuelven un problema organizativo, no técnico. El monolito modular es la respuesta por defecto; cuando extraigas un servicio, asume la factura completa: timeouts con presupuesto por salto, circuit breaker, outbox transaccional, idempotencia y sagas con compensación. El hilo conductor del capítulo es ese: cada herramienta duplica estado, y dominarla consiste en saber qué pasa si ese estado se pierde y qué pasa si se procesa dos veces.
11.14 Recursos adicionales
- NestJS · Caching — la referencia para
@nestjs/cache-manager: registro del módulo, almacenes, interceptor automático y personalización de la clave. Consúltala siempre en la versión que tengas instalada, porque las opciones de almacén y las unidades del TTL han cambiado entre versiones mayores. - NestJS · Queues — integración con BullMQ:
@Processor,WorkerHost, eventos del worker, colas con flujos y separación de productores y consumidores. - BullMQ · documentación oficial — el detalle que la documentación de Nest no cubre: estados de un trabajo, trabajos atascados, rate limiting, job schedulers, deduplicación, flujos padre-hijo y guía de buenas prácticas de producción.
- NestJS · Task Scheduling —
@Cron,@Interval,@Timeouty elSchedulerRegistrypara programación dinámica, con la sintaxis exacta de las expresiones y las zonas horarias. - Redis · documentación — comandos, expiración de claves, políticas de
maxmemory, contadores atómicos y patrones de cerrojo distribuido. Es la pieza que sostiene la caché, las colas, el limitador de peticiones y la difusión entre gateways de este capítulo. - NestJS · WebSockets Gateways — decoradores, ciclo de vida del gateway, uso de pipes, guards y filtros en mensajes, y adaptadores para sustituir Socket.IO o escalar horizontalmente.
- Socket.IO · Redis adapter — cómo funciona la propagación por pub/sub entre instancias, qué operaciones pasan a ser asíncronas y la comparación con el resto de adaptadores disponibles.
- MDN · Using server-sent events — el formato del flujo,
EventSource, reconexión automática yLast-Event-ID: todo lo que necesitas para decidir si te basta SSE en lugar de WebSocket. - NestJS · GraphQL — enfoque code-first, resolvers,
@ResolveField, suscripciones, complejidad de consultas y la integración con el driver de Apollo. - Apollo Server · documentación — el servidor que hay debajo del driver por defecto: plugins, formato de errores, consultas persistidas, caché de respuestas y métricas por operación.
- DataLoader — la librería que resuelve el N+1, con el contrato exacto de la función de lote (mismo tamaño y orden que las claves) y las advertencias sobre el ámbito de su caché.
- NestJS · Microservices —
@MessagePattern,@EventPattern,ClientProxy, aplicaciones híbridas y la documentación específica de cada transporte (Redis, NATS, MQTT, RabbitMQ, Kafka y gRPC). - Chris Richardson · patrón Saga — la descripción canónica de la alternativa a las transacciones distribuidas, con las variantes coreografiada y orquestada y sus contrapartidas. En el mismo catálogo está el outbox transaccional, imprescindible para publicar sin perder ni inventar mensajes.
- Opossum — implementación de circuit breaker para Node con métricas y respuestas de reserva, la que se usa habitualmente para envolver clientes de servicios remotos.
- RFC 9111 · HTTP Caching — la especificación vigente de
Cache-Control,ETagy validación condicional: la capa de caché que suele estar sin explotar antes de montar Redis.