Arquitectura de microservicios, mensajería y observabilidad
Este es el módulo donde se separa al programador del ingeniero de software. No trata de «cómo partir una aplicación en trozos», sino de qué problemas organizativos y de escala justifican pagar el precio de un sistema distribuido, y de las técnicas concretas —DDD, arquitectura hexagonal, resiliencia, saga, outbox, Kafka, trazas distribuidas— con las que ese precio se vuelve manejable. Verás mucho código de Spring Boot 3 con Java 21, muchos diagramas y, sobre todo, criterios de decisión, incluido el más importante: cuándo no hacer microservicios.
1 · Arquitecturas de software: del monolito a lo serverless
1.1 El espectro completo (no es una escalera de «peor a mejor»)
El error mental número uno es pensar que existe una progresión evolutiva monolito → microservicios → serverless en la que cada paso es una mejora. No lo es: es un eje de compromisos en el que se cambia simplicidad operativa por independencia de despliegue y escalado. Cada punto del eje es la respuesta correcta para algún contexto.
MENOS coste operativo MÁS independencia
MÁS simple de razonar MÁS aislamiento de fallos
◄──────────────────────────────────────────────────────────────────────────►
┌────────────┐ ┌────────────────┐ ┌────────┐ ┌────────────────┐ ┌────────────┐
│ MONOLITO │ │ MONOLITO │ │ SOA │ │ MICROSERVICIOS │ │ SERVERLESS │
│ │ │ MODULAR │ │ │ │ │ │ (FaaS) │
├────────────┤ ├────────────────┤ ├────────┤ ├────────────────┤ ├────────────┤
│ 1 proceso │ │ 1 proceso │ │ N srv. │ │ N procesos │ │ N funciones│
│ 1 BD │ │ 1 BD, esquemas │ │ 1 ESB │ │ 1 BD/servicio │ │ efímeras │
│ 1 desplieg.│ │ separados │ │ 1 BD │ │ N despliegues │ │ escala a 0 │
│ paquetes │ │ módulos con │ │ común │ │ red entre todos│ │ eventos │
│ acoplados │ │ API interna │ │ (!) │ │ │ │ │
└────────────┘ └────────────────┘ └────────┘ └────────────────┘ └────────────┘
│ │ │ │ │
Startup, MVP, El 80 % de las Herencia de Muchos equipos, Cargas a ráfagas,
equipos < 20 empresas los 2000 escalado dispar, glue code, ETL,
deberían estar > 50 personas webhooks
aquí
| Estilo | Qué es exactamente | Unidad de despliegue | Datos |
|---|---|---|---|
| Monolito | Toda la funcionalidad en un artefacto (un .jar), con llamadas en memoria y una sola base de datos. Sin fronteras internas obligatorias. |
1 artefacto | 1 esquema compartido |
| Monolito modular | Un artefacto, pero con módulos con fronteras explícitas y verificadas: cada módulo expone una API interna, tiene sus propias tablas y no accede a las de otros. Es un microservicio sin la red. | 1 artefacto | 1 BD, N esquemas lógicos |
| SOA clásica | Servicios grandes coordinados por un Enterprise Service Bus con lógica de negocio y transformaciones dentro del bus. Contratos SOAP/WSDL y gobierno centralizado. | N servicios grandes | Frecuentemente compartida |
| Microservicios | Servicios pequeños alineados a capacidades de negocio, con despliegue independiente, base de datos propia, propiedad de un equipo y comunicación por red. Tuberías tontas, extremos inteligentes. | N artefactos | 1 BD por servicio |
| Serverless / FaaS | Funciones sin servidor gestionado, con escalado a cero y facturación por invocación. El proveedor asume runtime, escalado y disponibilidad. | N funciones | Servicios gestionados |
1.2 Tabla honesta de ventajas y desventajas
| Criterio | Monolito | Monolito modular | Microservicios | Serverless |
|---|---|---|---|---|
| Tiempo hasta el primer despliegue | Horas | Días | Semanas (plataforma incluida) | Horas si el proveedor ya está |
| Coste cognitivo de una feature | Bajo | Bajo | Alto: varios repos, contratos, versiones | Medio |
| Refactor entre fronteras | Trivial (lo hace el IDE) | Fácil (compilador + ArchUnit) | Caro: cambio de contrato y despliegue coordinado | Caro |
| Transacciones ACID | Sí, gratis | Sí, gratis | No: saga y consistencia eventual | No |
| Depuración de un fallo | Un stack trace | Igual | Trazas distribuidas obligatorias | Muy difícil sin observabilidad |
| Escalado | Todo o nada | Todo o nada | Por servicio, muy fino | Automático y a cero |
| Aislamiento de fallos | Nulo: una fuga tumba todo | Bajo | Alto si hay resiliencia; si no, peor que el monolito | Alto |
| Despliegue independiente por equipo | No | No (mismo binario) | Sí — su razón de ser | Sí |
| Heterogeneidad tecnológica | No | No | Sí (y suele ser una trampa) | Sí |
| Latencia interna | Nanosegundos (llamada de método) | Nanosegundos | Milisegundos por salto, y se acumulan | Milisegundos más arranques en frío |
| Coste de infraestructura | Bajo | Bajo | Alto: réplicas, malla, observabilidad, brokers | Barato a ráfagas, caro sostenido |
| Personas mínimas para operarlo bien | 2–3 | 3–8 | ~20+ con alguien dedicado a plataforma | 3–8 y dependencia del proveedor |
| Onboarding de un junior | 1 semana | 1 semana | 1–2 meses | Variable |
1.3 La ley de Conway: la arquitectura que ya tienes
Melvin Conway, 1967: «Las organizaciones que diseñan sistemas están abocadas a producir diseños que son copias de sus estructuras de comunicación». No es una metáfora ni un chiste: es una observación empírica que se cumple sin excepción.
ORGANIZACIÓN ARQUITECTURA QUE PRODUCE
───────────────────────────────────────── ────────────────────────────────────────
Equipo de frontend │ Equipo de backend Frontend ──HTTP──► Backend ──SQL──► BD
Equipo de BD │ (tres capas, tres equipos, tres colas
│ de tickets entre ellos)
Equipo "Pedidos" │ Equipo "Pagos" Servicio Pedidos ──evento──► Servicio Pagos
Equipo "Envíos" │ Servicio Envíos
(fronteras técnicas = fronteras de equipo)
3 equipos en 3 zonas horarias, con poca 3 servicios con contratos rígidos,
comunicación entre ellos versionados y coordinación mínima
1 equipo de 5 personas obligado a hacer Un MONOLITO DISTRIBUIDO: N despliegues
"microservicios" porque toca que hay que actualizar a la vez
La consecuencia práctica se conoce como maniobra inversa de Conway (popularizada por Team Topologies): si quieres una arquitectura determinada, reorganiza primero los equipos para que su estructura de comunicación se parezca a la arquitectura deseada. Dicho en negativo: no puedes tener microservicios de verdad si todas las decisiones pasan por un único comité de arquitectura y todos los equipos comparten el mismo backlog.
1.4 El coste real de los microservicios
Nadie te vende esta lista en la charla de la conferencia. Cada punto es dinero, tiempo o noches sin dormir.
| Coste | En qué se traduce | Mitigación (que también cuesta) |
|---|---|---|
| Operación | N pipelines, N imágenes, N configuraciones, N secretos, N conjuntos de alertas, N runbooks. La plataforma se convierte en un producto interno con su propio equipo. | Plantillas de proyecto, golden path, GitOps, un equipo de plataforma real. |
| Latencia | Una llamada de método pasa de ~20 ns a 1–5 ms por salto (serialización, red, deserialización). Cinco saltos en serie son 5–25 ms de suelo antes de hacer nada útil. | Menos saltos, llamadas en paralelo, caché, asincronía, colocalidad. |
| Consistencia | Se pierde la transacción ACID que cruzaba módulos. Aparecen estados intermedios visibles para el usuario y hay que diseñar compensaciones. | Saga, outbox, idempotencia y conversaciones incómodas con negocio. |
| Depuración | El stack trace deja de contar la historia completa. Sin traza distribuida, un fallo en producción puede costar horas de búsqueda en cinco sistemas de logs. | OpenTelemetry desde el día 1, correlación obligatoria, logs estructurados. |
| Pruebas | Los tests de integración de verdad necesitan varios servicios arriba. La combinatoria de versiones explota. | Testcontainers, dobles de prueba y contract testing (ver módulo 07). |
| Equipos y carga cognitiva | Cada persona debe conocer la topología, no solo su código. Se dispara si un equipo posee más de 3–4 servicios. | Propiedad clara, documentación viva (C4, ADRs), límite de servicios por equipo. |
| Disponibilidad compuesta | Si una petición depende de 10 servicios al 99,9 %, la disponibilidad resultante es 0,99910 ≈ 99,0 %: pasas de 43 min/mes de caída a más de 7 h/mes. | Resiliencia (sección 5), degradación elegante, asincronía. |
| Coste económico | Mínimo dos réplicas por servicio, malla de servicios, brokers y almacenamiento de trazas y logs (el mayor sobrecoste oculto de la observabilidad). | Muestreo, retención agresiva, métricas de baja cardinalidad. |
EFECTO "COLA LARGA" (tail at scale) — por qué el p99 empeora al dividir
Un servicio con p99 = 100 ms parece rápido.
Si una petición del usuario abanica a 10 servicios en paralelo y espera a todos:
P(al menos uno cae en su p99) = 1 - 0,99^10 = 9,6 %
→ casi 1 de cada 10 peticiones del usuario sufre el peor caso de ALGÚN servicio.
→ El p99 del sistema NO es el p99 de sus partes: es bastante peor.
Y en serie, las latencias se SUMAN:
Gateway ──2ms──► Pedidos ──8ms──► Clientes ──6ms──► Precios ──9ms──► Inventario
Suelo de red y saltos: ~25 ms + el trabajo real de cada servicio.
1.5 «Monolito modular primero»: la recomendación por defecto
Por defecto —y «por defecto» significa salvo que puedas justificar lo contrario por escrito— empieza con un monolito modular. No es una concesión ni un paso intermedio de segunda: es una arquitectura legítima que, bien hecha, sostiene productos de millones de usuarios.
MONOLITO MODULAR — un despliegue, fronteras reales
┌──────────────────────────────────────────────────────────────────────────┐
│ aplicacion.jar │
│ │
│ ┌───────────────┐ API interna ┌───────────────┐ API interna │
│ │ MÓDULO │ ──────────────► │ MÓDULO │ ─────────────┐ │
│ │ pedidos │ (interfaz + │ facturacion │ │ │
│ │ │ evento) │ │ ▼ │
│ │ · api/ │ │ · api/ │ ┌───────────┐ │
│ │ · dominio/ │ ◄────evento──── │ · dominio/ │ │ MÓDULO │ │
│ │ · infra/ │ │ · infra/ │ │ envios │ │
│ └───────┬───────┘ └───────┬───────┘ └─────┬─────┘ │
│ │ │ │ │
│ esquema esquema esquema │
│ pedidos facturacion envios │
│ ────────┴─────────────────────────────────┴────────────────────┴────── │
│ UNA sola base de datos PostgreSQL │
│ (esquemas separados; NINGÚN JOIN entre esquemas de módulos) │
└──────────────────────────────────────────────────────────────────────────┘
Reglas que lo hacen "modular" de verdad (si no, es un monolito normal):
1. Cada módulo expone SOLO su paquete `api`; el resto es package-private.
2. Ningún módulo consulta las tablas de otro. Ni un SELECT. Ni "solo para leer".
3. La comunicación entre módulos es por interfaz o por evento de aplicación.
4. Las reglas 1–3 se verifican en CI con ArchUnit o con Spring Modulith.
5. Cada módulo tiene sus propios tests, que arrancan sin los demás módulos.
// Verificación automática de las fronteras con Spring Modulith (Spring Boot 3.x)
// dependencia: org.springframework.modulith:spring-modulith-starter-core
@ApplicationModuleTest // arranca SOLO el módulo bajo prueba
class PedidosModuloTest {
@Test
void las_fronteras_entre_modulos_se_respetan() {
ApplicationModules.of(Aplicacion.class).verify(); // falla si un módulo espía a otro
}
@Test
void documenta_la_arquitectura() {
new Documenter(ApplicationModules.of(Aplicacion.class))
.writeDocumentation(); // genera diagramas C4 y tablas en target/
}
}
1.6 Cuándo dividir de verdad: los cuatro criterios
Extrae un servicio de tu monolito modular cuando puedas responder «sí, claramente» a al menos uno de estos cuatro criterios. Si la respuesta es «es que queda más limpio», la respuesta es no.
| Criterio | Señal objetiva y medible | Ejemplo real |
|---|---|---|
| 1. Equipos independientes | Dos equipos distintos bloquean el despliegue del otro; hay conflictos de merge constantes y la coordinación de entregas aparece en las retrospectivas. | El equipo de Pagos no puede sacar una pasarela nueva porque el monolito se despliega los martes con todo el catálogo. |
| 2. Escalado dispar | Un módulo consume un orden de magnitud más CPU, memoria o E/S que el resto, y pagas réplicas completas para escalar solo esa parte. | El motor de recomendaciones necesita 8 vCPU y 16 GB; el resto funciona con 1 vCPU y 1 GB. |
| 3. Ciclos de vida distintos | Un módulo cambia varias veces al día y otro dos veces al año, o tienen requisitos regulatorios y tecnológicos incompatibles. | El motor de reglas de precios se toca a diario; el módulo de contabilidad certificado se congela y se audita. |
| 4. Aislamiento de fallos | Un componente inestable o dependiente de terceros arrastra a todo el proceso: fugas, saturación de hilos, GC. | La generación de PDF con una librería nativa provoca OOM y tumba también el checkout. |
1.7 El antipatrón: el monolito distribuido
Es el peor de todos los mundos: pagas el coste íntegro de un sistema distribuido y no obtienes ninguna de sus ventajas. Se reconoce por estos síntomas:
- Despliegues coordinados: para sacar una feature hay que desplegar tres servicios en un orden concreto. Si esto pasa, no tienes servicios independientes: tienes un monolito con latencia de red.
- Base de datos compartida: dos o más servicios leen o escriben las mismas tablas. El esquema es un contrato oculto que nadie versiona.
- Cadenas síncronas largas: A llama a B, que llama a C, que llama a D, y el usuario espera. La disponibilidad se multiplica y el fallo se propaga hacia arriba.
- Librería «común» con el dominio dentro: un
commons-modelocon las entidades compartidas, que obliga a subir versión y redesplegar todo al cambiar un campo. - Un cambio funcional toca N repositorios: la métrica más honesta. Si la media de repositorios por historia de usuario supera 1,5, las fronteras están mal.
MONOLITO DISTRIBUIDO MICROSERVICIOS DE VERDAD
────────────────────────────────────── ──────────────────────────────────────
Feature "descuento por fidelidad": Feature "descuento por fidelidad":
1. Cambiar DTO en commons-modelo v3.4 1. Cambiar el servicio Precios.
2. Subir versión en Pedidos, Precios, 2. Desplegarlo.
Clientes y Gateway. 3. Fin.
3. Desplegar en orden: Precios → Clientes
→ Pedidos → Gateway. Los demás servicios ni se enteran: el
4. Si falla el 3.º, rollback de los cuatro. contrato es compatible hacia atrás y
el evento tiene campos opcionales.
Tiempo: 2 semanas y una ventana nocturna.
Tiempo: 2 horas, en horario laboral.
1.8 Cómo se descompone un monolito, paso a paso
Nunca con una reescritura desde cero (el famoso big bang rewrite, que fracasa la inmensa mayoría de las veces porque durante dos años no entregas valor y el sistema viejo sigue cambiando). Se hace con el patrón strangler fig (higuera estranguladora), nombre que Martin Fowler tomó de las higueras australianas que crecen alrededor de un árbol hasta sustituirlo mientras el árbol original sigue vivo.
PATRÓN STRANGLER FIG — sustitución incremental sin apagar nada
FASE 0 — punto de partida
Cliente ──────────────────────────────────► MONOLITO ────► BD
FASE 1 — interponer una fachada (todavía no cambia nada funcionalmente)
Cliente ────► GATEWAY / proxy ─────────────► MONOLITO ────► BD
(100 % del tráfico pasa) (sin cambios)
FASE 2 — extraer la primera capacidad, con doble escritura o sincronización
┌── /api/envios/** (5 % del tráfico, canary) ──┐
Cliente ────► GATEWAY ─┤ ▼
└── resto ─────────► MONOLITO ────► BD SERVICIO ENVIOS
│ │
└── evento/CDC ──────►└──► BD envíos
FASE 3 — mover el 100 % del tráfico de esa ruta y borrar el código del monolito
┌── /api/envios/** ─────────────────► SERVICIO ENVIOS ──► BD
Cliente ────► GATEWAY ─┤
└── resto ─────────► MONOLITO ────► BD (código de envíos BORRADO)
FASE N — repetir. El monolito adelgaza. Puede quedarse para siempre como
"servicio núcleo", y eso está PERFECTAMENTE BIEN.
REGLA DE ORO: no se avanza de fase sin poder VOLVER ATRÁS con un cambio de ruta.
Orden de extracción recomendado, de menor a mayor riesgo:
- Capacidades sin estado y periféricas: envío de notificaciones, generación de PDF, exportaciones. No tienen datos críticos y su fallo es tolerable.
- Capacidades con datos propios y poco acoplamiento: catálogo, búsqueda, gestión documental.
- Capacidades con escalado dispar: recomendaciones, procesamiento de imágenes, importaciones masivas.
- El núcleo transaccional (pedidos, pagos, contabilidad): el último, si acaso. Aquí es donde la consistencia eventual duele de verdad.
La capa anticorrupción (ACL)
Al extraer un servicio nuevo, el peligro es que el modelo legado —con sus tablas de 80 columnas, sus
banderas flag_tipo_2 y sus nombres de los noventa— se filtre al modelo limpio. La
capa anticorrupción es un traductor explícito en la frontera: el servicio nuevo nunca ve
el modelo viejo.
CAPA ANTICORRUPCIÓN (Anti-Corruption Layer)
┌──────────────────────────┐ ┌───────────────┐ ┌────────────────────────┐
│ MONOLITO LEGADO │ │ ACL │ │ SERVICIO ENVIOS │
│ │ │ │ │ (modelo limpio) │
│ TB_PED_CAB │ ─────► │ Traductor │ ─────► │ Envio │
│ COD_PED VARCHAR(12) │ SOAP │ · mapea │ DTO │ id: EnvioId │
│ IND_EST CHAR(1) │ o CDC │ · valida │ limpio │ estado: EstadoEnvio │
│ FEC_ALT NUMBER(8) │ │ · normaliza │ │ creadoEn: Instant │
│ FLG_URG CHAR(1) │ │ · rechaza lo │ │ urgente: boolean │
│ │ │ inválido │ │ │
└──────────────────────────┘ └───────────────┘ └────────────────────────┘
Vocabulario legado Punto ÚNICO de contagio Lenguaje ubicuo nuevo
Sin ACL, el "20080131" numérico y el CHAR(1) acaban en tu dominio nuevo, y en dos años
tu servicio "limpio" es indistinguible del legado.
// La ACL es un adaptador de salida: vive en infraestructura, NUNCA en el dominio.
package com.ejemplo.envios.infraestructura.salida.legado;
@Component
class TraductorPedidoLegado {
private static final DateTimeFormatter FORMATO_LEGADO = DateTimeFormatter.ofPattern("yyyyMMdd");
/** Traduce el modelo del monolito al lenguaje ubicuo de Envíos. Aquí muere el legado. */
SolicitudEnvio traducir(PedidoLegadoDto legado) {
Objects.requireNonNull(legado, "pedido legado");
EstadoEnvio estado = switch (legado.indEst()) {
case "P" -> EstadoEnvio.PENDIENTE;
case "E" -> EstadoEnvio.EN_TRANSITO;
case "F" -> EstadoEnvio.ENTREGADO;
case "A" -> EstadoEnvio.CANCELADO;
default -> throw new TraduccionImposibleException(
"Estado legado desconocido: " + legado.indEst()); // fallar rápido, no adivinar
};
LocalDate alta = LocalDate.parse(String.valueOf(legado.fecAlt()), FORMATO_LEGADO);
return new SolicitudEnvio(
new PedidoId(legado.codPed().strip()),
estado,
alta.atStartOfDay(ZoneOffset.UTC).toInstant(),
"S".equals(legado.flgUrg()));
}
}
Checklist — decisión arquitectónica
2 · Diseño guiado por el dominio (DDD)
DDD no es «poner los objetos en un paquete domain». Es una disciplina para
alinear el software con el negocio, y su aportación decisiva a los microservicios es la
respuesta a la pregunta más difícil: ¿por dónde corto?. La respuesta de DDD es: por los
contextos delimitados, no por capas técnicas ni por entidades de la base de datos.
2.1 Lenguaje ubicuo: el requisito previo a todo lo demás
El lenguaje ubicuo es un vocabulario compartido, riguroso y sin traducción
entre negocio, producto y código. Si el analista dice «póliza», la clase se llama Poliza, la
tabla poliza, el evento PolizaEmitida y el endpoint /polizas. Cada
traducción mental que un desarrollador tiene que hacer es una oportunidad de bug.
| Síntoma de lenguaje NO ubicuo | Coste real | Corrección |
|---|---|---|
Negocio dice «reserva», el código dice BookingEntity y la BD TB_RSV. | Cada conversación necesita un traductor humano; los malentendidos llegan a producción. | Un solo término, en un solo idioma, en todas las capas. |
Clases llamadas Manager, Helper, Processor, Data, Info. | Nombres sin significado de negocio: nadie sabe qué hacen sin leerlas. | Nombrar por la intención del dominio: CalculadoraDePrima, PoliticaDeCancelacion. |
| La misma palabra significa cosas distintas según el equipo. | Modelo imposible: se acaba con una clase gigante que intenta contentar a todos. | Es la señal de que hay dos contextos delimitados. Sección 2.2. |
| Un glosario en Confluence que nadie actualiza. | Documentación muerta. | El glosario es el código: los tipos y sus nombres son la única versión viva. |
2.2 Subdominios y contextos delimitados
Un subdominio es una parte del problema de negocio. Un contexto delimitado (bounded context) es una frontera de la solución dentro de la cual un modelo y un lenguaje son consistentes y sin ambigüedad. La regla más útil de todo DDD es esta:
LA PALABRA "CLIENTE" EN UNA ASEGURADORA — un concepto, cuatro modelos
┌───────────────────────┐ ┌───────────────────────┐ ┌───────────────────────┐
│ CONTEXTO: VENTAS │ │ CONTEXTO: PÓLIZAS │ │ CONTEXTO: SINIESTROS │
├───────────────────────┤ ├───────────────────────┤ ├───────────────────────┤
│ Cliente = │ │ Cliente = │ │ Cliente = │
│ · lead / prospecto │ │ · tomador │ │ · perjudicado │
│ · canal de captación │ │ · asegurado │ │ · parte contraria │
│ · probabilidad cierre│ │ · beneficiario │ │ · perito asignado │
│ · NO tiene pólizas │ │ · riesgo suscrito │ │ · histórico de partes│
└───────────────────────┘ └───────────────────────┘ └───────────────────────┘
▲ ▲ ▲
└──────────── mismo ID de persona ────────────────────┘
(la identidad se comparte; el MODELO no)
┌───────────────────────┐
│ CONTEXTO: FACTURACIÓN │ Intentar UNA clase Cliente con todos los campos
├───────────────────────┤ produce una entidad de 60 atributos donde el 70 %
│ Cliente = │ es null en cada uso, con validaciones condicionales
│ · pagador │ imposibles y sin ninguna invariante real.
│ · forma de pago │ Eso es el "Big Ball of Mud".
│ · morosidad │
└───────────────────────┘
Los subdominios no valen todos lo mismo, y eso determina dónde inviertes a tu mejor gente:
| Tipo de subdominio | Qué es | Estrategia | Ejemplo en una tienda online |
|---|---|---|---|
| Núcleo (core) | La razón por la que la empresa gana dinero y se diferencia. | Desarrollo propio, DDD táctico completo, el mejor equipo y el mejor testing. | Motor de precios dinámicos y recomendaciones. |
| De soporte | Necesario pero no diferencial. | Desarrollo propio sencillo, sin sobreingeniería. | Gestión de devoluciones, catálogo. |
| Genérico | Resuelto igual en toda la industria. | Comprar, no construir. | Facturación fiscal, envío de emails, autenticación, pasarela de pago. |
2.3 Mapa de contextos: las relaciones importan más que las cajas
El mapa de contextos documenta cómo se relacionan los contextos y, sobre todo, quién manda en cada relación. Es el documento de arquitectura más valioso que puedes tener, y cabe en un folio.
MAPA DE CONTEXTOS — tienda online (U = aguas arriba, D = aguas abajo)
┌─────────────────────┐
│ CATÁLOGO │ (U)
│ Open Host Service │
└──────────┬──────────┘
│ API pública versionada + eventos
┌──────────────────┼──────────────────┐
▼ ▼ ▼
┌───────────────┐ ┌───────────────┐ ┌───────────────┐
(D) │ PEDIDOS │ │ BÚSQUEDA │ │ RECOMENDA- │ (D)
│ │◄─┤ (Conformist) │ │ CIONES │
└───────┬───────┘ └───────────────┘ └───────────────┘
│ Customer/Supplier
│ (Pedidos negocia el contrato con Pagos)
▼
┌───────────────┐ ACL ┌──────────────────────────┐
(U) │ PAGOS │ ───────────────► │ PASARELA EXTERNA (SaaS) │
│ │ Anticorruption │ no podemos influir │
└───────┬───────┘ └──────────────────────────┘
│ evento PagoConfirmado (Published Language)
▼
┌───────────────┐ Shared Kernel (¡úsalo poco!)
│ FACTURACIÓN │ ◄─────────── tipos Dinero, Iva, PaisFiscal
└───────────────┘
| Patrón de relación | Significado | Cuándo usarlo |
|---|---|---|
| Partnership | Dos equipos se coordinan y triunfan o fracasan juntos. | Contextos que evolucionan a la vez. Poco escalable: puede ser señal de que deberían ser uno solo. |
| Shared Kernel | Un trozo de modelo compartido en código. | Solo para tipos verdaderamente universales e inmutables (Dinero, Cif). Cada línea compartida es acoplamiento permanente. |
| Customer / Supplier | El de aguas abajo es cliente y puede negociar prioridades con el de arriba. | La relación sana por defecto dentro de una misma empresa. |
| Conformist | El de abajo acepta el modelo del de arriba tal cual, sin traducir. | Cuando el modelo ajeno es bueno y traducirlo no aporta. Barato, pero te ata. |
| Anticorruption Layer | Traducción defensiva en la frontera. | Sistemas legados, SaaS externos y todo lo que no controlas. |
| Open Host Service | El de arriba publica una API pensada para muchos consumidores. | Contextos con muchos clientes: catálogo, identidad. |
| Published Language | Un formato de intercambio bien definido y versionado (Avro, protobuf, JSON Schema). | Siempre que haya eventos: es lo que evita que cada consumidor invente su interpretación. |
| Separate Ways | No integrar. Duplicar a propósito. | Cuando el coste de integrar supera al de duplicar. Es una decisión válida y a menudo la mejor. |
2.4 Agregados y límites transaccionales
Un agregado es un grupo de objetos que se tratan como una unidad de consistencia. Tiene una raíz (la única entidad accesible desde fuera) y protege unas invariantes. La regla que hay que grabar a fuego:
Las cuatro reglas de diseño de agregados (Vaughn Vernon):
- Protege invariantes verdaderas dentro de la frontera. Si una regla tolera unos segundos de retraso, no es una invariante: es una política, y va fuera.
- Diseña agregados pequeños. Un agregado grande (un
Clientecon sus 10.000 pedidos) provoca cargas enormes, bloqueos largos y conflictos de concurrencia constantes. - Referencia a otros agregados solo por identidad (
ClienteId, noCliente). Esto rompe el grafo de objetos y hace posible dividir después. - Actualiza otros agregados con consistencia eventual, publicando un evento de dominio. Nunca modifiques dos agregados en la misma transacción.
DISEÑO DE AGREGADOS — la frontera es la transacción
❌ AGREGADO GIGANTE (todo consistente al instante = todo bloqueado a la vez)
┌──────────────────────────────────────────────────────────┐
│ Cliente (raíz) │
│ ├── Direcciones[] │
│ ├── Pedidos[] ← 10.000 filas cargadas para cambiar │
│ │ └── Lineas[] el teléfono del cliente │
│ └── Facturas[] │
└──────────────────────────────────────────────────────────┘
Consecuencias: LazyInitializationException, bloqueos optimistas que fallan
sin parar, memoria desperdiciada y ninguna posibilidad de dividir.
✅ AGREGADOS PEQUEÑOS, RELACIONADOS POR ID
┌───────────────────┐ ┌───────────────────┐ ┌───────────────────┐
│ Cliente (raíz) │ │ Pedido (raíz) │ │ Factura (raíz) │
│ id: ClienteId │◄─ ─ │ clienteId │ │ pedidoId │
│ nombre │ id │ lineas[] │ │ total │
│ direcciones[] │ │ estado │ └───────────────────┘
└───────────────────┘ │ total ← INVARIANTE: suma de líneas │
└────────────────────────────────────────────┘
Transacción 1: confirmar Pedido → publica PedidoConfirmado
Transacción 2 (asíncrona): crear Factura al recibir el evento
INVARIANTE dentro de Pedido: total == suma de líneas Y estado válido
POLÍTICA fuera de Pedido: "todo pedido confirmado genera factura"
2.5 Entidades frente a objetos de valor
| Entidad | Objeto de valor (value object) | |
|---|---|---|
| Identidad | Tiene identidad propia que persiste aunque cambien todos sus atributos. | Se define solo por sus atributos: dos con los mismos valores son el mismo. |
| Mutabilidad | Su estado evoluciona en el tiempo. | Inmutable siempre. «Cambiar» es crear otro. |
| Igualdad | Por identificador. | Por todos los componentes (record lo da gratis). |
| Ejemplos | Pedido, Cliente, Poliza. | Dinero, Email, Direccion, Periodo, PedidoId. |
| Dónde van las validaciones | En los métodos que cambian el estado. | En el constructor: es imposible construir uno inválido. |
String email por Email
email y BigDecimal importe por Dinero importe elimina de golpe familias
enteras de bugs: parámetros intercambiados, validaciones olvidadas, monedas mezcladas y redondeos
inconsistentes. En Java 21 cuesta una línea: un record con constructor compacto.
2.6 Eventos de dominio
Un evento de dominio es un hecho relevante para el negocio que ya ha
ocurrido. Se nombra siempre en pasado (PedidoConfirmado, no
ConfirmarPedido, que sería un comando) y es inmutable: no se puede rechazar
el pasado.
| Comando | Evento | Consulta | |
|---|---|---|---|
| Intención | «Haz esto» | «Esto ha pasado» | «Dime esto» |
| Nombre | Imperativo: ConfirmarPedido | Pasado: PedidoConfirmado | ObtenerPedido |
| Destinatarios | Exactamente uno | Cero o muchos | Uno |
| ¿Se puede rechazar? | Sí (validación) | No, ya ocurrió | — |
| Acoplamiento | El emisor conoce al receptor | El emisor no conoce a los receptores | Directo |
2.7 Servicios de dominio, de aplicación y repositorios
| Elemento | Responsabilidad | Qué NO hace | ¿Transaccional? |
|---|---|---|---|
| Servicio de dominio | Lógica de negocio que no encaja de forma natural en una entidad porque implica varias: PoliticaDeDescuentos, CalculadoraDeIva. |
No accede a infraestructura, no conoce HTTP ni JPA, no abre transacciones. | No |
| Servicio de aplicación (caso de uso) | Orquesta: carga agregados por el repositorio, invoca al dominio, guarda, publica eventos, delimita la transacción y traduce excepciones. | No contiene reglas de negocio. Si tiene if de negocio, la regla está en el sitio equivocado. |
Sí: aquí va @Transactional |
| Repositorio | Interfaz en el dominio que da la ilusión de una colección en memoria de raíces de agregado: guardar, porId, buscarPor…. |
No expone EntityManager, Criteria ni SQL. No devuelve entidades JPA al dominio. |
Participa |
| Factoría | Crear agregados complejos garantizando invariantes desde el primer instante. | No persiste. | No |
2.8 DDD táctico en Java 21: código completo
Objeto de valor: Dinero
package com.ejemplo.pedidos.dominio.modelo;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.Currency;
import java.util.Objects;
/**
* Objeto de valor inmutable. Imposible construir uno inválido e
* imposible sumar euros con dólares por accidente.
*/
public record Dinero(BigDecimal importe, Currency moneda) implements Comparable<Dinero> {
public static final Currency EUR = Currency.getInstance("EUR");
// Constructor compacto: valida y NORMALIZA. Se ejecuta en toda construcción.
public Dinero {
Objects.requireNonNull(importe, "importe");
Objects.requireNonNull(moneda, "moneda");
// Escala fija según la moneda (EUR = 2): evita que 10.5 y 10.50 sean "distintos"
importe = importe.setScale(moneda.getDefaultFractionDigits(), RoundingMode.HALF_UP);
}
public static Dinero euros(String cantidad) { // SIEMPRE desde String, nunca desde double
return new Dinero(new BigDecimal(cantidad), EUR);
}
public static Dinero cero(Currency moneda) {
return new Dinero(BigDecimal.ZERO, moneda);
}
public Dinero sumar(Dinero otro) {
exigirMismaMoneda(otro);
return new Dinero(importe.add(otro.importe), moneda);
}
public Dinero restar(Dinero otro) {
exigirMismaMoneda(otro);
return new Dinero(importe.subtract(otro.importe), moneda);
}
public Dinero multiplicar(int cantidad) {
if (cantidad < 0) throw new IllegalArgumentException("cantidad negativa: " + cantidad);
return new Dinero(importe.multiply(BigDecimal.valueOf(cantidad)), moneda);
}
public Dinero aplicarPorcentaje(BigDecimal porcentaje) {
return new Dinero(importe.multiply(porcentaje)
.divide(BigDecimal.valueOf(100), RoundingMode.HALF_UP), moneda);
}
public boolean esMayorQue(Dinero otro) { return compareTo(otro) > 0; }
public boolean esPositivo() { return importe.signum() > 0; }
@Override public int compareTo(Dinero otro) {
exigirMismaMoneda(otro);
return importe.compareTo(otro.importe);
}
private void exigirMismaMoneda(Dinero otro) {
if (!moneda.equals(otro.moneda)) {
throw new MonedasIncompatiblesException(moneda, otro.moneda); // excepción de DOMINIO
}
}
@Override public String toString() { return importe + " " + moneda.getCurrencyCode(); }
}
Identificadores tipados y otros objetos de valor
// ❌ Todos los identificadores son String: nada impide pasar el argumento equivocado
void confirmar(String pedidoId, String clienteId, String cupon) { }
confirmar(clienteId, pedidoId, cupon); // compila perfectamente. Bug en producción.
// ✅ Identificadores tipados: el compilador se convierte en tu revisor
public record PedidoId(UUID valor) {
public PedidoId { Objects.requireNonNull(valor, "id de pedido"); }
public static PedidoId nuevo() { return new PedidoId(UUID.randomUUID()); }
public static PedidoId de(String texto) { return new PedidoId(UUID.fromString(texto)); }
@Override public String toString() { return valor.toString(); }
}
public record ClienteId(UUID valor) { /* ídem */ }
void confirmar(PedidoId pedidoId, ClienteId clienteId, CodigoCupon cupon) { }
// confirmar(clienteId, pedidoId, cupon); // ← ERROR DE COMPILACIÓN ✅
// Objeto de valor con regla de negocio dentro
public record Email(String valor) {
private static final Pattern PATRON = Pattern.compile("^[^@\\s]+@[^@\\s]+\\.[a-zA-Z]{2,}$");
public Email {
Objects.requireNonNull(valor, "email");
valor = valor.strip().toLowerCase(Locale.ROOT); // normalización
if (!PATRON.matcher(valor).matches()) throw new EmailInvalidoException(valor);
}
public String dominio() { return valor.substring(valor.indexOf('@') + 1); }
}
// Objeto de valor compuesto
public record Periodo(LocalDate desde, LocalDate hasta) {
public Periodo {
if (hasta.isBefore(desde)) throw new IllegalArgumentException("periodo invertido");
}
public boolean contiene(LocalDate d) { return !d.isBefore(desde) && !d.isAfter(hasta); }
public long dias() { return ChronoUnit.DAYS.between(desde, hasta) + 1; }
}
Agregado Pedido con invariantes y eventos
package com.ejemplo.pedidos.dominio.modelo;
/**
* RAÍZ DE AGREGADO. Ninguna anotación de framework: ni @Entity, ni @Component,
* ni Jackson. Este fichero debe compilar sin Spring en el classpath.
*/
public class Pedido {
private final PedidoId id;
private final ClienteId clienteId; // otro agregado: SOLO por identidad
private final List<LineaPedido> lineas;
private EstadoPedido estado;
private Dinero total;
private final Instant creadoEn;
private long version; // bloqueo optimista
private final List<EventoDominio> eventos = new ArrayList<>();
private static final int MAX_LINEAS = 100;
// --- Construcción controlada -------------------------------------------------
private Pedido(PedidoId id, ClienteId clienteId, Instant creadoEn) {
this.id = Objects.requireNonNull(id);
this.clienteId = Objects.requireNonNull(clienteId);
this.lineas = new ArrayList<>();
this.estado = EstadoPedido.BORRADOR;
this.total = Dinero.cero(Dinero.EUR);
this.creadoEn = creadoEn;
}
/** Factoría: el único modo de crear un pedido nuevo, ya válido. */
public static Pedido crear(ClienteId clienteId, Clock reloj) {
var pedido = new Pedido(PedidoId.nuevo(), clienteId, reloj.instant());
pedido.registrar(new PedidoCreado(pedido.id, clienteId, pedido.creadoEn));
return pedido;
}
/** Reconstrucción desde persistencia: NO emite eventos. */
public static Pedido rehidratar(PedidoId id, ClienteId clienteId, List<LineaPedido> lineas,
EstadoPedido estado, Instant creadoEn, long version) {
var pedido = new Pedido(id, clienteId, creadoEn);
pedido.lineas.addAll(lineas);
pedido.estado = estado;
pedido.version = version;
pedido.total = pedido.calcularTotal();
return pedido;
}
// --- Comportamiento: aquí vive el negocio ------------------------------------
public void anadirLinea(Sku sku, int unidades, Dinero precioUnitario) {
exigirEstado(EstadoPedido.BORRADOR, "añadir líneas");
if (unidades <= 0) throw new ReglaDeNegocioException("Las unidades deben ser positivas");
if (lineas.size() >= MAX_LINEAS)
throw new ReglaDeNegocioException("Un pedido no puede superar " + MAX_LINEAS + " líneas");
// Invariante: no hay líneas duplicadas del mismo SKU; se acumulan
lineas.stream()
.filter(l -> l.sku().equals(sku))
.findFirst()
.ifPresentOrElse(
existente -> reemplazar(existente, existente.conMasUnidades(unidades)),
() -> lineas.add(new LineaPedido(sku, unidades, precioUnitario)));
this.total = calcularTotal(); // INVARIANTE: total == suma de líneas, siempre
}
public void quitarLinea(Sku sku) {
exigirEstado(EstadoPedido.BORRADOR, "quitar líneas");
boolean quitada = lineas.removeIf(l -> l.sku().equals(sku));
if (!quitada) throw new ReglaDeNegocioException("El pedido no contiene el SKU " + sku);
this.total = calcularTotal();
}
public void confirmar(Clock reloj) {
exigirEstado(EstadoPedido.BORRADOR, "confirmar");
if (lineas.isEmpty()) throw new ReglaDeNegocioException("No se puede confirmar un pedido vacío");
if (!total.esPositivo()) throw new ReglaDeNegocioException("El total debe ser positivo");
this.estado = EstadoPedido.CONFIRMADO;
registrar(new PedidoConfirmado(id, clienteId, total,
lineas.stream().map(LineaPedido::aResumen).toList(), reloj.instant()));
}
public void cancelar(MotivoCancelacion motivo, Clock reloj) {
if (estado == EstadoPedido.ENTREGADO)
throw new ReglaDeNegocioException("Un pedido entregado no se cancela: se devuelve");
if (estado == EstadoPedido.CANCELADO) return; // idempotente: cancelar dos veces no falla
this.estado = EstadoPedido.CANCELADO;
registrar(new PedidoCancelado(id, motivo, reloj.instant()));
}
// --- Eventos de dominio -------------------------------------------------------
private void registrar(EventoDominio evento) { eventos.add(evento); }
/** El servicio de aplicación los recoge tras guardar y los publica (patrón outbox). */
public List<EventoDominio> eventosPendientes() { return List.copyOf(eventos); }
public void limpiarEventos() { eventos.clear(); }
// --- Consultas y utilidades ---------------------------------------------------
private Dinero calcularTotal() {
return lineas.stream().map(LineaPedido::subtotal)
.reduce(Dinero.cero(Dinero.EUR), Dinero::sumar);
}
private void reemplazar(LineaPedido vieja, LineaPedido nueva) {
lineas.set(lineas.indexOf(vieja), nueva);
}
private void exigirEstado(EstadoPedido esperado, String accion) {
if (estado != esperado)
throw new TransicionInvalidaException(
"No se puede %s un pedido en estado %s".formatted(accion, estado));
}
public PedidoId id() { return id; }
public ClienteId clienteId() { return clienteId; }
public EstadoPedido estado() { return estado; }
public Dinero total() { return total; }
public long version() { return version; }
public List<LineaPedido> lineas() { return List.copyOf(lineas); } // copia defensiva
}
// Entidad interna del agregado: no se accede a ella desde fuera
public record LineaPedido(Sku sku, int unidades, Dinero precioUnitario) {
public LineaPedido {
Objects.requireNonNull(sku);
if (unidades <= 0) throw new IllegalArgumentException("unidades <= 0");
}
public Dinero subtotal() { return precioUnitario.multiplicar(unidades); }
public LineaPedido conMasUnidades(int extra) {
return new LineaPedido(sku, unidades + extra, precioUnitario);
}
public ResumenLinea aResumen() { return new ResumenLinea(sku.valor(), unidades, subtotal().importe()); }
}
Evento de dominio PedidoConfirmado
package com.ejemplo.pedidos.dominio.evento;
/** Contrato común de todos los eventos del dominio. */
public sealed interface EventoDominio
permits PedidoCreado, PedidoConfirmado, PedidoCancelado {
UUID eventoId(); // identidad del EVENTO (para deduplicar en el consumidor)
Instant ocurridoEn(); // cuándo pasó, no cuándo se publicó
String tipo(); // nombre estable para el esquema publicado
}
/**
* Evento de negocio. Es un CONTRATO PÚBLICO: cambiarlo rompe consumidores.
* Incluye los datos que los consumidores necesitan (event-carried state transfer)
* para no tener que llamar de vuelta a Pedidos.
*/
public record PedidoConfirmado(
UUID eventoId,
PedidoId pedidoId,
ClienteId clienteId,
Dinero total,
List<ResumenLinea> lineas,
Instant ocurridoEn) implements EventoDominio {
public PedidoConfirmado {
Objects.requireNonNull(pedidoId);
Objects.requireNonNull(clienteId);
lineas = List.copyOf(lineas); // inmutable de verdad
}
// Constructor de conveniencia: el id del evento lo genera el dominio
public PedidoConfirmado(PedidoId pedidoId, ClienteId clienteId, Dinero total,
List<ResumenLinea> lineas, Instant ocurridoEn) {
this(UUID.randomUUID(), pedidoId, clienteId, total, lineas, ocurridoEn);
}
@Override public String tipo() { return "pedidos.PedidoConfirmado.v1"; } // versión EN el nombre
}
public record ResumenLinea(String sku, int unidades, BigDecimal subtotal) { }
PedidoConfirmado con tipos ricos, interno al servicio) no es el
evento de integración que publicas en Kafka. El de integración es un DTO plano, versionado, con
esquema y compatibilidad hacia atrás. Mezclarlos significa que refactorizar tu dominio rompe a otros
equipos. Traduce en la frontera, siempre.
3 · Arquitectura hexagonal, puertos y adaptadores
La arquitectura hexagonal (Alistair Cockburn, 2005), también llamada puertos y adaptadores, y sus primas la arquitectura cebolla y la limpia (Robert C. Martin) dicen todas lo mismo con distintos dibujos: las dependencias apuntan hacia dentro, hacia el dominio; el dominio no depende de nada. La tecnología (HTTP, JPA, Kafka) es un detalle intercambiable que se enchufa por los bordes.
3.1 El concepto y la dirección de las dependencias
ARQUITECTURA HEXAGONAL — puertos (interfaces) y adaptadores (implementaciones)
LADO CONDUCTOR (driving) LADO CONDUCIDO (driven)
quién USA la aplicación qué USA la aplicación
┌────────────────┐ ┌────────────────────┐
│ Controlador │──┐ ┌──│ Repositorio JPA │
│ REST │ │ │ │ (PostgreSQL) │
└────────────────┘ │ │ └────────────────────┘
┌────────────────┐ │ ┌──────────────────────┐ │ ┌────────────────────┐
│ Listener Kafka │──┼──►│ PUERTO DE ENTRADA │ ├──│ Publicador Kafka │
└────────────────┘ │ │ ConfirmarPedidoUC │ │ └────────────────────┘
┌────────────────┐ │ ├──────────────────────┤ │ ┌────────────────────┐
│ Comando CLI │──┤ │ │ ├──│ Cliente HTTP pagos │
└────────────────┘ │ │ APLICACIÓN │ │ └────────────────────┘
┌────────────────┐ │ │ (casos de uso, │ │ ┌────────────────────┐
│ Tarea planif. │──┘ │ transacciones) │ └──│ Adaptador de email │
└────────────────┘ │ │ └────────────────────┘
│ ┌────────────────┐ │ ▲
│ │ DOMINIO │ │ │
│ │ Pedido │ │ ┌────────────────────┐
│ │ Dinero │ │ │ PUERTOS DE SALIDA │
│ │ Políticas │ │◄──────│ (interfaces EN el │
│ │ Eventos │ │ │ dominio) │
│ │ │ │ │ RepositorioPedidos │
│ │ SIN Spring │ │ │ PublicadorEventos │
│ │ SIN JPA │ │ │ PasarelaDePagos │
│ │ SIN Jackson │ │ └────────────────────┘
│ └────────────────┘ │
└──────────────────────┘
DIRECCIÓN DE LAS DEPENDENCIAS (la única regla que importa):
infraestructura ───► aplicación ───► dominio ───► (nada)
▲
y NUNCA una flecha saliendo del dominio
Se logra con INVERSIÓN DE DEPENDENCIAS: el dominio DECLARA la interfaz que
necesita (puerto de salida) y la infraestructura la IMPLEMENTA (adaptador).
infraestructura entero y que dominio y aplicacion
sigan compilando? ¿Puedes escribir un test del caso de uso que se ejecute en 5 ms sin
Spring, sin base de datos y sin Docker? Si la respuesta a ambas es sí, lo tienes. Si no, hay una flecha
en la dirección equivocada.
3.2 Estructura de paquetes concreta en un proyecto Spring Boot 3
src/main/java/com/ejemplo/pedidos/
│
├── dominio/ ← CERO dependencias de framework
│ ├── modelo/
│ │ ├── Pedido.java (raíz de agregado)
│ │ ├── LineaPedido.java
│ │ ├── PedidoId.java ClienteId.java Sku.java
│ │ ├── Dinero.java EstadoPedido.java
│ │ └── excepcion/
│ │ ├── ReglaDeNegocioException.java
│ │ ├── TransicionInvalidaException.java
│ │ └── PedidoNoEncontradoException.java
│ ├── evento/
│ │ └── EventoDominio.java PedidoConfirmado.java PedidoCancelado.java
│ ├── servicio/
│ │ └── PoliticaDeDescuentos.java (lógica que abarca varias entidades)
│ └── puerto/
│ └── salida/ ← INTERFACES que el dominio necesita
│ ├── RepositorioPedidos.java
│ ├── PublicadorEventos.java
│ ├── PasarelaDePagos.java
│ └── ConsultaInventario.java
│
├── aplicacion/ ← casos de uso; conoce el dominio, no la infra
│ ├── puerto/entrada/
│ │ ├── ConfirmarPedidoUseCase.java (interfaz del caso de uso)
│ │ └── ConsultarPedidoUseCase.java
│ ├── ConfirmarPedidoService.java (implementación + @Transactional)
│ ├── CancelarPedidoService.java
│ └── comando/
│ └── ConfirmarPedidoComando.java
│
├── infraestructura/ ← TODO lo que huele a tecnología
│ ├── entrada/
│ │ ├── rest/
│ │ │ ├── PedidoController.java
│ │ │ ├── dto/ CrearPedidoRequest.java PedidoResponse.java
│ │ │ ├── MapeadorRest.java
│ │ │ └── ManejadorErroresGlobal.java (@RestControllerAdvice)
│ │ └── mensajeria/
│ │ └── PagoConfirmadoListener.java (@KafkaListener)
│ ├── salida/
│ │ ├── persistencia/
│ │ │ ├── PedidoEntity.java LineaEntity.java (@Entity, jakarta.persistence)
│ │ │ ├── PedidoJpaRepository.java (Spring Data)
│ │ │ ├── RepositorioPedidosAdapter.java (implementa el PUERTO)
│ │ │ └── MapeadorPersistencia.java
│ │ ├── mensajeria/
│ │ │ ├── PublicadorEventosOutbox.java (implementa el PUERTO)
│ │ │ └── OutboxEntity.java
│ │ └── pagos/
│ │ ├── PasarelaDePagosHttp.java (implementa el PUERTO)
│ │ └── PagosApi.java (HTTP interface declarativa)
│ └── configuracion/
│ ├── BeansDominioConfig.java
│ ├── ObservabilidadConfig.java
│ └── ResilienciaConfig.java
│
└── PedidosApplication.java
src/test/java/com/ejemplo/pedidos/
├── dominio/ → tests unitarios puros, milisegundos, sin Spring
├── aplicacion/ → tests de caso de uso con dobles en memoria
├── arquitectura/ → ArchUnit: las reglas anteriores VERIFICADAS
└── infraestructura/ → @DataJpaTest, @WebMvcTest, Testcontainers
Pedido (dominio) y PedidoEntity (JPA), con un mapeador entre ellos. Mucha gente
lo considera duplicación innecesaria y anota el agregado con @Entity. Funciona… hasta que
Hibernate exige un constructor sin argumentos y setters, y tu agregado deja de poder proteger
sus invariantes. La separación cuesta un mapeador; la fusión cuesta el diseño.
3.3 Un caso de uso completo, de REST a la base de datos
// ───────────────────────── DOMINIO: puertos de salida ─────────────────────────
package com.ejemplo.pedidos.dominio.puerto.salida;
public interface RepositorioPedidos {
void guardar(Pedido pedido);
Optional<Pedido> porId(PedidoId id);
List<Pedido> pendientesDe(ClienteId cliente);
}
public interface PublicadorEventos {
void publicar(List<EventoDominio> eventos);
}
public interface ConsultaInventario {
/** @return true si hay stock suficiente para TODAS las líneas. */
boolean haySuficiente(List<LineaPedido> lineas);
}
// ───────────────────── APLICACIÓN: el caso de uso (orquestación) ─────────────────────
package com.ejemplo.pedidos.aplicacion;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@Service
public class ConfirmarPedidoService implements ConfirmarPedidoUseCase {
private static final Logger log = LoggerFactory.getLogger(ConfirmarPedidoService.class);
private final RepositorioPedidos repositorio;
private final PublicadorEventos publicador;
private final ConsultaInventario inventario;
private final Clock reloj;
// Inyección por constructor: sin @Autowired, campos final, testeable sin Spring
public ConfirmarPedidoService(RepositorioPedidos repositorio,
PublicadorEventos publicador,
ConsultaInventario inventario,
Clock reloj) {
this.repositorio = repositorio;
this.publicador = publicador;
this.inventario = inventario;
this.reloj = reloj;
}
/**
* La transacción vive AQUÍ, en el caso de uso: es la unidad de trabajo del negocio.
* Dentro de ella se guarda el agregado Y se escribe la tabla outbox (sección 6.5),
* de forma que el evento y el estado se confirman o se descartan juntos.
*/
@Override
@Transactional
public ResultadoConfirmacion ejecutar(ConfirmarPedidoComando comando) {
// 1. Cargar el agregado (o fallar con una excepción de dominio)
Pedido pedido = repositorio.porId(comando.pedidoId())
.orElseThrow(() -> new PedidoNoEncontradoException(comando.pedidoId()));
// 2. Consultar lo que el dominio necesita del exterior, a través de un PUERTO
if (!inventario.haySuficiente(pedido.lineas())) {
throw new StockInsuficienteException(pedido.id());
}
// 3. Ejecutar la regla de negocio: la decide el AGREGADO, no este servicio
pedido.confirmar(reloj);
// 4. Persistir y publicar
repositorio.guardar(pedido);
publicador.publicar(pedido.eventosPendientes());
pedido.limpiarEventos();
log.info("Pedido {} confirmado por {} con total {}",
pedido.id(), pedido.clienteId(), pedido.total());
return new ResultadoConfirmacion(pedido.id(), pedido.estado(), pedido.total());
}
}
// ─────────────── INFRAESTRUCTURA: adaptador de entrada (REST) ───────────────
package com.ejemplo.pedidos.infraestructura.entrada.rest;
@RestController
@RequestMapping("/api/v1/pedidos")
class PedidoController {
private final ConfirmarPedidoUseCase confirmarPedido; // depende del PUERTO, no del Service
PedidoController(ConfirmarPedidoUseCase confirmarPedido) {
this.confirmarPedido = confirmarPedido;
}
@PostMapping("/{id}/confirmacion")
ResponseEntity<PedidoResponse> confirmar(@PathVariable UUID id) {
var resultado = confirmarPedido.ejecutar(new ConfirmarPedidoComando(new PedidoId(id)));
return ResponseEntity.ok(PedidoResponse.desde(resultado));
}
}
/**
* Traducción de excepciones de dominio a HTTP. El dominio NO conoce códigos de estado;
* esta clase es el único sitio donde se decide qué es un 404 y qué es un 409.
* Formato RFC 7807 (application/problem+json), estándar en Spring Boot 3.
*/
@RestControllerAdvice
class ManejadorErroresGlobal {
@ExceptionHandler(PedidoNoEncontradoException.class)
ProblemDetail noEncontrado(PedidoNoEncontradoException e) {
var pd = ProblemDetail.forStatusAndDetail(HttpStatus.NOT_FOUND, e.getMessage());
pd.setTitle("Pedido no encontrado");
pd.setType(URI.create("https://errores.ejemplo.com/pedido-no-encontrado"));
return pd;
}
@ExceptionHandler({TransicionInvalidaException.class, StockInsuficienteException.class})
ProblemDetail conflicto(RuntimeException e) {
return ProblemDetail.forStatusAndDetail(HttpStatus.CONFLICT, e.getMessage());
}
@ExceptionHandler(ReglaDeNegocioException.class)
ProblemDetail reglaNegocio(ReglaDeNegocioException e) {
return ProblemDetail.forStatusAndDetail(HttpStatus.UNPROCESSABLE_ENTITY, e.getMessage());
}
}
// ─────────── INFRAESTRUCTURA: adaptador de salida (JPA) ───────────
package com.ejemplo.pedidos.infraestructura.salida.persistencia;
import jakarta.persistence.*;
@Entity
@Table(name = "pedido", schema = "pedidos")
class PedidoEntity {
@Id
private UUID id;
@Column(name = "cliente_id", nullable = false)
private UUID clienteId;
@Enumerated(EnumType.STRING) // NUNCA ORDINAL
@Column(nullable = false, length = 20)
private String estado;
@Column(name = "total_importe", nullable = false, precision = 19, scale = 2)
private BigDecimal totalImporte;
@Column(name = "total_moneda", nullable = false, length = 3)
private String totalMoneda;
@OneToMany(mappedBy = "pedido", cascade = CascadeType.ALL, orphanRemoval = true)
private List<LineaEntity> lineas = new ArrayList<>();
@Column(name = "creado_en", nullable = false)
private Instant creadoEn;
@Version // bloqueo optimista
private long version;
protected PedidoEntity() { } // exigido por JPA: aquí no molesta a nadie
// getters y setters de infraestructura...
}
interface PedidoJpaRepository extends JpaRepository<PedidoEntity, UUID> {
List<PedidoEntity> findByClienteIdAndEstado(UUID clienteId, String estado);
}
/** El ADAPTADOR: implementa el puerto del dominio usando Spring Data. */
@Repository
class RepositorioPedidosAdapter implements RepositorioPedidos {
private final PedidoJpaRepository jpa;
private final MapeadorPersistencia mapeador;
RepositorioPedidosAdapter(PedidoJpaRepository jpa, MapeadorPersistencia mapeador) {
this.jpa = jpa; this.mapeador = mapeador;
}
@Override public void guardar(Pedido pedido) {
jpa.save(mapeador.aEntidad(pedido));
}
@Override public Optional<Pedido> porId(PedidoId id) {
return jpa.findById(id.valor()).map(mapeador::aDominio); // JPA no sale de este paquete
}
@Override public List<Pedido> pendientesDe(ClienteId cliente) {
return jpa.findByClienteIdAndEstado(cliente.valor(), EstadoPedido.CONFIRMADO.name())
.stream().map(mapeador::aDominio).toList();
}
}
El test que demuestra que la arquitectura funciona
// Sin Spring, sin base de datos, sin Docker: 3 ms. Este es el premio de la hexagonal.
class ConfirmarPedidoServiceTest {
private final RepositorioPedidosEnMemoria repositorio = new RepositorioPedidosEnMemoria();
private final PublicadorEventosEspia publicador = new PublicadorEventosEspia();
private final Clock reloj = Clock.fixed(Instant.parse("2026-03-01T10:00:00Z"), ZoneOffset.UTC);
@Test
void confirma_el_pedido_y_publica_el_evento_cuando_hay_stock() {
var servicio = new ConfirmarPedidoService(repositorio, publicador, lineas -> true, reloj);
var pedido = unPedidoConUnaLinea();
repositorio.guardar(pedido);
var resultado = servicio.ejecutar(new ConfirmarPedidoComando(pedido.id()));
assertThat(resultado.estado()).isEqualTo(EstadoPedido.CONFIRMADO);
assertThat(publicador.publicados()).hasSize(1)
.first().isInstanceOf(PedidoConfirmado.class);
}
@Test
void rechaza_la_confirmacion_si_no_hay_stock_y_no_publica_nada() {
var servicio = new ConfirmarPedidoService(repositorio, publicador, lineas -> false, reloj);
var pedido = unPedidoConUnaLinea();
repositorio.guardar(pedido);
assertThatThrownBy(() -> servicio.ejecutar(new ConfirmarPedidoComando(pedido.id())))
.isInstanceOf(StockInsuficienteException.class);
assertThat(publicador.publicados()).isEmpty();
assertThat(repositorio.porId(pedido.id()).orElseThrow().estado())
.isEqualTo(EstadoPedido.BORRADOR);
}
}
3.4 Reglas verificables con ArchUnit
Una arquitectura que no se verifica automáticamente se degrada en tres sprints. ArchUnit convierte las reglas de la pizarra en tests que fallan en CI. Amplía esto en el módulo 07 · Testing.
package com.ejemplo.pedidos.arquitectura;
import com.tngtech.archunit.junit.AnalyzeClasses;
import com.tngtech.archunit.junit.ArchTest;
import com.tngtech.archunit.core.importer.ImportOption;
import static com.tngtech.archunit.lang.syntax.ArchRuleDefinition.*;
import static com.tngtech.archunit.library.Architectures.layeredArchitecture;
@AnalyzeClasses(packages = "com.ejemplo.pedidos",
importOptions = ImportOption.DoNotIncludeTests.class)
class ArquitecturaHexagonalTest {
// 1. El dominio no conoce Spring, ni JPA, ni Jackson, ni nada de fuera.
@ArchTest
static final ArchRule el_dominio_es_puro =
noClasses().that().resideInAPackage("..dominio..")
.should().dependOnClassesThat().resideInAnyPackage(
"..aplicacion..", "..infraestructura..",
"org.springframework..", "jakarta.persistence..",
"jakarta.servlet..", "com.fasterxml.jackson..",
"org.apache.kafka..", "org.hibernate..")
.because("el dominio debe compilar y testearse sin ningún framework");
// 2. El dominio tampoco lleva anotaciones de framework.
@ArchTest
static final ArchRule el_dominio_no_lleva_anotaciones_de_spring =
noClasses().that().resideInAPackage("..dominio..")
.should().beAnnotatedWith("org.springframework.stereotype.Component")
.orShould().beAnnotatedWith("jakarta.persistence.Entity");
// 3. Capas: la infraestructura puede depender de todo; el dominio, de nadie.
@ArchTest
static final ArchRule capas_respetadas = layeredArchitecture().consideringAllDependencies()
.layer("Dominio").definedBy("..dominio..")
.layer("Aplicacion").definedBy("..aplicacion..")
.layer("Infraestructura").definedBy("..infraestructura..")
.whereLayer("Infraestructura").mayNotBeAccessedByAnyLayer()
.whereLayer("Aplicacion").mayOnlyBeAccessedByLayers("Infraestructura")
.whereLayer("Dominio").mayOnlyBeAccessedByLayers("Aplicacion", "Infraestructura");
// 4. Los controladores no tocan repositorios: pasan siempre por un caso de uso.
@ArchTest
static final ArchRule los_controladores_usan_casos_de_uso =
noClasses().that().resideInAPackage("..infraestructura.entrada.rest..")
.should().dependOnClassesThat().resideInAPackage("..salida.persistencia..");
// 5. Las entidades JPA no salen de su paquete (nada de @Entity en la API REST).
@ArchTest
static final ArchRule las_entidades_jpa_no_se_filtran =
classes().that().areAnnotatedWith("jakarta.persistence.Entity")
.should().resideInAPackage("..infraestructura.salida.persistencia..");
// 6. Nada de println: se loguea con SLF4J.
@ArchTest
static final ArchRule sin_system_out = noClasses().should().accessStandardStreams();
// 7. Los puertos son interfaces.
@ArchTest
static final ArchRule los_puertos_son_interfaces =
classes().that().resideInAPackage("..dominio.puerto..").should().beInterfaces();
// 8. @Transactional solo en la capa de aplicación (nunca en controladores ni en el dominio).
@ArchTest
static final ArchRule transacciones_solo_en_aplicacion =
classes().that().areAnnotatedWith("org.springframework.transaction.annotation.Transactional")
.should().resideInAPackage("..aplicacion..");
}
3.5 Ventajas reales… y cuándo es sobreingeniería
Lo que ganas de verdad
- Tests rapidísimos: la lógica de negocio se prueba en milisegundos, sin contexto de Spring. Una suite de 2.000 tests de dominio tarda menos que 20 tests con
@SpringBootTest. - Cambiar de tecnología es local: pasar de REST a gRPC, de PostgreSQL a Mongo o de RabbitMQ a Kafka toca un adaptador, no el negocio.
- El negocio queda legible: se puede leer
Pedido.javacon un analista funcional al lado. - Facilita la extracción a microservicio: los puertos ya son la frontera; solo cambia el adaptador de local a remoto.
- Retrasa decisiones: puedes empezar con un repositorio en memoria y elegir la base de datos en la semana 4.
Cuándo NO merece la pena
- CRUD puro sin reglas: si el «dominio» es copiar el request a una tabla, la hexagonal solo añade tres clases por entidad. Usa Spring Data y sé feliz.
- Prototipos y pruebas de concepto con fecha de caducidad.
- Servicios de integración finos (transformar y reenviar): apenas hay dominio que proteger.
- Aplicarla como dogma: interfaces con una sola implementación para todo es ceremonia sin valor.
- Equipo sin criterio todavía: mal aplicada produce mapeadores infinitos y un «dominio» que es un DTO con otro nombre.
Checklist — DDD y arquitectura hexagonal
4 · Comunicación entre servicios
4.1 Síncrono frente a asíncrono: la decisión con más consecuencias
Antes de elegir REST, gRPC o Kafka hay una decisión más profunda: ¿el emisor necesita la respuesta para continuar? Casi siempre la respuesta honesta es «no», y aun así se implementa síncrono por inercia.
| Criterio | Síncrono (petición/respuesta) | Asíncrono (mensajes y eventos) |
|---|---|---|
| Acoplamiento temporal | Alto: los dos deben estar vivos a la vez | Nulo: el receptor puede estar caído y recuperar después |
| Disponibilidad resultante | Se multiplica: 0,999 × 0,999 × … | Se aísla: el broker desacopla |
| Latencia percibida | La suma de toda la cadena | Respuesta inmediata («aceptado»); el trabajo va detrás |
| Consistencia | Inmediata dentro del servicio | Eventual: hay una ventana de incoherencia |
| Complejidad de depuración | Media: la traza es una cascada | Alta: hay que correlacionar, ordenar y deduplicar |
| Manejo de errores | El error vuelve al llamante al instante | Reintentos, DLQ y alertas; el emisor ya se fue |
| Contrapresión | Mala: el llamante satura al llamado | Natural: la cola absorbe los picos |
| Usa esto cuando… | El usuario espera el resultado y lo necesita para decidir: consultar stock, autenticar, validar un pago en el checkout. | El resultado es un efecto secundario del negocio: facturar, notificar, indexar, calcular puntos, sincronizar sistemas. |
EL MISMO CASO DE USO, DOS DISEÑOS
❌ CADENA SÍNCRONA (el usuario espera 470 ms y el fallo de cualquiera lo tira todo)
Usuario ──► Pedidos ──► Inventario ──► Precios ──► Pagos ──► Facturación ──► Email
250 ms 40 ms 60 ms 50 ms 180 ms 60 ms 80 ms
└─► total 470 ms
Disponibilidad = 0,999^6 = 99,40 % → 4,3 horas de caída al mes
✅ NÚCLEO SÍNCRONO MÍNIMO + RESTO POR EVENTOS (el usuario espera 90 ms)
Usuario ──► Pedidos ──► Inventario (reserva) 90 ms → 202 Accepted
│ 40 ms
└── publica PedidoConfirmado ──► [ KAFKA ]
├──► Pagos (async)
├──► Facturación (async)
├──► Email (async)
└──► Analítica (async)
Disponibilidad de la ruta crítica = 0,999^2 = 99,80 % → 86 min/mes
Si Email está caído 2 horas, el usuario NI SE ENTERA: los mensajes esperan.
4.2 REST sobre HTTP: contratos, versionado y compatibilidad
REST sigue siendo el estándar de facto para comunicación síncrona: universal, depurable con
curl, cacheable y comprendido por todo el mundo. Su punto débil es la disciplina: nada te
obliga a mantener un contrato estable.
| Estrategia de versionado | Ejemplo | Pros | Contras |
|---|---|---|---|
| En la ruta (la más usada) | /api/v1/pedidos | Explícito, cacheable, trivial de enrutar en el gateway | Duplica rutas; los puristas REST protestan |
| Cabecera de tipo de medio | Accept: application/vnd.ejemplo.pedido.v2+json | URLs estables, versionado por recurso | Invisible en el navegador; más difícil de probar y cachear |
| Parámetro de consulta | /pedidos?version=2 | Simple | Se olvida y ensucia la caché |
| Sin versión, solo evolución compatible | /api/pedidos | Lo ideal: nunca rompes | Exige una disciplina que pocos equipos mantienen |
COMPATIBILIDAD HACIA ATRÁS DE UNA API REST
CAMBIOS SEGUROS (no rompen a nadie) CAMBIOS QUE ROMPEN (exigen nueva versión)
───────────────────────────────────── ──────────────────────────────────────────
+ Añadir un campo OPCIONAL a la respuesta − Eliminar o renombrar un campo
+ Añadir un endpoint nuevo − Cambiar el tipo de un campo (int → string)
+ Añadir un parámetro opcional − Hacer obligatorio un campo opcional
+ Añadir un valor a un enum SI el cliente − Cambiar el significado de un valor
lo tolera (documéntalo desde el día 1) − Cambiar el código HTTP de una respuesta
+ Relajar una validación − Endurecer una validación
+ Añadir una cabecera − Cambiar la estructura de errores
− Cambiar la semántica (idempotente → no)
REGLA DE POSTEL, aplicada con cabeza: sé conservador en lo que envías y liberal en lo
que aceptas. En la práctica: IGNORA los campos desconocidos al deserializar.
// Cliente tolerante: ignorar campos desconocidos es lo que permite al productor evolucionar
@JsonIgnoreProperties(ignoreUnknown = true) // por defecto en Spring Boot: no lo desactives
public record PedidoDto(UUID id, String estado, BigDecimal total, String moneda) { }
// Configuración global equivalente
@Bean
Jackson2ObjectMapperBuilderCustomizer tolerante() {
return builder -> builder
.featuresToDisable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES)
.serializationInclusion(JsonInclude.Include.NON_NULL); // no envíes nulls inútiles
}
// Deprecación explícita y observable de un endpoint antiguo
@GetMapping(value = "/api/v1/pedidos/{id}")
@Deprecated(since = "2026-02-01", forRemoval = true)
ResponseEntity<PedidoV1Response> obtenerV1(@PathVariable UUID id) {
contadorUsoV1.increment(); // MIDE quién sigue usándolo: sin datos no puedes apagarlo
return ResponseEntity.ok()
.header("Deprecation", "true")
.header("Sunset", "Wed, 01 Jul 2026 00:00:00 GMT") // RFC 8594
.header("Link", "<https://api.ejemplo.com/api/v2/pedidos>; rel=\"successor-version\"")
.body(servicio.obtenerV1(id));
}
4.3 gRPC: cuando el rendimiento y el contrato mandan
gRPC usa HTTP/2 como transporte y Protocol Buffers como formato binario, con generación
de código a partir de un fichero .proto que es la única fuente de verdad del
contrato.
// inventario.proto — el contrato es un fichero versionado en Git
syntax = "proto3";
package inventario.v1;
option java_multiple_files = true;
option java_package = "com.ejemplo.inventario.grpc";
service Inventario {
rpc ConsultarStock (ConsultaStockRequest) returns (StockResponse); // unario
rpc ObservarStock (ConsultaStockRequest) returns (stream StockResponse); // server stream
rpc ReportarLecturas (stream Lectura) returns (ResumenResponse); // client stream
rpc Sincronizar (stream Lectura) returns (stream StockResponse); // bidireccional
}
message ConsultaStockRequest {
string sku = 1; // el NÚMERO de campo es el contrato, no el nombre
string almacen_id = 2;
}
message StockResponse {
string sku = 1;
int32 disponible = 2;
int32 reservado = 3;
google.protobuf.Timestamp actualizado_en = 4;
reserved 5, 6; // números retirados: NUNCA se reutilizan
reserved "cantidad_antigua";
}
| Aspecto | REST + JSON | gRPC + Protobuf |
|---|---|---|
| Tamaño del mensaje | Referencia (texto, con nombres de campo repetidos) | 3–10× menor (binario, campos por número) |
| CPU de serialización | Alta (parseo de texto) | Mucho menor |
| Latencia típica intra-clúster | ~2–5 ms | ~0,5–2 ms (multiplexado HTTP/2, conexión persistente) |
| Contrato | OpenAPI opcional, a menudo desactualizado | .proto obligatorio; genera cliente y servidor |
| Streaming | SSE o WebSocket, añadido aparte | Nativo en cuatro modos |
| Depuración manual | Trivial (curl, navegador) | Necesita grpcurl; el binario no se lee |
| Navegadores | Directo | Requiere gRPC-Web y un proxy |
| Balanceo de carga | Sencillo (L4/L7, conexión corta) | Delicado: HTTP/2 mantiene la conexión y el balanceo L4 concentra el tráfico; necesita proxy L7 o balanceo en cliente |
| Caché HTTP | Sí, estándar | No |
| Curva de aprendizaje | Mínima; todo el mundo lo conoce | Media; plugin de compilación y generación de código en el build |
Elige gRPC cuando la comunicación es interna entre servicios, el volumen es alto (miles de peticiones por segundo), la latencia importa, quieres un contrato fuerte con generación de código o necesitas streaming bidireccional. Quédate en REST cuando la API es pública o la consumen navegadores y terceros, el volumen es moderado o la simplicidad operativa vale más que unos milisegundos. Un patrón muy sano es REST hacia fuera, gRPC hacia dentro.
4.4 GraphQL y el patrón BFF
GraphQL da al cliente el poder de pedir exactamente los campos que necesita, en una sola petición, resolviendo dos problemas clásicos de REST: el over-fetching (te llegan 40 campos y usas 3) y el under-fetching (necesitas 5 llamadas encadenadas para pintar una pantalla).
SIN GRAPHQL — la app móvil hace 4 viajes (y en 3G se nota)
App ──► GET /pedidos/42 (40 campos, usa 4)
──► GET /clientes/7 (30 campos, usa 2)
──► GET /productos?ids=1,2,3 (3 × 25 campos, usa 2 de cada uno)
──► GET /envios?pedido=42 (18 campos, usa 1)
CON GRAPHQL / BFF — un viaje con exactamente lo necesario
App ──► POST /graphql
query { pedido(id:42) { total estado
cliente { nombre }
lineas { producto { nombre } unidades } } }
┌──────────────┐
│ GraphQL │──► Pedidos ← el BFF hace el abanico DENTRO del centro de datos,
│ BFF / móvil │──► Clientes donde la latencia es de 1 ms, no de 200 ms
│ │──► Catálogo
└──────────────┘──► Envíos
// Spring for GraphQL (Spring Boot 3): resolver por lotes con @BatchMapping
@Controller
class PedidoGraphQlController {
private final ConsultarPedidoUseCase pedidos;
private final ClientesApi clientes;
@QueryMapping
PedidoVista pedido(@Argument UUID id) {
return pedidos.porId(id);
}
/**
* @BatchMapping resuelve el campo `cliente` de TODOS los pedidos de la consulta
* en UNA sola llamada. Este es el DataLoader de Spring GraphQL.
*/
@BatchMapping(typeName = "Pedido", field = "cliente")
Map<PedidoVista, ClienteVista> cliente(List<PedidoVista> lote) {
var ids = lote.stream().map(PedidoVista::clienteId).distinct().toList();
Map<ClienteId, ClienteVista> porId = clientes.porIds(ids); // 1 llamada, no N
return lote.stream().collect(Collectors.toMap(p -> p, p -> porId.get(p.clienteId())));
}
}
spring:
graphql:
schema:
locations: classpath:graphql/**/
graphiql:
enabled: false # NUNCA en producción
# Límites obligatorios: sin ellos, una consulta anidada de 20 niveles tumba el servicio.
# Se aplican con instrumentaciones: MaxQueryDepthInstrumentation y
# MaxQueryComplexityInstrumentation, además de consultas persistidas (APQ).
Cuándo NO usar GraphQL:
- API pública sin control de los clientes: alguien escribirá una consulta anidada que tumbe la base de datos. Necesitas límites de profundidad, complejidad, coste y consultas persistidas.
- Cuando la caché HTTP es tu principal aliado: todo va por
POSTa una misma URL; pierdes la caché de CDN y de navegador. - Comunicación servicio a servicio: ahí quieres contratos rígidos, no flexibilidad. Usa REST o gRPC.
- Un solo cliente con necesidades estables: el coste del esquema, los resolvers y el tooling no se amortiza.
- Si tu problema real es una API REST mal diseñada: GraphQL no arregla un mal modelo, lo esconde.
4.5 WebSockets y Server-Sent Events
SSE (text/event-stream) | WebSocket | |
|---|---|---|
| Dirección | Solo servidor → cliente | Bidireccional |
| Protocolo | HTTP normal: funciona con proxies, CDN y HTTP/2 | Upgrade a ws://: infraestructura específica |
| Reconexión | Automática en el navegador, con Last-Event-ID | Manual: la implementas tú |
| Formato | Texto | Texto y binario |
| Casos típicos | Notificaciones, progreso de una tarea, precios, tokens de un LLM, feeds | Chat, colaboración en vivo, juegos, trading, terminales |
ulimit y cuenta con miles, no millones, por réplica. Con virtual threads
de Java 21 (ver módulo 03) el coste por conexión baja mucho.
// SSE en Spring MVC: sencillo, resistente y suficiente para el 80 % de los casos
@GetMapping(value = "/api/pedidos/{id}/eventos", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
SseEmitter seguir(@PathVariable UUID id) {
var emitter = new SseEmitter(Duration.ofMinutes(15).toMillis());
registro.suscribir(id, emitter);
emitter.onCompletion(() -> registro.baja(id, emitter));
emitter.onTimeout(() -> registro.baja(id, emitter));
emitter.onError(e -> registro.baja(id, emitter));
return emitter;
}
void notificar(UUID pedidoId, EstadoPedido estado) throws IOException {
for (SseEmitter e : registro.de(pedidoId)) {
e.send(SseEmitter.event()
.id(String.valueOf(System.currentTimeMillis())) // permite reanudar
.name("estado")
.data(Map.of("pedidoId", pedidoId, "estado", estado)));
}
}
4.6 Mensajería: colas frente a topics
COLA (queue) — reparto de trabajo, cada mensaje lo procesa UN consumidor
┌──────────────┐
┌───────────────┐ ┌───►│ Worker 1 │ m1, m4
Productor ──► m1 m2 m3 →│ COLA │────┼───►│ Worker 2 │ m2, m5
m4 m5 └───────────────┘ └───►│ Worker 3 │ m3
└──────────────┘
Objetivo: ESCALAR el procesamiento. Semántica: "que alguien haga esto".
Ejemplos: RabbitMQ (queue), SQS, JMS queue.
TOPIC (publicación/suscripción) — cada suscriptor recibe TODOS los mensajes
┌──────────────┐
┌───────────────┐ ┌───►│ Facturación │ m1 m2 m3
Productor ──► m1 m2 m3 →│ TOPIC │────┼───►│ Email │ m1 m2 m3
└───────────────┘ └───►│ Analítica │ m1 m2 m3
└──────────────┘
Objetivo: DESACOPLAR. Semántica: "ha pasado esto, por si a alguien le interesa".
Ejemplos: Kafka (topic + consumer groups), RabbitMQ (fanout exchange), SNS.
KAFKA COMBINA LAS DOS: un topic con N particiones y varios consumer groups.
· Dentro de un grupo → comportamiento de COLA (reparto por particiones).
· Entre grupos → comportamiento de TOPIC (todos reciben todo).
4.7 Comparativa: REST, gRPC, mensajería y GraphQL
| Criterio | REST/JSON | gRPC | Mensajería (Kafka/AMQP) | GraphQL |
|---|---|---|---|---|
| Modelo | Petición/respuesta | Petición/respuesta y streaming | Publicación/suscripción, fire-and-forget | Petición/respuesta con consulta declarativa |
| Acoplamiento temporal | Alto | Alto | Ninguno | Alto |
| Rendimiento | Medio | Alto | Muy alto en volumen agregado | Medio (depende de los resolvers) |
| Contrato | OpenAPI (opcional) | .proto (obligatorio) | Esquema Avro/JSON con registro | SDL (obligatorio) |
| Evolución | Disciplina manual | Reglas de protobuf | Compatibilidad del registro de esquemas | Campos @deprecated |
| Depuración | Muy fácil | Media | Difícil (asíncrono) | Media |
| Entrega garantizada | No (se reintenta) | No | Sí, con persistencia y reintentos | No |
| Idóneo para | API públicas, integración, CRUD | Tráfico interno de alto volumen, streaming | Eventos de negocio, integración desacoplada, ETL | BFF para móvil y web con muchas vistas |
| Peor uso | Chatty APIs con 20 llamadas por pantalla | API pública para terceros y navegadores | Cuando el usuario necesita el resultado ya | Comunicación servicio a servicio |
4.8 Formatos y compatibilidad de esquemas
| Formato | Tamaño | Esquema | Legible | Cuándo |
|---|---|---|---|---|
| JSON | Grande | Opcional (JSON Schema) | Sí | API públicas, depuración y volúmenes moderados. Por defecto, salvo que midas un problema. |
| Protobuf | Muy pequeño | Obligatorio, en el fichero .proto | No | gRPC, tráfico interno intenso, contratos con generación de código. |
| Avro | Muy pequeño | Obligatorio, viaja por referencia a un registro | No | Kafka a gran escala y data lakes: el registro de esquemas hace cumplir la compatibilidad. |
| MessagePack / CBOR | Pequeño | No | No | JSON binario cuando solo quieres ahorrar bytes sin cambiar el modelo. |
REGISTRO DE ESQUEMAS (Schema Registry) — quién valida que no rompas a nadie
Productor Schema Registry Consumidor
───────── ─────────────── ──────────
1. registra el esquema ──► valida compatibilidad
con la versión anterior
(BACKWARD / FORWARD / FULL)
│
si NO es compatible: RECHAZA
el despliegue del productor ✅
│
2. serializa el mensaje ◄── devuelve schemaId (int)
[magic byte][schemaId 4B][payload Avro binario]
│
3. envía a Kafka ──────────────────┼──────────────────► 4. lee el schemaId
│ 5. descarga el esquema (cacheado)
└───────────────────► 6. deserializa con el esquema de
ESCRITURA + el de LECTURA
| Modo | Qué garantiza | Cambios permitidos | Orden de despliegue |
|---|---|---|---|
| BACKWARD (el más común) | Un consumidor con el esquema nuevo puede leer datos escritos con el anterior. | Borrar campos; añadir campos opcionales con valor por defecto. | Consumidores primero, luego productores. |
| FORWARD | Un consumidor con el esquema anterior puede leer datos escritos con el nuevo. | Añadir campos; borrar campos opcionales. | Productores primero, luego consumidores. |
| FULL | Las dos cosas a la vez. | Solo añadir o quitar campos opcionales con valor por defecto. | Indiferente. Es el modo que querrás si no controlas a los consumidores. |
| *_TRANSITIVE | La comprobación se hace contra todas las versiones anteriores, no solo la última. | Más restrictivo. | Recomendado si hay consumidores lentos en actualizarse o reprocesado histórico. |
| NONE | Nada: el registro solo almacena. | Todos. | No lo uses en producción. |
// pedido-confirmado-v2.avsc — evolución compatible en modo FULL
{
"type": "record",
"name": "PedidoConfirmado",
"namespace": "com.ejemplo.pedidos.eventos",
"fields": [
{"name": "eventoId", "type": {"type":"string","logicalType":"uuid"}},
{"name": "pedidoId", "type": {"type":"string","logicalType":"uuid"}},
{"name": "clienteId", "type": {"type":"string","logicalType":"uuid"}},
{"name": "totalCentimos", "type": "long"},
{"name": "moneda", "type": "string", "default": "EUR"},
{"name": "ocurridoEn", "type": {"type":"long","logicalType":"timestamp-millis"}},
// ✅ Campo NUEVO: unión con null y default null → compatible en ambos sentidos
{"name": "canalVenta", "type": ["null","string"], "default": null}
// ❌ Lo que NO puedes hacer sin romper:
// · renombrar "moneda" (usa "aliases" si es imprescindible)
// · cambiar totalCentimos de long a string
// · añadir un campo obligatorio sin default
]
}
5 · Patrones de resiliencia
En un monolito, una llamada de método o funciona o lanza una excepción determinista. En un sistema distribuido existe un tercer resultado, mucho peor: no sabes qué ha pasado. La petición pudo no llegar, pudo llegar y procesarse pero perderse la respuesta, o pudo seguir en curso mientras tú ya te rendiste. Toda esta sección trata de convivir con esa incertidumbre.
5.1 Las ocho falacias de la computación distribuida
Formuladas en Sun Microsystems (Peter Deutsch y James Gosling, 1994–1997). Cada una es una suposición falsa que todo desarrollador hace la primera vez, y cada una tiene su factura correspondiente.
| # | Falacia | Realidad | Qué debes hacer |
|---|---|---|---|
| 1 | La red es fiable | Se pierden paquetes, se caen switches, hay particiones de red y reinicios de pods a diario. | Reintentos con backoff, idempotencia, circuit breakers. |
| 2 | La latencia es cero | 1 ms dentro del centro de datos, 100–300 ms entre continentes, y la varianza es peor que la media. | Menos saltos, llamadas en paralelo, caché, asincronía. |
| 3 | El ancho de banda es infinito | Devolver 50 MB de JSON satura el enlace y a los clientes móviles. | Paginación, compresión, proyecciones, formatos binarios. |
| 4 | La red es segura | Dentro del clúster también hay atacantes y errores de configuración. | mTLS, confianza cero, cifrado en tránsito y en reposo. |
| 5 | La topología no cambia | Los pods se mueven, las IP cambian y el autoescalado añade y quita réplicas cada minuto. | Descubrimiento de servicios, DNS con TTL corto, nada de IP fijas. |
| 6 | Hay un solo administrador | Cinco equipos, tres proveedores cloud y un SaaS que actualiza sin avisarte. | Contratos explícitos, versionado, ACL, degradación. |
| 7 | El coste de transporte es cero | Serializar y deserializar consume CPU; el tráfico entre zonas se factura. | Medir el coste por salto, colocalidad de zona, menos conversación innecesaria. |
| 8 | La red es homogénea | Conviven HTTP/1.1, HTTP/2, gRPC, MQTT, VPN y móviles en 3G. | Estándares, contratos explícitos, pruebas en condiciones degradadas. |
5.2 Timeouts: la primera línea de defensa
CÓMO ELEGIR UN TIMEOUT — con datos, no con números redondos
1. Mide la latencia REAL del dependiente por percentiles (no la media).
p50 = 20 ms p95 = 80 ms p99 = 150 ms p99.9 = 600 ms
2. Elige un poco por encima del p99 (o del p99.9 si el reintento es caro):
timeout ≈ p99 × 1,5 → ~225 ms
3. Comprueba el PRESUPUESTO TOTAL de la petición del usuario (deadline).
Presupuesto de la API: 1.000 ms
├─ Gateway 50 ms
├─ Pedidos 150 ms
│ ├─ Inventario 225 ms ┐
│ └─ Precios 200 ms ├─ en PARALELO: cuenta el mayor (225 ms)
└─ margen y red 200 ms ┘
Total: 50 + 150 + 225 + 200 = 625 ms ✅ cabe
4. LOS TIMEOUTS DEBEN DECRECER HACIA ABAJO. Nunca al revés.
✅ CORRECTO ❌ INCORRECTO (el error más común)
Cliente 10 s Cliente 2 s
Gateway 8 s Gateway 5 s
Pedidos 5 s Pedidos 10 s
Inv. 2 s Inv. 30 s
El de arriba siempre espera MÁS El cliente se rinde a los 2 s, pero abajo
que el de abajo: el fallo se detecta siguen 28 s de trabajo consumiendo hilos,
donde ocurre y se puede degradar. conexiones y CPU PARA NADIE. Así se produce
un colapso por saturación.
// Spring Boot 3: timeouts explícitos en RestClient (y en RestTemplate/WebClient)
@Bean
RestClient inventarioRestClient(RestClient.Builder builder) {
var factory = new SimpleClientHttpRequestFactory();
factory.setConnectTimeout(Duration.ofMillis(500)); // establecer la conexión TCP+TLS
factory.setReadTimeout(Duration.ofMillis(2_000)); // esperar la respuesta
return builder.baseUrl("http://inventario").requestFactory(factory).build();
}
// Con Apache HttpClient 5 (recomendado: pool configurable y timeout de adquisición)
@Bean
ClientHttpRequestFactory factoriaConPool() {
var pool = PoolingHttpClientConnectionManagerBuilder.create()
.setMaxConnTotal(200)
.setMaxConnPerRoute(50) // por host: evita monopolizar el pool
.build();
var config = RequestConfig.custom()
.setConnectionRequestTimeout(Timeout.ofMilliseconds(200)) // ¡el olvidado!
.setResponseTimeout(Timeout.ofMilliseconds(2_000))
.build();
var cliente = HttpClients.custom()
.setConnectionManager(pool)
.setDefaultRequestConfig(config)
.evictIdleConnections(TimeValue.ofSeconds(30))
.build();
return new HttpComponentsClientHttpRequestFactory(cliente);
}
# Los OTROS timeouts que también hay que poner y casi nadie pone
spring:
datasource:
hikari:
connection-timeout: 3000 # esperar una conexión libre del pool
validation-timeout: 2000
max-lifetime: 1800000
jpa:
properties:
jakarta.persistence.query.timeout: 5000 # ms: mata consultas eternas
kafka:
consumer:
properties:
max.poll.interval.ms: 300000 # si tardas más procesando, te expulsan del grupo
data:
redis:
timeout: 500ms
connect-timeout: 500ms
server:
tomcat:
connection-timeout: 20s
keep-alive-timeout: 20s
threads:
max: 200 # el techo real de concurrencia de tu servicio
shutdown: graceful # deja terminar las peticiones en curso al desplegar
5.3 Reintentos: potentes y peligrosos
POST /pagos que expiró por timeout puede cobrar dos veces al cliente. El
timeout no significa «no se hizo»: significa «no sé si se hizo». Sin idempotencia, no hay
reintentos.
REINTENTOS: LOS CUATRO ELEMENTOS OBLIGATORIOS
1. SOLO ERRORES TRANSITORIOS
✅ reintentar: timeout, 502/503/504, error de conexión, 429 (respetando Retry-After)
❌ NO reintentar: 400, 401, 403, 404, 409, 422 → el resultado será el mismo
2. BACKOFF EXPONENCIAL
intento 1 → espera 100 ms
intento 2 → espera 200 ms
intento 3 → espera 400 ms
intento 4 → espera 800 ms (con tope: nunca más de, por ejemplo, 5 s)
3. JITTER (aleatoriedad) — SIN ESTO SE PRODUCE LA "TORMENTA DE REINTENTOS"
Sin jitter: 1.000 clientes fallan en el mismo instante → los 1.000 reintentan
exactamente a los 100 ms → nuevo pico idéntico → el servicio nunca se recupera.
Sin jitter ▲ ▲ ▲ (picos sincronizados)
───┴──────────┴──────────┴────►
Con jitter ▁▂▃▂▁▂▃▂▁▃▂▁▂▃▁▂▃▂▁▃▂▁▂▃▁▂▃ (carga repartida)
─────────────────────────────►
espera = aleatorio(0, min(tope, base × 2^intento)) ← "full jitter" (AWS)
4. LÍMITE DE INTENTOS Y PRESUPUESTO
· Máximo 3 intentos en llamadas de usuario (más, y agotas el deadline).
· Presupuesto global: si más del 10 % del tráfico son reintentos, PÁRALOS.
· NUNCA reintentes en varios niveles a la vez: 3 × 3 × 3 = 27 llamadas reales.
❌ AMPLIFICACIÓN EN CASCADA
Gateway (3 intentos) → Pedidos (3) → Inventario (3) = 27 peticiones al servicio
que ya estaba saturado. Los reintentos MATARON al enfermo.
✅ Reintenta en UN solo nivel, el más cercano al fallo, y usa circuit breaker.
5.4 Circuit breaker (cortacircuitos)
Cuando un dependiente está caído, seguir llamándolo es peor que inútil: consume tus hilos, alarga tus latencias y le impide recuperarse. El cortacircuitos deja de llamar durante un tiempo y falla rápido, exactamente igual que el diferencial de tu casa.
MÁQUINA DE ESTADOS DEL CIRCUIT BREAKER
tasa de fallo >= umbral
(con nº mínimo de llamadas)
┌─────────────┐ ────────────────────────────► ┌────────────┐
│ │ │ │
│ CLOSED │ │ OPEN │
│ │ ◄───────────────────────────── │ │
│ pasan todas │ tasa de fallo < umbral │ rechaza YA │
│ las llamadas│ en el estado HALF_OPEN │ sin llamar │
└─────────────┘ └─────┬──────┘
▲ │
│ │ pasa
│ │ waitDurationInOpenState
│ ┌───────────────────┐ │ (por ejemplo, 30 s)
└────────│ HALF_OPEN │◄────────────────┘
│ │
éxito │ deja pasar N │──────► si fallan: vuelve a OPEN
suficiente │ llamadas de prueba│
└───────────────────┘
ESTADOS EXTRA DE RESILIENCE4J:
· DISABLED → siempre cerrado, sin decisión automática
· FORCED_OPEN → siempre abierto (útil para simular caídas en chaos testing)
· METRICS_ONLY → mide pero nunca abre (perfecto para estrenarlo en producción)
QUÉ CUENTA COMO "FALLO":
✅ Timeout, excepción de E/S, 5xx, llamada LENTA (slowCallRateThreshold)
❌ 4xx de negocio: un 404 NO es un fallo del dependiente. Si los cuentas,
abrirás el circuito por peticiones perfectamente correctas.
5.5 Bulkhead: compartimentos estancos
El nombre viene de los mamparos de un barco: si una sección se inunda, el resto sigue a flote. En software se trata de que un dependiente lento no se lleve por delante todos tus hilos.
SIN BULKHEAD — un dependiente lento consume TODO el pool de 200 hilos
┌────────────────────── Pool de hilos del servidor (200) ──────────────────────┐
│ ███████████████████████████████████████████████████████████████████████████ │
│ 198 hilos bloqueados esperando al servicio de RECOMENDACIONES (¡opcional!) │
│ ██ │
│ 2 hilos para el CHECKOUT, que es lo que da dinero │
└──────────────────────────────────────────────────────────────────────────────┘
Resultado: el 99 % del servicio cae por culpa de una funcionalidad prescindible.
CON BULKHEAD — cada dependiente tiene un presupuesto máximo de concurrencia
┌─────────────────┬─────────────────┬─────────────────┬───────────────────────┐
│ Pagos (crítico) │ Inventario │ Recomendaciones │ Resto de peticiones │
│ máx 40 llamadas │ máx 40 llamadas │ máx 10 llamadas │ │
│ ████████ │ ██████ │ ██████████ LLENO│ ████████████ │
└─────────────────┴─────────────────┴─────────────────┴───────────────────────┘
Recomendaciones está saturado → sus llamadas se rechazan al instante
(BulkheadFullException) y se DEGRADA esa sección. El checkout ni se entera.
DOS IMPLEMENTACIONES:
· Semáforo (barato): limita llamadas concurrentes en el hilo del llamante.
· Pool de hilos (aislamiento total): ejecuta en un pool propio; añade cambio de
contexto y complica la propagación del contexto (traza, seguridad).
Con virtual threads (Java 21) el semáforo suele ser suficiente y es lo recomendable.
5.6 Rate limiting y throttling
| Algoritmo | Cómo funciona | Ráfagas | Uso típico |
|---|---|---|---|
| Token bucket | Un cubo con capacidad B se rellena a R tokens por segundo. Cada petición consume un token; si no hay, se rechaza o espera. | Sí, hasta B de golpe | El más usado en API públicas: permite picos legítimos y limita la media. Es lo que usan Spring Cloud Gateway y Resilience4j. |
| Leaky bucket | Las peticiones entran en una cola y salen a ritmo constante. Si la cola se llena, se descartan. | No: alisa la salida | Proteger un recurso que no tolera picos: un sistema legado, un SaaS con cuota estricta. |
| Ventana fija | N peticiones por minuto natural. | Sí, mal: 2N en la frontera entre ventanas | Simple, aceptable para cuotas groseras. |
| Ventana deslizante | Cuenta las peticiones de los últimos 60 s reales. | Controladas | Más justo; más caro de calcular. |
| Concurrencia (semáforo) | Limita peticiones simultáneas, no por segundo. | — | Proteger recursos escasos (conexiones a BD). Es un bulkhead. |
429 Too Many Requests con la cabecera
Retry-After y, si puedes, RateLimit-Limit, RateLimit-Remaining y
RateLimit-Reset. Un 500 o un cierre de conexión hace que el cliente reintente
de inmediato y empeore la situación. Y limita por cliente (API key, usuario, tenant),
no globalmente: si no, un abusador deja fuera a todos los demás.
5.7 Fallback y degradación elegante
Fallar rápido está bien; fallar con dignidad está mucho mejor. Ante un dependiente caído, ordena las alternativas de mejor a peor:
- Valor cacheado, aunque esté algo obsoleto (stale-while-revalidate). Casi siempre es mejor un precio de hace 5 minutos que un error.
- Valor por defecto seguro: sin recomendaciones, muestra los más vendidos; sin cálculo de gastos de envío, aplica la tarifa estándar.
- Funcionalidad reducida: oculta la sección, muestra un aviso discreto, deshabilita el botón.
- Diferir: acepta la petición (
202 Accepted) y procésala cuando el dependiente vuelva. - Error explícito y honesto, con
traceIdy sin exponer detalles internos. Siempre mejor que colgarse 60 segundos.
saldo = 0 cuando el servicio de saldo está
caído (el usuario cree que le han robado), autorizar un pago porque el antifraude no responde (fallo
abierto en una decisión de seguridad) o devolver una lista vacía cuando en realidad hay datos
(el usuario cree que perdió sus pedidos). Regla: en decisiones de seguridad y dinero se falla
cerrado; en decisiones de experiencia, se falla abierto.
5.8 Load shedding y contrapresión
LOAD SHEDDING — rechazar pronto para poder servir a alguien
Sin load shedding: la latencia explota y NADIE recibe respuesta útil
carga ─────────────────────────────►
p99 ▁▁▁▂▂▃▄▆████████████████████ todas las peticiones tardan 30 s
éxito ████████████▇▅▃▁▁▁▁▁▁▁▁▁▁▁▁▁ y luego expiran en el cliente
▲
punto de saturación: la cola crece más rápido de lo que se vacía
Con load shedding: por encima del umbral se rechaza con 503 + Retry-After
éxito ███████████████████████████ el 80 % se sirve bien y rápido
429/503 ░░░░░░░░░░░░░░░░░░ el 20 % se rechaza en 1 ms
CÓMO DECIDIR QUÉ TIRAR (por orden):
1. Peticiones cuyo deadline ya venció (trabajo garantizadamente inútil).
2. Tráfico de baja prioridad: informes, prefetch, bots, reintentos.
3. Clientes que exceden su cuota.
4. Y solo al final: tráfico de usuario interactivo.
BACKPRESSURE (contrapresión): en lugar de descartar, se propaga "voy lento" hacia
el productor para que reduzca el ritmo.
· TCP lo hace con la ventana de recepción.
· Reactive Streams lo hace con request(n).
· Kafka lo hace de forma natural: el consumidor va a su ritmo y el lag crece.
· Una cola acotada + rechazo es la versión simple y honesta del backpressure.
LEY DE LITTLE — para dimensionar sin adivinar: L = λ × W
L = peticiones en el sistema (concurrencia)
λ = tasa de llegada (req/s)
W = tiempo de respuesta (s)
Ejemplo: 500 req/s × 0,2 s = 100 peticiones concurrentes → tu pool necesita ≥ 100
hilos (o virtual threads) y el pool de BD debe soportar la parte que llega a ella.
5.9 Idempotencia: el concepto que lo sostiene todo
Una operación es idempotente si ejecutarla N veces produce el mismo efecto que ejecutarla una vez. Es el requisito previo de los reintentos, de at-least-once y de casi todo lo que viene en las secciones 6 y 7.
| Operación | ¿Idempotente? | Por qué |
|---|---|---|
GET /pedidos/42 | Sí | Solo lee. |
PUT /pedidos/42 {estado:"ENVIADO"} | Sí | Fija un estado absoluto. |
DELETE /pedidos/42 | Sí | La segunda vez ya no está: devuelve 204 o 404, pero el efecto es el mismo. |
POST /pedidos | No | Crea un recurso nuevo cada vez → pedidos duplicados. |
PATCH /cuenta {saldo: -50} | No | Es un incremento relativo: dos veces resta 100. |
POST /pagos | No, y aquí duele | Doble cobro. Necesita clave de idempotencia obligatoriamente. |
CLAVE DE IDEMPOTENCIA — cómo se hace un POST seguro de reintentar
Cliente Servicio de Pagos BD
─────── ───────────────── ──
1. genera UUID una sola vez
Idempotency-Key: 7f3a...9c
POST /pagos {pedido:42, importe:100}
──────────────────────────────────►
2. INSERT en idempotencia
(clave, hash_peticion, EN_CURSO)
── si viola la PK → ya existe ──┐
3. ejecuta el cobro real │
4. guarda la respuesta y COMPLETADA │
◄────── 201 Created {pagoId} ────── │
│
5. TIMEOUT en el cliente: no sabe si se hizo │
REINTENTA con la MISMA clave 7f3a...9c │
──────────────────────────────────► │
6. la clave ya existe ◄─────────────┘
· si COMPLETADA → devuelve la
MISMA respuesta guardada (201)
· si EN_CURSO → 409 Conflict y
que reintente más tarde
· si el hash de la petición NO
coincide → 422: misma clave con
cuerpo distinto = error del cliente
◄────── 201 Created {pagoId} ────── (el mismo pagoId, sin cobrar dos veces)
-- La tabla que hace posible todo lo anterior. La clave primaria es el mecanismo:
-- la unicidad la garantiza la base de datos, no el código.
CREATE TABLE idempotencia (
clave VARCHAR(64) PRIMARY KEY,
endpoint VARCHAR(120) NOT NULL,
hash_peticion CHAR(64) NOT NULL, -- SHA-256 del cuerpo normalizado
estado VARCHAR(16) NOT NULL, -- EN_CURSO | COMPLETADA | FALLIDA
codigo_http SMALLINT,
respuesta JSONB,
creado_en TIMESTAMPTZ NOT NULL DEFAULT now(),
expira_en TIMESTAMPTZ NOT NULL -- 24-72 h; después se purga
);
CREATE INDEX idx_idempotencia_expira ON idempotencia (expira_en);
// Endpoint de pagos idempotente, completo y listo para copiar
@RestController
@RequestMapping("/api/v1/pagos")
class PagoController {
private final ServicioIdempotencia idempotencia;
private final CobrarUseCase cobrar;
@PostMapping
ResponseEntity<?> pagar(@RequestHeader("Idempotency-Key") @Size(min = 16, max = 64) String clave,
@Valid @RequestBody SolicitudPago solicitud) {
return idempotencia.ejecutar(clave, "POST /api/v1/pagos", solicitud, () -> {
var resultado = cobrar.ejecutar(solicitud.aComando());
return ResponseEntity.status(HttpStatus.CREATED).body(PagoResponse.desde(resultado));
});
}
}
@Service
class ServicioIdempotencia {
private final IdempotenciaRepository repo;
private final TransactionTemplate tx;
<T> ResponseEntity<?> ejecutar(String clave, String endpoint, Object peticion,
Supplier<ResponseEntity<T>> operacion) {
String hash = sha256(escribirCanonico(peticion));
Optional<RegistroIdempotencia> existente = repo.findById(clave);
if (existente.isPresent()) {
var reg = existente.get();
if (!reg.hashPeticion().equals(hash)) {
throw new ClaveIdempotenciaReutilizadaException(clave); // → 422
}
return switch (reg.estado()) {
case COMPLETADA -> ResponseEntity.status(reg.codigoHttp())
.header("Idempotent-Replay", "true")
.body(reg.respuesta());
case EN_CURSO -> ResponseEntity.status(HttpStatus.CONFLICT)
.header("Retry-After", "1").build();
case FALLIDA -> ResponseEntity.status(reg.codigoHttp()).body(reg.respuesta());
};
}
try {
// La restricción PRIMARY KEY resuelve la carrera entre dos peticiones simultáneas.
// REQUIRES_NEW: el marcador debe persistir aunque la operación de negocio falle.
tx.executeWithoutResult(s -> repo.insertarEnCurso(clave, endpoint, hash));
} catch (DataIntegrityViolationException carrera) {
return ResponseEntity.status(HttpStatus.CONFLICT).header("Retry-After", "1").build();
}
try {
ResponseEntity<T> respuesta = operacion.get();
repo.completar(clave, respuesta.getStatusCode().value(), escribir(respuesta.getBody()));
return respuesta;
} catch (ReglaDeNegocioException e) {
repo.fallar(clave, 422, escribir(Map.of("detail", e.getMessage()))); // error determinista
throw e;
} catch (RuntimeException e) {
repo.borrar(clave); // error transitorio: que el cliente pueda reintentar de verdad
throw e;
}
}
}
UNIQUE de negocio (UNIQUE(pedido_id, tipo_pago)) o con que el
cliente genere el identificador del recurso (PUT /pagos/{uuid}), lo que
convierte la creación en idempotente por diseño. Es la solución más elegante: no hay estado extra que
mantener ni purgar.
5.10 Resilience4j en Spring Boot 3: configuración completa
Resilience4j es la librería estándar en el ecosistema Spring desde que Hystrix entró en mantenimiento. Es modular, funcional, sin dependencias pesadas y se integra con Micrometer y con Spring Boot mediante anotaciones.
<!-- pom.xml — Spring Boot 3.x + Java 21 -->
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-spring-boot3</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-aop</artifactId> <!-- necesario para las anotaciones -->
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId> <!-- métricas y health -->
</dependency>
resilience4j:
circuitbreaker:
configs:
default:
sliding-window-type: COUNT_BASED # o TIME_BASED (últimos N segundos)
sliding-window-size: 100 # últimas 100 llamadas
minimum-number-of-calls: 20 # no decidir con 2 llamadas: sería ruido
failure-rate-threshold: 50 # % de fallos que abre el circuito
slow-call-duration-threshold: 2s # una llamada más lenta cuenta como "lenta"
slow-call-rate-threshold: 80 # % de llamadas lentas que también abre
wait-duration-in-open-state: 30s
permitted-number-of-calls-in-half-open-state: 5
automatic-transition-from-open-to-half-open-enabled: true
register-health-indicator: true # aparece en /actuator/health
record-exceptions: # QUÉ cuenta como fallo
- java.io.IOException
- java.util.concurrent.TimeoutException
- org.springframework.web.client.HttpServerErrorException
ignore-exceptions: # QUÉ NO cuenta (¡los 4xx de negocio!)
- com.ejemplo.dominio.ReglaDeNegocioException
- com.ejemplo.dominio.PedidoNoEncontradoException
instances:
inventario:
base-config: default
pagos:
base-config: default
failure-rate-threshold: 30 # con el dinero somos más conservadores
wait-duration-in-open-state: 60s
retry:
configs:
default:
max-attempts: 3 # 1 original + 2 reintentos
wait-duration: 200ms
enable-exponential-backoff: true
exponential-backoff-multiplier: 2 # 200 ms, 400 ms, 800 ms
exponential-max-wait-duration: 5s
enable-randomized-wait: true # JITTER: imprescindible
randomized-wait-factor: 0.5
retry-exceptions:
- java.io.IOException
- java.util.concurrent.TimeoutException
ignore-exceptions:
- com.ejemplo.dominio.ReglaDeNegocioException
instances:
inventario: { base-config: default }
pagos: { max-attempts: 1 } # ❗ no reintentamos cobros sin idempotencia
timelimiter:
configs:
default:
timeout-duration: 2s
cancel-running-future: true
instances:
inventario: { base-config: default }
pagos: { timeout-duration: 5s }
bulkhead: # semáforo: limita concurrencia
instances:
recomendaciones:
max-concurrent-calls: 10
max-wait-duration: 0 # 0 = rechazar de inmediato, no encolar
inventario:
max-concurrent-calls: 40
max-wait-duration: 20ms
thread-pool-bulkhead: # aislamiento total con pool propio
instances:
informes:
max-thread-pool-size: 8
core-thread-pool-size: 4
queue-capacity: 20
ratelimiter:
instances:
apiExterna:
limit-for-period: 100 # 100 llamadas...
limit-refresh-period: 1s # ...por segundo (token bucket)
timeout-duration: 0 # no esperar: fallar rápido
register-health-indicator: true
management:
endpoints.web.exposure.include: health,metrics,prometheus,circuitbreakers,circuitbreakerevents
endpoint.health.show-details: always
health.circuitbreakers.enabled: true
metrics.distribution.percentiles-histogram.resilience4j.circuitbreaker.calls: true
// Uso con anotaciones. El ORDEN de los decoradores importa muchísimo.
@Component
class InventarioClienteResiliente implements ConsultaInventario {
private static final Logger log = LoggerFactory.getLogger(InventarioClienteResiliente.class);
private final InventarioApi api; // HTTP interface declarativa
private final CacheStock cache;
/**
* Orden por defecto de los aspectos en Spring (de fuera hacia dentro):
*
* Retry ( CircuitBreaker ( RateLimiter ( TimeLimiter ( Bulkhead ( llamada ) ) ) ) )
*
* Es el orden correcto y casi nunca hay que cambiarlo:
* · El bulkhead protege primero el recurso.
* · El timeout corta la llamada individual.
* · El circuit breaker ve el resultado de cada intento REAL.
* · El retry es lo más externo: reintenta la operación completa.
* Si pusieras Retry DENTRO del CircuitBreaker, el breaker vería 1 fallo donde hubo 3
* y tardaría el triple en abrirse.
* Se ajusta con resilience4j.retry.retry-aspect-order y equivalentes.
*/
@Retry(name = "inventario")
@CircuitBreaker(name = "inventario", fallbackMethod = "stockDesdeCache")
@Bulkhead(name = "inventario")
@Override
public boolean haySuficiente(List<LineaPedido> lineas) {
return api.consultarStock(lineas.stream().map(l -> l.sku().valor()).toList())
.stream().allMatch(StockDto::disponible);
}
/**
* El fallback debe tener la MISMA firma más un parámetro Throwable al final.
* Puedes declarar varios, del más específico al más genérico.
*/
@SuppressWarnings("unused")
private boolean stockDesdeCache(List<LineaPedido> lineas, CallNotPermittedException abierto) {
log.warn("Circuito de inventario ABIERTO: sirviendo stock cacheado");
return cache.optimista(lineas); // degradación consciente y medida
}
@SuppressWarnings("unused")
private boolean stockDesdeCache(List<LineaPedido> lineas, Throwable t) {
log.warn("Fallo consultando inventario ({}): sirviendo stock cacheado", t.toString());
return cache.optimista(lineas);
}
}
// Uso programático: más control, sin AOP, y funciona en llamadas internas
// (recuerda: las anotaciones de Spring NO se aplican a llamadas dentro de la misma clase)
@Configuration
class ResilienciaProgramaticaConfig {
@Bean
Decoradores decoradores(CircuitBreakerRegistry cbRegistry,
RetryRegistry retryRegistry,
BulkheadRegistry bhRegistry) {
CircuitBreaker cb = cbRegistry.circuitBreaker("pagos");
Retry retry = retryRegistry.retry("pagos");
Bulkhead bh = bhRegistry.bulkhead("pagos");
// Eventos: aquí es donde se entera la alerta de que algo va mal
cb.getEventPublisher()
.onStateTransition(e -> LoggerFactory.getLogger("resiliencia")
.warn("Circuit breaker {}: {} -> {}", e.getCircuitBreakerName(),
e.getStateTransition().getFromState(), e.getStateTransition().getToState()))
.onCallNotPermitted(e -> contadorRechazos.increment());
return new Decoradores(cb, retry, bh);
}
}
// Composición explícita (la ventaja es que el orden se ve)
Supplier<Recibo> decorada = Decorators
.ofSupplier(() -> pasarela.cobrar(solicitud))
.withBulkhead(bulkhead)
.withCircuitBreaker(circuitBreaker)
.withRetry(retry)
.withFallback(List.of(CallNotPermittedException.class), e -> Recibo.diferido(solicitud))
.decorate();
Recibo recibo = decorada.get();
@TimeLimiter: solo funciona sobre métodos que devuelven
CompletableFuture o tipos reactivos, porque necesita poder cancelar la ejecución desde
fuera. En un método bloqueante y síncrono, @TimeLimiter no hace nada: el
timeout debe configurarse en el cliente HTTP. Es uno de los fallos más frecuentes en revisiones de
código.
5.11 Probar la resiliencia: WireMock y chaos engineering
Un patrón de resiliencia no probado es una decoración. Los tres escenarios que hay que tener automatizados son: dependiente lento, dependiente que devuelve 500 y dependiente que no responde en absoluto.
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
@AutoConfigureWireMock(port = 0) // spring-cloud-contract-wiremock
class InventarioResilienciaTest {
@Autowired ConsultaInventario inventario;
@Autowired CircuitBreakerRegistry registry;
@BeforeEach void reset() { registry.circuitBreaker("inventario").reset(); }
@Test
void aplica_el_timeout_cuando_el_dependiente_es_lento() {
stubFor(post(urlEqualTo("/api/stock"))
.willReturn(aResponse().withStatus(200)
.withFixedDelay(5_000))); // 5 s frente a un timeout de 2 s
long inicio = System.nanoTime();
boolean resultado = inventario.haySuficiente(List.of(unaLinea()));
Duration transcurrido = Duration.ofNanos(System.nanoTime() - inicio);
assertThat(transcurrido).isLessThan(Duration.ofSeconds(3)); // no esperó los 5 s
assertThat(resultado).isTrue(); // sirvió el fallback
}
@Test
void abre_el_circuito_tras_superar_el_umbral_de_fallos() {
stubFor(post(urlEqualTo("/api/stock")).willReturn(aResponse().withStatus(503)));
for (int i = 0; i < 25; i++) inventario.haySuficiente(List.of(unaLinea()));
assertThat(registry.circuitBreaker("inventario").getState())
.isEqualTo(CircuitBreaker.State.OPEN);
}
@Test
void deja_de_llamar_al_dependiente_mientras_el_circuito_esta_abierto() {
registry.circuitBreaker("inventario").transitionToOpenState();
inventario.haySuficiente(List.of(unaLinea()));
verify(0, postRequestedFor(urlEqualTo("/api/stock"))); // ni una sola llamada real
}
@Test
void reintenta_solo_los_errores_transitorios() {
stubFor(post(urlEqualTo("/api/stock")).inScenario("flaky")
.whenScenarioStateIs(STARTED)
.willReturn(aResponse().withStatus(503))
.willSetStateTo("segundo"));
stubFor(post(urlEqualTo("/api/stock")).inScenario("flaky")
.whenScenarioStateIs("segundo")
.willReturn(okJson("[{\"sku\":\"A1\",\"disponible\":true}]")));
assertThat(inventario.haySuficiente(List.of(unaLinea()))).isTrue();
verify(2, postRequestedFor(urlEqualTo("/api/stock"))); // reintentó exactamente una vez
}
}
# Chaos engineering básico: experimentos con hipótesis, no "romper cosas por diversión"
#
# 1. HIPÓTESIS: "si el servicio de recomendaciones cae, el checkout sigue funcionando
# con p99 < 500 ms y una tasa de error < 0,1 %".
# 2. RADIO DE EXPLOSIÓN pequeño: un pod, un 1 % del tráfico, en horario laboral,
# con el equipo delante y un botón de parada.
# 3. MEDIR antes, durante y después. 4. AUTOMATIZAR el experimento en CI.
# a) Matar un pod al azar
kubectl get pod -l app=recomendaciones --field-selector=status.phase=Running \
-o name | shuf -n1 | xargs kubectl delete
# b) Inyectar latencia con Istio (VirtualService)
# fault: { delay: { percentage: { value: 50 }, fixedDelay: 5s } }
# c) Chaos Monkey for Spring Boot: latencia y excepciones DENTRO de la aplicación
# (perfil de pruebas, jamás en el perfil de producción)
java -jar app.jar \
--spring.profiles.active=chaos-monkey \
--chaos.monkey.enabled=true \
--chaos.monkey.watcher.service=true \
--chaos.monkey.assaults.level=5 \
--chaos.monkey.assaults.latencyActive=true \
--chaos.monkey.assaults.latencyRangeStartInMillis=2000 \
--chaos.monkey.assaults.exceptionsActive=true
# d) Cortar la red a una dependencia con Toxiproxy en un test de integración
# (Testcontainers tiene módulo: ToxiproxyContainer)
Checklist — comunicación y resiliencia
6 · Datos en sistemas distribuidos
Aquí está el 90 % de la dificultad real de los microservicios. La comunicación se resuelve con librerías; los datos distribuidos se resuelven con decisiones de negocio.
6.1 Una base de datos por servicio
Es el principio no negociable: si dos servicios comparten tablas, no son dos servicios. La base de datos es el detalle de implementación más íntimo de un servicio; exponerla equivale a hacer públicos todos sus campos privados y renunciar para siempre a refactorizar.
❌ BASE DE DATOS COMPARTIDA — el acoplamiento invisible
┌───────────┐ ┌───────────┐ ┌───────────┐
│ Pedidos │ │ Pagos │ │ Envíos │
└─────┬─────┘ └─────┬─────┘ └─────┬─────┘
│ SELECT/UPDATE │ │ Cada servicio hace SELECT y UPDATE
└────────┬───────┴───────────────┘ sobre las tablas de los demás.
▼
┌────────────────────┐ Consecuencias REALES:
│ BD compartida │ · Añadir una columna NOT NULL rompe a otros dos equipos.
│ pedido │ · Renombrar una tabla es un proyecto de 3 meses.
│ pago │ · Un bloqueo de Pagos frena a Pedidos.
│ envio │ · Nadie sabe quién escribe qué: la integridad es un mito.
│ cliente │ · No puedes escalar ni migrar una parte por separado.
└────────────────────┘ · Es un MONOLITO con latencia de red. Sin ventajas.
✅ UNA BASE DE DATOS POR SERVICIO
┌───────────┐ ┌───────────┐ ┌───────────┐
│ Pedidos │──API/──│ Pagos │─evento─│ Envíos │
└─────┬─────┘ evento └─────┬─────┘ └─────┬─────┘
▼ ▼ ▼
┌────────────┐ ┌───────────┐ ┌───────────┐
│ BD pedidos │ │ BD pagos │ │ BD envíos │ Cada equipo elige motor,
│ PostgreSQL │ │PostgreSQL │ │ MongoDB │ esquema, índices y migra
└────────────┘ └───────────┘ └───────────┘ cuando quiere.
Réplica de datos ajenos: Envíos guarda una COPIA de los datos del pedido que
necesita (dirección, líneas), recibida por evento. No es "duplicar datos mal":
es la forma correcta de eliminar una dependencia en tiempo de ejecución.
JOIN».
Correcto, y esa es exactamente la incomodidad que te avisa de que la frontera puede estar mal puesta. Si
necesitas hacer joins constantemente entre dos servicios, probablemente sean un solo
servicio. Si solo lo necesitas para informes, la respuesta es una vista materializada o un
almacén analítico, no romper el aislamiento.
6.2 Consistencia eventual y sus implicaciones para el negocio
Consistencia eventual significa que, si dejas de escribir, todas las réplicas convergerán al mismo valor… en algún momento. Ese «algún momento» suele ser de milisegundos, pero puede ser de minutos si algo va mal, y el diseño debe contemplarlo.
EJEMPLO CONCRETO DE UX — el usuario cancela un pedido
t=0 ms Usuario pulsa "Cancelar"
t=15 ms Pedidos marca CANCELADO y responde 200 OK. Publica PedidoCancelado.
t=18 ms La UI muestra "Pedido cancelado" ✅
t=40 ms Facturación consume el evento y anula la factura.
t=95 ms Envíos consume el evento y cancela la etiqueta de transporte.
t=110 ms Analítica actualiza el panel.
⚠️ VENTANA DE INCOHERENCIA: entre t=18 y t=110 ms, si el usuario navega a
"Mis facturas", PUEDE VER la factura del pedido que acaba de cancelar.
CÓMO SE GESTIONA (de mejor a peor):
1. DISEÑAR LA UI PARA ELLO: "Cancelación en curso. La factura se anulará en unos
segundos." El usuario tolera perfectamente lo que se le explica.
2. LEER TUS PROPIAS ESCRITURAS (read-your-writes): tras escribir, lee del servicio
propietario o de la réplica primaria durante unos segundos.
3. BLOQUEO OPTIMISTA EN LA UI: pintar el estado esperado localmente y reconciliar.
4. ❌ Poner un sleep(2000) y rezar.
CONVERSACIÓN OBLIGATORIA CON NEGOCIO:
"¿Cuánto tiempo puede pasar entre que se cancela un pedido y se anula la factura?"
· "Debe ser instantáneo, es legal" → MISMO servicio, misma transacción.
· "Unos segundos está bien" → eventos, consistencia eventual.
· "Con que sea el mismo día, vale" → proceso por lotes nocturno; aún más simple.
Esta pregunta define la arquitectura. No la contestes tú solo.
6.3 Por qué 2PC casi nunca es la respuesta
El commit en dos fases (2PC, con XA/JTA) promete transacciones ACID entre varios recursos. En teoría es la solución perfecta; en la práctica se abandonó por buenas razones.
COMMIT EN DOS FASES (2PC)
FASE 1 — preparación FASE 2 — confirmación
Coordinador Coordinador
│──── prepare ────► BD Pedidos │──── commit ────► BD Pedidos
│──── prepare ────► BD Pagos │──── commit ────► BD Pagos
│──── prepare ────► Broker │──── commit ────► Broker
│◄─── listo ─────── │◄─── ok ─────────
(todos bloquean recursos (si alguno falla aquí,
y esperan) el sistema queda INCIERTO)
PROBLEMAS QUE LO DESCARTAN EN MICROSERVICIOS:
1. BLOQUEOS LARGOS: los recursos quedan bloqueados durante toda la ida y vuelta.
El rendimiento cae un orden de magnitud.
2. EL COORDINADOR ES UN PUNTO ÚNICO DE FALLO: si muere entre fase 1 y 2, los
participantes quedan "in doubt", bloqueados, hasta intervención MANUAL.
3. DISPONIBILIDAD MULTIPLICATIVA: hace falta que TODOS estén vivos a la vez.
4. SOPORTE ESCASO: Kafka, la mayoría de NoSQL y casi todas las API REST no
hablan XA. Solo lo soportan bases relacionales y algunos brokers JMS.
5. ACOPLAMIENTO TEMPORAL TOTAL: justo lo contrario de lo que buscas.
¿CUÁNDO SÍ? Dentro de un mismo servicio, con 2 recursos que soportan XA, con volumen
bajo, en un entorno controlado y sin alternativa. Es decir: casi nunca.
La alternativa moderna se llama SAGA.
6.4 El patrón Saga
Una saga es una secuencia de transacciones locales. Cada paso confirma en su propia base de datos y publica un evento que dispara el siguiente. Si un paso falla, se ejecutan transacciones de compensación que deshacen semánticamente lo hecho.
SAGA COREOGRAFIADA — cada servicio reacciona a eventos; no hay director
Pedidos Pagos Inventario Envíos
│ │ │ │
│ PedidoConfirmado │ │
├───────────────►│ │ │
│ │ cobra │ │
│ │ PagoRealizado │ │
│ ├─────────────────►│ │
│ │ │ reserva stock │
│ │ │ StockReservado │
│ │ ├───────────────►│
│ │ │ │ crea envío
│ │ │ │ EnvioCreado
│◄───────────────┴──────────────────┴────────────────┤
│ marca COMPLETADO │
── CAMINO DE COMPENSACIÓN (no hay stock) ──
│ │ │ StockNoDisponible
│ │◄─────────────────┤
│ │ reembolsa │
│ │ PagoReembolsado │
│◄───────────────┤ │
│ marca CANCELADO_SIN_STOCK │
✅ Sin punto único de fallo, muy desacoplado, fácil de empezar.
❌ La lógica del proceso está REPARTIDA: nadie puede responder "¿en qué punto
está la saga del pedido 42?" sin mirar cinco servicios. Con más de 4 pasos
se vuelve inmanejable y aparecen ciclos de eventos difíciles de razonar.
SAGA ORQUESTADA — un coordinador explícito dirige y conoce el estado
┌──────────────────────────────────┐
│ ORQUESTADOR DE PEDIDO (saga) │
│ estado: ESPERANDO_PAGO │
│ máquina de estados persistida │
└───┬───────┬──────────┬───────────┘
1. CobrarPago │ │ │ 3. CrearEnvio
┌────────────────┘ │ └──────────────┐
▼ ▼ 2. ReservarStock ▼
┌─────────┐ ┌────────────┐ ┌─────────┐
│ Pagos │ │ Inventario │ │ Envíos │
└────┬────┘ └─────┬──────┘ └────┬────┘
│ PagoRealizado │ StockReservado │ EnvioCreado
└────────────────────────┴────────────────────────┘
(respuestas al orquestador)
COMPENSACIONES en orden INVERSO al de ejecución:
falla CrearEnvio → LiberarStock → ReembolsarPago → marcar pedido CANCELADO
✅ El estado de la saga es CONSULTABLE y observable; la lógica está en un sitio;
los timeouts por paso son fáciles; se puede reintentar un paso concreto.
❌ Un componente más que mantener; riesgo de convertirse en un "dios" con toda
la lógica de negocio dentro (mantenlo como coordinador, no como cerebro).
| Criterio | Coreografiada | Orquestada |
|---|---|---|
| Nº de pasos recomendado | 2–4 | 4 o más |
| Acoplamiento | Mínimo | El orquestador conoce a todos |
| Visibilidad del estado | Baja: hay que reconstruirla con trazas | Alta: una fila en una tabla |
| Riesgo | Ciclos de eventos, lógica difusa | Orquestador anémico o «dios» |
| Timeouts por paso | Complicados | Triviales |
| Herramientas | Kafka y listeners | Spring Statemachine, Temporal, Camunda/Zeebe, Axon |
// Orquestador de saga con máquina de estados persistida. Simplificado pero completo.
@Service
class OrquestadorSagaPedido {
private static final Logger log = LoggerFactory.getLogger(OrquestadorSagaPedido.class);
private final SagaRepository sagas;
private final ComandoPublisher comandos;
/** Paso 0: la saga arranca al confirmarse el pedido. */
@Transactional
@KafkaListener(topics = "pedidos.eventos", groupId = "saga-pedido")
public void alConfirmarPedido(PedidoConfirmado evento) {
var saga = SagaPedido.iniciar(evento.pedidoId(), evento.total(), evento.lineas());
sagas.save(saga);
comandos.enviar(new CobrarPago(saga.id(), evento.pedidoId(), evento.total()));
}
/** Paso 1 OK → paso 2. */
@Transactional
@KafkaListener(topics = "pagos.eventos", groupId = "saga-pedido")
public void alRealizarsePago(PagoRealizado evento) {
var saga = cargar(evento.sagaId());
if (!saga.enEstado(EstadoSaga.ESPERANDO_PAGO)) {
log.info("Evento fuera de orden o duplicado para la saga {}: ignorado", saga.id());
return; // IDEMPOTENCIA: los eventos se repiten
}
saga.registrarPagoRealizado(evento.pagoId());
comandos.enviar(new ReservarStock(saga.id(), saga.lineas()));
}
/** Paso 1 KO → no hay nada previo que compensar. */
@Transactional
@KafkaListener(topics = "pagos.eventos", groupId = "saga-pedido")
public void alRechazarsePago(PagoRechazado evento) {
var saga = cargar(evento.sagaId());
saga.fallar(MotivoFallo.PAGO_RECHAZADO);
comandos.enviar(new CancelarPedido(saga.pedidoId(), MotivoCancelacion.PAGO_RECHAZADO));
}
/** Paso 2 KO → COMPENSAR el paso 1 en orden inverso. */
@Transactional
@KafkaListener(topics = "inventario.eventos", groupId = "saga-pedido")
public void alFallarLaReserva(StockNoDisponible evento) {
var saga = cargar(evento.sagaId());
saga.iniciarCompensacion(MotivoFallo.SIN_STOCK);
comandos.enviar(new ReembolsarPago(saga.id(), saga.pagoId(),
"Sin stock para el pedido " + saga.pedidoId()));
}
/**
* Vigilante de timeouts: sin esto, una saga puede quedarse colgada para siempre
* porque un evento se perdió. Es una pieza OBLIGATORIA, no un extra.
*/
@Scheduled(fixedDelay = 30_000)
@Transactional
public void compensarSagasAtascadas() {
Instant limite = Instant.now().minus(Duration.ofMinutes(5));
for (SagaPedido saga : sagas.findAtascadasAntesDe(limite)) {
log.error("Saga {} atascada en {} desde {}: compensando",
saga.id(), saga.estado(), saga.actualizadaEn());
saga.iniciarCompensacion(MotivoFallo.TIMEOUT);
compensarSegunEstado(saga);
}
}
private SagaPedido cargar(UUID id) {
return sagas.findById(id).orElseThrow(() -> new SagaNoEncontradaException(id));
}
}
-- El estado de la saga es una tabla: consultable, auditable y reparable a mano si hace falta
CREATE TABLE saga_pedido (
id UUID PRIMARY KEY,
pedido_id UUID NOT NULL UNIQUE,
estado VARCHAR(32) NOT NULL, -- INICIADA, ESPERANDO_PAGO, ESPERANDO_STOCK,
-- ESPERANDO_ENVIO, COMPLETADA, COMPENSANDO, FALLIDA
paso_actual SMALLINT NOT NULL,
pago_id UUID,
reserva_id UUID,
envio_id UUID,
motivo_fallo VARCHAR(64),
intentos SMALLINT NOT NULL DEFAULT 0,
creada_en TIMESTAMPTZ NOT NULL DEFAULT now(),
actualizada_en TIMESTAMPTZ NOT NULL DEFAULT now(),
version BIGINT NOT NULL DEFAULT 0 -- bloqueo optimista
);
-- Índice para el vigilante de sagas atascadas
CREATE INDEX idx_saga_atascadas ON saga_pedido (actualizada_en)
WHERE estado NOT IN ('COMPLETADA','FALLIDA');
6.5 El patrón Outbox transaccional
Este es, probablemente, el patrón más importante de todo el módulo, porque resuelve un problema que aparece siempre y que casi todos los equipos descubren tarde: la doble escritura (dual write).
EL PROBLEMA DEL DUAL WRITE — no hay forma de hacer esto bien "a mano"
@Transactional
void confirmar(PedidoId id) {
pedido.confirmar();
repositorio.guardar(pedido); // (1) escribe en la BD
kafkaTemplate.send("pedidos", evento); // (2) escribe en Kafka ← ¡otro sistema!
}
Cuatro finales posibles, dos de ellos catastróficos:
┌─────────────┬─────────────┬───────────────────────────────────────────────┐
│ BD │ Kafka │ Resultado │
├─────────────┼─────────────┼───────────────────────────────────────────────┤
│ ✅ commit │ ✅ enviado │ Correcto. │
│ ❌ rollback │ ❌ no envío │ Correcto (no pasó nada). │
│ ✅ commit │ ❌ FALLA │ 💥 El pedido está confirmado y NADIE se entera│
│ │ │ Nunca se factura ni se envía. │
│ ❌ rollback │ ✅ ENVIADO │ 💥 Se factura y se envía un pedido que NO │
│ │ │ existe en la base de datos. │
└─────────────┴─────────────┴───────────────────────────────────────────────┘
Invertir el orden no ayuda. Poner el send() fuera de la transacción tampoco.
Reintentar tampoco: el proceso puede morir justo entre (1) y (2).
NO EXISTE una solución con dos escrituras a dos sistemas sin un protocolo extra.
LA SOLUCIÓN: OUTBOX — escribir en UN solo sistema transaccional
┌──────────────────── UNA SOLA TRANSACCIÓN DE BASE DE DATOS ────────────────────┐
│ │
│ INSERT/UPDATE pedido + INSERT INTO outbox (evento) │
│ ──────────────────── ───────────────────────── │
│ Ambas cosas confirman juntas o se descartan juntas. Atomicidad REAL. │
└───────────────────────────────────────────────────────────────────────────────┘
│
▼
┌──────────────────────────────────────┐
│ PUBLICADOR (uno de dos sabores) │
├──────────────────────────────────────┤
│ A) Polling: cada 200 ms hace │
│ SELECT ... WHERE publicado IS NULL│
│ FOR UPDATE SKIP LOCKED │
│ → envía a Kafka → marca publicado │
│ │
│ B) CDC con Debezium: lee el WAL de │
│ PostgreSQL y publica sin tocar la │
│ aplicación. Menos latencia y │
│ menos carga en la BD. │
└───────────────┬──────────────────────┘
▼
[ KAFKA ] → consumidores
GARANTÍA RESULTANTE: at-least-once. El evento puede publicarse DOS veces si el
publicador muere justo después de enviar y antes de marcar. Por eso el consumidor
DEBE ser idempotente (sección 6.6).
CREATE TABLE outbox (
id UUID PRIMARY KEY,
agregado_tipo VARCHAR(64) NOT NULL, -- "Pedido" → sirve de topic/routing
agregado_id VARCHAR(64) NOT NULL, -- clave de partición: ORDEN por agregado
tipo_evento VARCHAR(128) NOT NULL, -- "pedidos.PedidoConfirmado.v1"
payload JSONB NOT NULL,
cabeceras JSONB NOT NULL DEFAULT '{}', -- traceId, tenant, versión
creado_en TIMESTAMPTZ NOT NULL DEFAULT now(),
publicado_en TIMESTAMPTZ,
intentos SMALLINT NOT NULL DEFAULT 0
);
-- Índice PARCIAL: solo indexa lo pendiente. La tabla puede tener millones de filas
-- publicadas y el índice sigue siendo diminuto.
CREATE INDEX idx_outbox_pendiente ON outbox (creado_en) WHERE publicado_en IS NULL;
-- Purga: sin esto la tabla crece sin límite (una partición por día es aún mejor)
DELETE FROM outbox WHERE publicado_en < now() - INTERVAL '7 days';
// ─── Adaptador de salida: implementa el puerto PublicadorEventos escribiendo en la outbox ───
@Component
class PublicadorEventosOutbox implements PublicadorEventos {
private final OutboxRepository outbox;
private final ObjectMapper mapper;
private final Tracer tracer;
@Override
public void publicar(List<EventoDominio> eventos) {
// Se ejecuta DENTRO de la transacción del caso de uso: ese es todo el truco.
for (EventoDominio evento : eventos) {
outbox.save(new OutboxEntity(
evento.eventoId(),
"Pedido",
idDeAgregado(evento),
evento.tipo(),
serializar(evento),
cabecerasConTraza())); // propagar el traceId al consumidor
}
}
private Map<String, String> cabecerasConTraza() {
var span = tracer.currentSpan();
return span == null ? Map.of()
: Map.of("traceparent", "00-%s-%s-01".formatted(
span.context().traceId(), span.context().spanId()));
}
}
// ─── Publicador por polling: sencillo, sin infraestructura extra ───
@Component
class PublicadorOutboxScheduler {
private static final Logger log = LoggerFactory.getLogger(PublicadorOutboxScheduler.class);
private static final int LOTE = 100;
private final OutboxRepository outbox;
private final KafkaTemplate<String, byte[]> kafka;
/**
* fixedDelay corto: la latencia de publicación es, como mucho, este intervalo.
* SKIP LOCKED permite que varias réplicas trabajen en paralelo SIN duplicar
* ni bloquearse entre ellas: cada una coge filas distintas.
*/
@Scheduled(fixedDelay = 200)
@Transactional
public void publicarPendientes() {
List<OutboxEntity> pendientes = outbox.bloquearPendientes(LOTE);
if (pendientes.isEmpty()) return;
for (OutboxEntity fila : pendientes) {
var registro = new ProducerRecord<>(
"pedidos.eventos",
null, // partición: la decide la clave
fila.agregadoId(), // CLAVE = id del agregado → orden garantizado
fila.payload());
fila.cabeceras().forEach((k, v) ->
registro.headers().add(k, v.getBytes(StandardCharsets.UTF_8)));
registro.headers().add("tipo", fila.tipoEvento().getBytes(StandardCharsets.UTF_8));
registro.headers().add("id", fila.id().toString().getBytes(StandardCharsets.UTF_8));
try {
kafka.send(registro).get(5, TimeUnit.SECONDS); // esperamos confirmación real
fila.marcarPublicado(Instant.now());
} catch (Exception e) {
fila.incrementarIntentos();
log.error("No se pudo publicar el evento {} (intento {})", fila.id(), fila.intentos(), e);
break; // no seguimos con el lote: preservamos el ORDEN
}
}
}
}
interface OutboxRepository extends JpaRepository<OutboxEntity, UUID> {
@Query(value = """
SELECT * FROM outbox
WHERE publicado_en IS NULL
ORDER BY creado_en
LIMIT :limite
FOR UPDATE SKIP LOCKED
""", nativeQuery = true)
List<OutboxEntity> bloquearPendientes(@Param("limite") int limite);
}
# ─── Alternativa B: CDC con Debezium. Sin polling y sin carga extra en la BD ───
# Conector de Debezium para PostgreSQL con el SMT de outbox
name: pedidos-outbox-connector
config:
connector.class: io.debezium.connector.postgresql.PostgresConnector
database.hostname: postgres-pedidos
database.dbname: pedidos
plugin.name: pgoutput # decodificación lógica nativa de PostgreSQL 10+
slot.name: pedidos_outbox_slot
publication.autocreate.mode: filtered
table.include.list: pedidos.outbox
tombstones.on.delete: "false"
# El SMT de outbox convierte cada fila insertada en un evento con la forma correcta
transforms: outbox
transforms.outbox.type: io.debezium.transforms.outbox.EventRouter
transforms.outbox.table.field.event.id: id
transforms.outbox.table.field.event.key: agregado_id # → clave de partición
transforms.outbox.table.field.event.type: tipo_evento
transforms.outbox.table.field.event.payload: payload
transforms.outbox.route.by.field: agregado_tipo
transforms.outbox.route.topic.replacement: ${routedByValue}.eventos
# ⚠️ Coste de CDC: un clúster de Kafka Connect que operar, un slot de replicación que
# VIGILAR (si el conector se para, el WAL crece hasta llenar el disco de la BD y
# tumbar el servicio) y permisos de replicación en producción.
# Empieza con polling; pasa a CDC cuando el volumen o la latencia lo justifiquen.
6.6 Inbox y deduplicación en el consumidor
La outbox garantiza at-least-once: el mensaje llega al menos una vez, y puede llegar repetido. El complemento obligatorio es hacer el consumidor idempotente, y la forma más directa es una tabla inbox de mensajes ya procesados.
// ─── Inbox: registro de mensajes procesados, en la MISMA transacción que el efecto ───
@Component
class PedidoConfirmadoListener {
private static final Logger log = LoggerFactory.getLogger(PedidoConfirmadoListener.class);
private final InboxRepository inbox;
private final CrearFacturaUseCase crearFactura;
@KafkaListener(topics = "pedidos.eventos", groupId = "facturacion")
@Transactional
public void escuchar(@Payload PedidoConfirmadoDto evento,
@Header("id") String idMensaje) {
// 1. ¿Ya lo procesamos? La PK de la tabla hace de candado distribuido.
if (!inbox.registrarSiEsNuevo(idMensaje, "facturacion")) {
log.debug("Mensaje {} ya procesado por facturacion: descartado", idMensaje);
return; // duplicado: salimos sin efecto secundario
}
// 2. Efecto de negocio, en la MISMA transacción que el registro del inbox.
// Si esto falla y hay rollback, el registro del inbox también desaparece
// y el mensaje se podrá reprocesar. Atomicidad de nuevo.
crearFactura.ejecutar(evento.aComando());
}
}
interface InboxRepository extends Repository<InboxEntity, String> {
/** @return 1 si se insertó (mensaje nuevo); 0 si ya existía. */
@Modifying
@Query(value = """
INSERT INTO mensajes_procesados (mensaje_id, consumidor, procesado_en)
VALUES (:id, :consumidor, now())
ON CONFLICT (mensaje_id, consumidor) DO NOTHING
""", nativeQuery = true)
int intentarRegistrar(@Param("id") String id, @Param("consumidor") String consumidor);
default boolean registrarSiEsNuevo(String id, String consumidor) {
return intentarRegistrar(id, consumidor) == 1;
}
}
UPSERT por clave de
negocio, o fijar un estado absoluto en lugar de incrementarlo;
(2) una restricción UNIQUE de negocio (UNIQUE(pedido_id) en la tabla de
facturas) que rechaza el duplicado en la base de datos;
(3) comprobar la versión del agregado y descartar eventos antiguos. Usa la tabla inbox
cuando el efecto es externo (enviar un email, llamar a un tercero) y no hay clave natural.
6.7 CQRS: separar el modelo de escritura del de lectura
CQRS (Command Query Responsibility Segregation) es simplemente esto: el modelo que escribe y el que lee no tienen por qué ser el mismo. Nada más. No implica dos bases de datos, ni event sourcing, ni mensajería; eso son variantes que se añaden cuando hacen falta.
CQRS — tres niveles de intensidad; elige el menor que resuelva tu problema
NIVEL 1 — separar clases (coste casi cero; hazlo siempre)
Escritura: agregado Pedido con invariantes → RepositorioPedidos
Lectura: PedidoVista (record plano) → ConsultasPedido (SQL a medida)
Misma BD, mismas tablas. Solo dejas de forzar al agregado a servir consultas.
NIVEL 2 — vistas materializadas en la MISMA base de datos (coste bajo)
┌──────────────┐ transacción ┌────────────────────────────────────────┐
│ Comandos │───────────────►│ tablas normalizadas (fuente de verdad) │
└──────────────┘ └──────────────┬─────────────────────────┘
│ trigger / evento / job
▼
┌──────────────┐ SELECT ┌────────────────────────────────────────┐
│ Consultas │◄───────────────│ vista materializada desnormalizada │
└──────────────┘ │ (0 joins, índices a medida) │
└────────────────────────────────────────┘
NIVEL 3 — almacenes separados (coste alto; solo con motivo demostrado)
┌───────────────┐ comandos ┌────────────┐ eventos ┌──────────────┐
│ API escritura │────────────►│ PostgreSQL │────────────►│ Proyector │
└───────────────┘ └────────────┘ [Kafka] └──────┬───────┘
▼
┌───────────────┐ consultas ┌──────────────────┐
│ API lectura │────────────────────────────────────►│ Elasticsearch / │
└───────────────┘ │ Redis / MongoDB │
└──────────────────┘
Ventajas: escalar lecturas y escrituras por separado, motor óptimo para cada uso.
Coste: consistencia eventual visible, proyectores que mantener, reconstrucción de
proyecciones y doble operación. NO empieces aquí.
| Usa CQRS cuando… | NO uses CQRS cuando… |
|---|---|
| La relación lecturas/escrituras es muy asimétrica (1000:1). | Es un CRUD normal con consultas simples. |
| Las consultas necesitan un modelo radicalmente distinto: búsqueda facetada, agregados, informes. | Puedes resolverlo con un índice o una consulta mejor escrita. |
| Las consultas cruzan varios agregados o servicios. | El equipo aún no domina el modelo de dominio. |
| Los joins de las consultas están matando al modelo transaccional. | El negocio exige leer siempre el último dato al instante. |
| Necesitas escalar lecturas de forma independiente. | «Porque es lo moderno». |
6.8 Event Sourcing: qué es y cuándo NO
En event sourcing no se guarda el estado actual, sino la secuencia inmutable de eventos que lo produjeron. El estado se obtiene reproduciéndolos. La analogía perfecta es la contabilidad por partida doble: no se borra un asiento, se emite otro que lo corrige.
ESTADO vs EVENTOS
MODELO CLÁSICO (solo el "ahora") EVENT SOURCING (toda la historia)
┌──────────────────────────┐ ┌──────────────────────────────────────┐
│ cuenta │ │ eventos (append-only, inmutable) │
│ id = 42 │ │ 1 CuentaAbierta {saldo: 0} │
│ saldo = 150 │ │ 2 DineroIngresado {200} │
│ version = 4 │ │ 3 DineroRetirado {80} │
└──────────────────────────┘ │ 4 DineroIngresado {30} │
¿Por qué el saldo es 150? └──────────────────────────────────────┘
NO SE SABE. Se perdió. saldo = 0 + 200 - 80 + 30 = 150 ✅
y además: cuándo, en qué orden y por qué.
SNAPSHOT (optimización): cada N eventos se guarda el estado calculado para no
reproducir 500.000 eventos al cargar el agregado.
[snapshot v1000: saldo=4210] + eventos 1001..1043 → estado actual
Lo que ganas
- Auditoría perfecta y gratuita: la historia completa es el propio modelo. Oro puro en banca, seguros y salud.
- Consultas temporales: «¿cuál era el saldo el 3 de marzo a las 10:00?» es trivial.
- Depuración excepcional: puedes reproducir el estado exacto que provocó un bug.
- Proyecciones nuevas sobre datos viejos: creas una vista nueva y la rellenas reproduciendo el histórico.
- Encaja de forma natural con CQRS y con la integración por eventos.
Lo que cuesta de verdad
- Versionado de eventos eterno: un evento de 2019 debe seguir siendo legible en 2030. No puedes «migrar y olvidar».
- Consultas difíciles: sin proyecciones no puedes hacer un simple
WHERE estado = 'X'. - GDPR y derecho al olvido: un log inmutable con datos personales es un problema legal serio; se resuelve con crypto-shredding, que hay que diseñar desde el día 1.
- Curva de aprendizaje muy alta: todo el equipo debe entenderlo, no solo quien lo introdujo.
- Herramientas y operación: EventStoreDB, Axon o una tabla propia; snapshots, reproyecciones y su monitorización.
6.9 Consultas que cruzan servicios: composición o vista materializada
"Necesito una pantalla con: pedido + datos del cliente + estado del envío"
OPCIÓN A — API COMPOSITION (el BFF hace el abanico en tiempo real)
┌──────┐ 1 petición ┌──────────┐──►Pedidos (30 ms) ┐
│ App │────────────────►│ BFF │──►Clientes (25 ms) ├ en PARALELO: 40 ms
└──────┘ └──────────┘──►Envíos (40 ms) ┘
✅ Siempre datos frescos; sin almacenamiento extra; simple de entender.
❌ Latencia = la del más lento; disponibilidad = producto de todas;
imposible filtrar/ordenar/paginar por campos de OTRO servicio
("dame los pedidos de clientes VIP ordenados por ciudad" → inviable).
OPCIÓN B — VISTA MATERIALIZADA (un proyector escucha eventos y mantiene una tabla)
Pedidos ──evento──┐
Clientes ─evento──┼──► PROYECTOR ──► tabla vista_pedido_completo
Envíos ───evento──┘ (pedido + cliente + envío, desnormalizado)
┌──────┐ 1 petición ┌──────────┐ 1 SELECT sin joins (3 ms)
│ App │────────────────►│ BFF │──────────────────────► vista
└──────┘ └──────────┘
✅ Latencia mínima y constante; disponible aunque los otros servicios caigan;
permite filtrar, ordenar y paginar por CUALQUIER campo.
❌ Consistencia eventual; hay que mantener el proyector; hay que poder
RECONSTRUIR la vista desde cero (guion de reproyección: obligatorio).
CÓMO ELEGIR:
· Pocas consultas, datos que deben ser frescos, 2-3 servicios → composición.
· Pantallas de listado con filtros, paginación y alto tráfico → vista materializada.
· Informes y analítica → ni una ni otra:
replica a un almacén analítico y olvídate de tocar los servicios.
7 · Kafka y mensajería
7.1 Conceptos: broker, topic, partición, offset, réplicas
Kafka no es una cola: es un log distribuido, particionado y replicado. Entender esa frase completa evita el 80 % de los errores que se cometen con él. Los mensajes no se «consumen» y desaparecen: se leen por posición, y siguen ahí hasta que caduquen.
ANATOMÍA DE UN TOPIC
TOPIC "pedidos.eventos" (retención: 7 días · replication.factor = 3)
Partición 0 ┌────┬────┬────┬────┬────┬────┬────┐
│ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │ 6 │◄── se escribe SIEMPRE al final
└────┴────┴────┴────┴────┴────┴────┘ (append-only)
▲ ▲
offset commiteado log end offset
del grupo "envios" LAG = 6 - 3 = 3 mensajes
Partición 1 ┌────┬────┬────┬────┐
│ 0 │ 1 │ 2 │ 3 │
└────┴────┴────┴────┘
Partición 2 ┌────┬────┬────┬────┬────┐
│ 0 │ 1 │ 2 │ 3 │ 4 │
└────┴────┴────┴────┴────┘
· El OFFSET es único y creciente DENTRO de cada partición, no en el topic.
· El ORDEN solo está garantizado DENTRO de una partición. Nunca entre particiones.
· La CLAVE decide la partición: particion = hash(clave) % numero_de_particiones.
→ misma clave = misma partición = ORDEN GARANTIZADO para esa clave.
→ clave null = reparto (sticky partitioning) = SIN garantía de orden.
RÉPLICAS Y TOLERANCIA A FALLOS (replication.factor = 3)
Broker 1 Broker 2 Broker 3
┌─────────┐ ┌─────────┐ ┌─────────┐
│ P0 LÍDER│◄─────►│ P0 segui│◄─────►│ P0 segui│ Escrituras y lecturas van al
│ P1 segui│ │ P1 LÍDER│ │ P1 segui│ LÍDER de cada partición.
│ P2 segui│ │ P2 segui│ │ P2 LÍDER│ Los seguidores replican.
└─────────┘ └─────────┘ └─────────┘
ISR (In-Sync Replicas) = réplicas al día con el líder.
Si un líder cae, se elige uno nuevo ENTRE LOS ISR → cero pérdida.
LA COMBINACIÓN CORRECTA PARA NO PERDER DATOS:
replication.factor = 3
min.insync.replicas = 2 (a nivel de topic)
acks = all (en el productor)
→ una escritura solo se confirma si está en al menos 2 réplicas.
→ se tolera la caída de 1 broker sin perder datos ni disponibilidad de escritura.
⚠️ min.insync.replicas=2 con RF=2 NO tolera ninguna caída: el topic deja de aceptar
escrituras en cuanto pierdes un broker. Error de configuración muy común.
RETENCIÓN vs COMPACTACIÓN
cleanup.policy=delete (por defecto): borra por tiempo (retention.ms, 7 días) o
por tamaño (retention.bytes). Es un LOG DE EVENTOS.
cleanup.policy=compact: conserva, para cada CLAVE, al menos el ÚLTIMO valor.
Es una TABLA de estado. Un valor null es una "lápida" (tombstone) que borra
la clave. Ideal para topics de configuración, catálogos o snapshots de estado.
KRaft: desde Kafka 3.3 el modo KRaft (sin ZooKeeper) es apto para producción, y en
Kafka 4.0 ZooKeeper se eliminó por completo. Los metadatos viven en un quórum de
controladores dentro del propio Kafka.
7.2 El productor
| Propiedad | Valor recomendado | Por qué |
|---|---|---|
acks | all | 0 = «dispara y olvida» (pérdida garantizada); 1 = solo el líder (pierdes si cae antes de replicar); all = confirmado por los ISR. |
enable.idempotence | true (por defecto desde 3.0) | El broker deduplica reenvíos usando un identificador de productor y un número de secuencia por partición. Elimina los duplicados del reintento y preserva el orden. |
max.in.flight.requests.per.connection | 5 o menos | Con idempotencia activada, hasta 5 mantiene el orden. Sin idempotencia, más de 1 puede reordenar mensajes al reintentar. |
delivery.timeout.ms | 120000 | Es el que manda de verdad: tiempo total para dar un mensaje por entregado, reintentos incluidos. |
linger.ms | 5–50 | Esperar unos milisegundos para agrupar mensajes en lotes multiplica el rendimiento. 0 minimiza latencia y desperdicia red. |
batch.size | 32768–131072 | Tamaño del lote por partición. |
compression.type | zstd (o lz4) | Se comprime el lote entero: con JSON es habitual reducir 5–10× el tráfico y el almacenamiento a cambio de poca CPU. |
transactional.id | Único y estable por instancia | Solo si necesitas transacciones: escritura atómica en varios topics más el commit de offsets. |
pedidoId): garantiza que todos los eventos de un mismo pedido se procesan en orden. Si
eliges una clave con poca variedad (pais, un tenantId con un cliente enorme,
tipoEvento), tendrás particiones calientes: una saturada mientras las demás
están vacías, y no podrás escalar por más consumidores que añadas.
7.3 El consumidor: grupos, offsets y garantías
CONSUMER GROUPS — el reparto de particiones
TOPIC con 4 particiones
┌──────┬──────┬──────┬──────┐
│ P0 │ P1 │ P2 │ P3 │
└──┬───┴──┬───┴──┬───┴──┬───┘
│ │ │ │
┌──▼──────▼───┐ │ │ GRUPO "facturacion" con 2 consumidores:
│ Consumidor A│ │ │ A lee P0 y P1; B lee P2 y P3.
└─────────────┘ │ │
┌────────────────▼──────▼───┐
│ Consumidor B │
└───────────────────────────┘
REGLA DE ORO: el paralelismo máximo de un grupo = NÚMERO DE PARTICIONES.
· 4 particiones y 6 consumidores → 2 consumidores IDLE, sin hacer nada.
· 4 particiones y 2 consumidores → cada uno lleva 2. Correcto.
· Añadir particiones es fácil; QUITARLAS es imposible. Dimensiona con holgura
(2-3× el paralelismo previsto) pero sin pasarte: cada partición cuesta
descriptores de fichero, memoria y tiempo de failover.
OTRO GRUPO ("analitica") recibe TODOS los mensajes de nuevo, con sus propios
offsets. Añadir un consumidor nuevo NO afecta a los existentes: esa es la magia.
REBALANCEO — cuando entra o sale un consumidor
· Estrategia moderna: CooperativeStickyAssignor → reasigna solo las particiones
necesarias, sin parar a todo el grupo (el "stop-the-world" del rebalanceo eager
era el gran dolor clásico).
· group.instance.id (static membership) evita rebalanceos en reinicios
planificados, como los despliegues rodantes.
· Si tardas más de max.poll.interval.ms (5 min por defecto) en procesar un lote,
el broker te da por muerto y rebalancea → duplicados y consumer lag. Solución:
reducir max.poll.records o mover el trabajo pesado a otro hilo.
| Garantía | Cómo se consigue | Riesgo | Cuándo usarla |
|---|---|---|---|
| At-most-once | Commit del offset antes de procesar. | Si falla el procesamiento, el mensaje se pierde. | Métricas, telemetría, logs: perder uno no importa. |
| At-least-once (el 95 % de los casos) | Commit después de procesar con éxito. | Duplicados si el proceso muere entre el efecto y el commit. | Todo lo demás. Exige consumidor idempotente. |
| «Exactly-once» | Transacciones de Kafka: sendOffsetsToTransaction con consumidor en isolation.level=read_committed. |
Complejidad, menor rendimiento y una garantía limitada. | Procesamiento Kafka → Kafka (Kafka Streams). Ver el matiz de abajo. |
7.4 Errores, reintentos y dead letter queue
GESTIÓN DE ERRORES EN EL CONSUMO — retry topics y DLT
┌──────────────────────┐
│ pedidos.eventos │
└──────────┬───────────┘
▼
┌─────────────┐ éxito
│ Consumidor │──────────► commit y siguiente
└──────┬──────┘
│ excepción
┌─────────────┴──────────────┐
│ │
¿es RECUPERABLE? ¿es NO recuperable?
(timeout, 503, BD caída) (deserialización, validación,
│ regla de negocio violada)
▼ ▼
┌──────────────────────────┐ ┌──────────────────────┐
│ pedidos.eventos-retry-0 │ │ pedidos.eventos-DLT │
│ (espera 1 s) │ │ directamente, │
└────────────┬─────────────┘ │ sin reintentar │
▼ └──────────────────────┘
┌──────────────────────────┐ Reintentar un mensaje que
│ pedidos.eventos-retry-1 │ SIEMPRE va a fallar es
│ (espera 2 s) │ quemar CPU y retrasar
└────────────┬─────────────┘ a los mensajes buenos.
▼
┌──────────────────────────┐
│ pedidos.eventos-retry-2 │
│ (espera 4 s) │
└────────────┬─────────────┘
▼
┌──────────────────────────┐
│ pedidos.eventos-DLT │◄── ALERTA: un mensaje en la DLT es un incidente,
│ (retención 30 días) │ no un dato. Alerta si DLT > 0 durante 5 min.
└──────────────────────────┘
⚠️ POR QUÉ NO SE REINTENTA "EN SITIO" (bloqueando la partición):
Kafka entrega en ORDEN dentro de una partición. Si te quedas 30 s reintentando
el mensaje 5, los mensajes 6..5000 de esa partición ESPERAN. Un solo mensaje
envenenado paraliza a todos los demás clientes de esa partición. Los retry
topics separados desbloquean la partición principal a costa de perder el orden
para los mensajes reintentados: es un intercambio consciente.
MENSAJE ENVENENADO (poison pill): un mensaje que hace fallar al consumidor SIEMPRE
(JSON corrupto, esquema incompatible). Sin ErrorHandlingDeserializer, el consumidor
entra en bucle infinito: falla, no commitea, vuelve a leer el mismo mensaje... para
siempre, con el lag creciendo. Incidente clásico de las 3 de la mañana.
7.5 Spring Kafka: configuración y código completos
spring:
application:
name: servicio-facturacion
kafka:
bootstrap-servers: ${KAFKA_BOOTSTRAP:localhost:9092}
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
delivery.timeout.ms: 120000
request.timeout.ms: 30000
linger.ms: 20
compression.type: zstd
spring.json.add.type.headers: false # ❗ no acoples el consumidor a TUS clases
consumer:
group-id: facturacion
auto-offset-reset: earliest # 'latest' se salta el histórico: cuidado
enable-auto-commit: false # SIEMPRE manual o gestionado por el contenedor
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
properties:
spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer
spring.json.trusted.packages: "com.ejemplo.eventos" # nunca "*": es un riesgo real
spring.json.value.default.type: com.ejemplo.eventos.PedidoConfirmadoDto
isolation.level: read_committed # no leer mensajes de transacciones abiertas
max.poll.records: 100 # menos registros = menos riesgo de expulsión
max.poll.interval.ms: 300000
session.timeout.ms: 45000
heartbeat.interval.ms: 3000
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor
listener:
ack-mode: record # commit tras cada registro procesado con éxito
concurrency: 3 # 3 hilos consumidores: NUNCA más que particiones
observation-enabled: true # trazas de Micrometer en producción y consumo
admin:
auto-create: false # los topics se crean con IaC, no por sorpresa
logging:
level:
org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: INFO # ver rebalanceos
// ─── Productor ───
@Component
class PublicadorPedidos {
private static final Logger log = LoggerFactory.getLogger(PublicadorPedidos.class);
private final KafkaTemplate<String, Object> kafka;
private final MeterRegistry metricas;
/** Envío asíncrono con manejo explícito del resultado: no ignores el CompletableFuture. */
void publicar(PedidoConfirmadoDto evento) {
var mensaje = MessageBuilder.withPayload(evento)
.setHeader(KafkaHeaders.TOPIC, "pedidos.eventos")
.setHeader(KafkaHeaders.KEY, evento.pedidoId().toString()) // orden por pedido
.setHeader("tipo", "pedidos.PedidoConfirmado.v1")
.setHeader("id", evento.eventoId().toString())
.build();
kafka.send(mensaje).whenComplete((resultado, error) -> {
if (error != null) {
log.error("Fallo publicando {}", evento.eventoId(), error);
metricas.counter("eventos.publicacion.fallida").increment();
} else {
var md = resultado.getRecordMetadata();
log.debug("Publicado en {}-{} offset {}", md.topic(), md.partition(), md.offset());
}
});
}
}
// ─── Consumidor con reintentos escalonados y DLT ───
@Component
class FacturacionListener {
private static final Logger log = LoggerFactory.getLogger(FacturacionListener.class);
private final CrearFacturaUseCase crearFactura;
private final InboxRepository inbox;
@RetryableTopic(
attempts = "4", // 1 original + 3 reintentos
backoff = @Backoff(delay = 1000, multiplier = 2.0, maxDelay = 10_000),
dltStrategy = DltStrategy.FAIL_ON_ERROR,
topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
exclude = { ReglaDeNegocioException.class, // deterministas: directos a la DLT
DeserializationException.class })
@KafkaListener(topics = "pedidos.eventos", groupId = "facturacion")
@Transactional
public void escuchar(@Payload PedidoConfirmadoDto evento,
@Header(KafkaHeaders.RECEIVED_KEY) String clave,
@Header(name = "id", required = false) String idMensaje,
@Header(KafkaHeaders.RECEIVED_PARTITION) int particion,
@Header(KafkaHeaders.OFFSET) long offset) {
log.debug("Recibido {} p{} offset {}", clave, particion, offset);
if (!inbox.registrarSiEsNuevo(idMensaje, "facturacion")) return; // idempotencia
crearFactura.ejecutar(evento.aComando());
}
/** Punto final: aquí llegan los mensajes que no se pudieron procesar. */
@DltHandler
public void enDlt(@Payload PedidoConfirmadoDto evento,
@Header(KafkaHeaders.ORIGINAL_TOPIC) String topicOriginal,
@Header(KafkaHeaders.EXCEPTION_MESSAGE) String error) {
log.error("MENSAJE EN DLT desde {}: pedido={} error={}",
topicOriginal, evento.pedidoId(), error);
alertas.enviar(Severidad.ALTA, "Evento en DLT: " + evento.pedidoId());
repositorioDlt.guardarParaRevision(evento, topicOriginal, error);
}
}
// ─── Manejador de errores global (alternativa a @RetryableTopic, con DLT única) ───
@Configuration
class KafkaErrorConfig {
@Bean
DefaultErrorHandler errorHandler(KafkaTemplate<String, Object> template) {
// Enruta el fallo al topic .DLT conservando la MISMA partición: útil para depurar
var recuperador = new DeadLetterPublishingRecoverer(template,
(registro, ex) -> new TopicPartition(registro.topic() + ".DLT", registro.partition()));
var backoff = new ExponentialBackOffWithMaxRetries(3);
backoff.setInitialInterval(500L);
backoff.setMultiplier(2.0);
backoff.setMaxInterval(10_000L);
var handler = new DefaultErrorHandler(recuperador, backoff);
// Errores deterministas: no tiene sentido reintentarlos, van directos a la DLT
handler.addNotRetryableExceptions(
DeserializationException.class,
MessageConversionException.class,
MethodArgumentNotValidException.class,
ReglaDeNegocioException.class);
handler.setRetryListeners((registro, ex, intento) ->
LoggerFactory.getLogger("kafka.retry")
.warn("Reintento {} de {}-{}@{}: {}", intento, registro.topic(),
registro.partition(), registro.offset(), ex.getMessage()));
return handler;
}
/** Topics como código: replicación y retención explícitas, nada de auto-creación. */
@Bean
KafkaAdmin.NewTopics topics() {
return new KafkaAdmin.NewTopics(
TopicBuilder.name("pedidos.eventos")
.partitions(12).replicas(3)
.config(TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, "2")
.config(TopicConfig.RETENTION_MS_CONFIG,
String.valueOf(Duration.ofDays(7).toMillis()))
.build(),
TopicBuilder.name("pedidos.eventos.DLT")
.partitions(12).replicas(3)
.config(TopicConfig.RETENTION_MS_CONFIG,
String.valueOf(Duration.ofDays(30).toMillis()))
.build(),
TopicBuilder.name("catalogo.productos") // topic de ESTADO: compactado
.partitions(6).replicas(3)
.config(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT)
.build());
}
}
// ─── Test de integración real con Testcontainers (Spring Boot 3.1+ y @ServiceConnection) ───
@SpringBootTest
@Testcontainers
class FacturacionKafkaIT {
@Container
@ServiceConnection // configura spring.kafka.bootstrap-servers solo
static final ConfluentKafkaContainer KAFKA =
new ConfluentKafkaContainer("confluentinc/cp-kafka:7.6.1");
// Alternativa: new KafkaContainer(DockerImageName.parse("apache/kafka:3.8.0")) — KRaft nativo
@Container
@ServiceConnection
static final PostgreSQLContainer<?> POSTGRES =
new PostgreSQLContainer<>("postgres:16-alpine");
@Autowired KafkaTemplate<String, Object> kafka;
@Autowired FacturaRepository facturas;
@Test
void crea_una_factura_al_recibir_un_pedido_confirmado() {
var evento = unPedidoConfirmado();
kafka.send("pedidos.eventos", evento.pedidoId().toString(), evento);
await().atMost(Duration.ofSeconds(10))
.untilAsserted(() -> assertThat(facturas.findByPedidoId(evento.pedidoId()))
.isPresent()
.get().extracting(Factura::total).isEqualTo(evento.total()));
}
@Test
void no_duplica_la_factura_si_el_evento_llega_dos_veces() {
var evento = unPedidoConfirmado();
kafka.send("pedidos.eventos", evento.pedidoId().toString(), evento);
kafka.send("pedidos.eventos", evento.pedidoId().toString(), evento); // duplicado exacto
await().during(Duration.ofSeconds(3))
.atMost(Duration.ofSeconds(10))
.untilAsserted(() -> assertThat(facturas.countByPedidoId(evento.pedidoId()))
.isEqualTo(1)); // la idempotencia funciona
}
@Test
void envia_a_la_DLT_un_mensaje_que_no_se_puede_procesar() {
kafka.send("pedidos.eventos", "malo", new PedidoConfirmadoDto(null, null, null, null, null));
var consumidor = consumidorDe("pedidos.eventos.DLT");
var registros = KafkaTestUtils.getRecords(consumidor, Duration.ofSeconds(15));
assertThat(registros.count()).isEqualTo(1);
}
}
7.6 Diseño de eventos: el contrato que más dura
DOS ESTILOS DE EVENTO — elige conscientemente, no por accidente
A) NOTIFICACIÓN (event notification) — "algo pasó, ve a buscar los detalles"
{ "tipo": "PedidoConfirmado", "pedidoId": "abc-123", "ocurridoEn": "..." }
✅ Mensaje diminuto; el productor no expone su modelo interno.
❌ El consumidor DEBE llamar de vuelta → acoplamiento en tiempo de ejecución,
carga extra en el productor y posible "tormenta de callbacks" (N eventos =
N llamadas de vuelta simultáneas).
❌ Condición de carrera: al llamar de vuelta puedes leer un estado MÁS NUEVO
que el del evento y procesar algo inconsistente.
B) TRANSFERENCIA DE ESTADO (event-carried state transfer) — el evento se basta solo
{ "tipo": "PedidoConfirmado", "pedidoId": "abc-123", "clienteId": "...",
"total": {"importe": 12050, "moneda": "EUR"},
"lineas": [{"sku":"A1","unidades":2,"subtotal":6025}],
"ocurridoEn": "2026-03-01T10:00:00Z", "version": 4 }
✅ El consumidor es AUTÓNOMO: puede procesar aunque el productor esté caído.
✅ Sin llamadas de vuelta ni carreras: el evento es una foto coherente.
❌ Mensajes más grandes; el contrato es más amplio y hay que versionarlo bien.
❌ Cuidado con datos personales: el evento viaja y se retiene 7 días o más.
RECOMENDACIÓN: por defecto, B. Incluye lo que los consumidores necesitan para hacer
su trabajo, ni un campo más. Si un consumidor necesita 40 campos del productor,
revisa las fronteras: probablemente están mal.
ANATOMÍA DE UN BUEN EVENTO
┌───────────────────────────────────────────────────────────────────┐
│ CABECERAS (metadatos, fuera del payload) │
│ id UUID único del evento → deduplicación │
│ tipo "pedidos.PedidoConfirmado.v1" │
│ traceparent contexto W3C → traza distribuida │
│ ocurridoEn instante del HECHO (no de la publicación) │
│ productor "servicio-pedidos@2.4.1" │
│ tenant multi-tenencia │
├───────────────────────────────────────────────────────────────────┤
│ CLAVE: pedidoId → partición → ORDEN por agregado │
├───────────────────────────────────────────────────────────────────┤
│ PAYLOAD: solo datos de NEGOCIO, con tipos explícitos │
│ · importes en céntimos (long) o con moneda explícita │
│ · fechas en ISO-8601 UTC │
│ · enums como texto, nunca como ordinal │
│ · version del agregado → permite descartar eventos antiguos │
└───────────────────────────────────────────────────────────────────┘
NOMENCLATURA DE TOPICS (elige una y documéntala)
dominio.agregado.tipo → pedidos.pedido.eventos
contexto.evento.vN → ventas.pedido-confirmado.v1
Evita: "eventos", "datos", "temp", "test-juan", mayúsculas mezcladas.
VERSIONADO DE EVENTOS
1. Cambios COMPATIBLES (añadir campo opcional): misma versión, sin drama.
2. Cambios INCOMPATIBLES: publica v2 EN PARALELO a v1 durante un tiempo (doble
publicación), migra consumidores uno a uno y retira v1 cuando las métricas de
consumo de v1 estén a cero. NUNCA rompas y avises después.
3. Registra la versión en el nombre del tipo y/o del topic. Nunca solo en el
payload: el consumidor debe poder decidir ANTES de deserializar.
7.7 Kafka Streams (mención)
Kafka Streams es una librería (no un clúster aparte) que se embebe en tu aplicación Java
para procesar topics como flujos: filtrar, transformar, agregar por ventanas de tiempo, unir dos flujos o
un flujo con una tabla (KTable), manteniendo el estado en un almacén local (RocksDB)
respaldado por un topic de changelog. Es la forma natural de construir proyecciones y agregados
en tiempo real, y es el único sitio donde el exactly-once de Kafka funciona de punta a punta
(processing.guarantee=exactly_once_v2).
// Ejemplo mínimo: total facturado por cliente en ventanas de 1 hora
@Bean
KStream<String, PedidoConfirmadoDto> totalPorCliente(StreamsBuilder builder) {
KStream<String, PedidoConfirmadoDto> pedidos =
builder.stream("pedidos.eventos", Consumed.with(Serdes.String(), pedidoSerde()));
pedidos.filter((k, v) -> v.total().importe() > 0)
.groupBy((k, v) -> v.clienteId().toString(), Grouped.with(Serdes.String(), pedidoSerde()))
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofHours(1), Duration.ofMinutes(5)))
.aggregate(() -> 0L,
(clave, pedido, acumulado) -> acumulado + pedido.total().importe(),
Materialized.with(Serdes.String(), Serdes.Long()))
.toStream()
.map((ventana, total) -> KeyValue.pair(ventana.key(),
new ResumenCliente(ventana.key(), total, ventana.window().startTime())))
.to("clientes.resumen-horario", Produced.with(Serdes.String(), resumenSerde()));
return pedidos;
}
7.8 Kafka, RabbitMQ, SQS/SNS y Pulsar
| Criterio | Kafka | RabbitMQ | SQS / SNS | Pulsar |
|---|---|---|---|---|
| Modelo | Log distribuido particionado | Broker de colas con enrutamiento (AMQP) | Cola y pub-sub gestionados | Log y colas, con almacenamiento separado (BookKeeper) |
| Rendimiento | Millones de mensajes/s | Decenas de miles/s por cola | Alto, elástico y gestionado | Comparable a Kafka |
| Retención y relectura | Sí: días o años; se puede reprocesar todo | No: al consumir, desaparece | SQS 14 días máx.; sin relectura | Sí, con almacenamiento por niveles (S3) |
| Orden | Por partición | Por cola (con un consumidor) | Solo en colas FIFO, con menor caudal | Por partición o por clave |
| Enrutamiento | Simple: topic y clave | Muy rico: direct, topic, fanout, headers | Básico (SNS a varios destinos) | Rico |
| Entrega retrasada | No nativa (se emula con retry topics) | Sí: plugin de delayed exchange, TTL con DLX | Sí, hasta 15 min | Sí, nativa |
| Multi-tenencia y geo | Con MirrorMaker | Federación y shovel | Por región | Nativa: tenants, namespaces, geo-replicación |
| Operación | Compleja (mejor con KRaft o gestionado) | Sencilla | Cero: es un servicio | La más compleja (broker más BookKeeper) |
| Úsalo para… | Eventos de negocio, streaming, integración a gran escala, reprocesado histórico | Colas de trabajo, RPC asíncrono, enrutamiento complejo, prioridades | Todo en AWS cuando no quieres operar nada | Multi-tenencia y geo fuertes; menos ecosistema |
| Evítalo si… | Solo necesitas una cola de tareas: es un cañón para una mosca | Necesitas relectura o volúmenes enormes | Necesitas orden global o latencia muy baja | No tienes equipo de plataforma |
RABBITMQ — el modelo de enrutamiento que Kafka NO tiene
Productor ──► EXCHANGE ──(binding con routing key)──► QUEUE ──► Consumidor
TIPOS DE EXCHANGE
· direct : routing key EXACTA "pedido.creado" → cola pedidos-creados
· topic : con comodines "pedido.*" → todas las de pedido
(* = una palabra, # = cero o más) "pedido.#.urgente"
· fanout : a TODAS las colas ligadas (ignora la routing key) = pub/sub
· headers : según cabeceras, no según routing key
┌──────────┐ "pedido.creado.es" ┌─────────────────┐
│ Productor│──────────────────────►│ EXCHANGE topic │
└──────────┘ │ "ventas" │
└───┬────┬────┬───┘
binding "pedido.#" │ │ │ binding "#.es"
┌───────────────────────┘ │ └──────────────────┐
▼ ▼ ▼
┌──────────────┐ ┌──────────────────┐ ┌───────────────┐
│ q.pedidos │ │ q.auditoria │ │ q.espana │
└──────┬───────┘ └──────────────────┘ └───────────────┘
│ nack sin requeue tras N intentos
▼
┌──────────────┐ x-dead-letter-exchange
│ q.pedidos.dlq│
└──────────────┘
// Spring AMQP: configuración típica de RabbitMQ con DLQ y reintentos
@Configuration
class RabbitConfig {
@Bean TopicExchange ventas() { return new TopicExchange("ventas", true, false); }
@Bean Queue pedidos() {
return QueueBuilder.durable("q.pedidos")
.withArgument("x-dead-letter-exchange", "ventas.dlx")
.withArgument("x-dead-letter-routing-key", "pedidos.fallidos")
.withArgument("x-message-ttl", 600_000) // 10 min
.withArgument("x-queue-type", "quorum") // replicada y tolerante a fallos
.build();
}
@Bean Binding bindPedidos(Queue pedidos, TopicExchange ventas) {
return BindingBuilder.bind(pedidos).to(ventas).with("pedido.#");
}
@Bean DirectExchange dlx() { return new DirectExchange("ventas.dlx", true, false); }
@Bean Queue dlq() { return QueueBuilder.durable("q.pedidos.dlq").build(); }
@Bean Binding bindDlq(Queue dlq, DirectExchange dlx) {
return BindingBuilder.bind(dlq).to(dlx).with("pedidos.fallidos");
}
}
spring:
rabbitmq:
host: rabbit
publisher-confirm-type: correlated # confirmación real de que el broker lo aceptó
publisher-returns: true # avisa si un mensaje no llega a ninguna cola
listener:
simple:
acknowledge-mode: manual
prefetch: 20 # ❗ el valor por defecto (250) provoca reparto injusto
concurrency: 3
max-concurrency: 10
default-requeue-rejected: false # sin esto, un mensaje malo se reencola ETERNAMENTE
retry:
enabled: true
max-attempts: 3
initial-interval: 1s
multiplier: 2
Checklist — datos distribuidos y Kafka
8 · Spring Cloud y el ecosistema de plataforma
Spring Cloud es un conjunto de proyectos que resuelven los problemas transversales de una arquitectura distribuida. Aviso importante antes de empezar: en 2026, si despliegas en Kubernetes, buena parte de Spring Cloud es redundante. El descubrimiento de servicios lo hace el DNS del clúster, la configuración la hacen ConfigMaps y Secrets, y el balanceo lo hace el Service. Usa Spring Cloud donde aporte, no por costumbre.
8.1 API Gateway con Spring Cloud Gateway
El gateway es la puerta única de entrada desde el exterior. Su valor no es enrutar (eso lo hace un Ingress), sino concentrar en un solo sitio lo que no debe repetirse en N servicios: autenticación, límites de tasa, cabeceras de seguridad, CORS, agregación y observabilidad de borde.
TOPOLOGÍA CON GATEWAY
Internet
│
▼
┌──────────┐ TLS, WAF, protección DDoS
│ CDN │
└────┬─────┘
▼
┌──────────────────────────────────────────────────────────────┐
│ API GATEWAY (Spring Cloud Gateway) │
│ · Autenticación: valida el JWT UNA vez │
│ · Autorización gruesa por ruta │
│ · Rate limiting por cliente (Redis + token bucket) │
│ · Enrutamiento y reescritura de rutas │
│ · Cabeceras: X-Request-Id, traceparent, CORS, HSTS │
│ · Circuit breaker por ruta y degradación │
│ · Métricas y trazas del borde │
└───┬──────────────┬──────────────┬──────────────┬─────────────┘
▼ ▼ ▼ ▼
┌────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐
│Pedidos │ │ Pagos │ │ Catálogo │ │ Envíos │
└────────┘ └──────────┘ └──────────┘ └──────────┘
Confían en las cabeceras firmadas del gateway,
PERO validan de nuevo lo crítico (defensa en profundidad)
⚠️ RIESGOS DEL GATEWAY
· Punto único de fallo → mínimo 3 réplicas y despliegue sin caída.
· Cuello de botella → todo el tráfico pasa por aquí: mide su latencia añadida.
· "Gateway obeso": si empieza a tener lógica de NEGOCIO, has creado un ESB.
El gateway enruta y protege; no decide descuentos.
spring:
cloud:
gateway:
default-filters:
- AddResponseHeader=X-Content-Type-Options, nosniff
- name: Retry
args:
retries: 2
statuses: BAD_GATEWAY,SERVICE_UNAVAILABLE
methods: GET # ❗ solo métodos idempotentes
backoff: { firstBackoff: 50ms, maxBackoff: 500ms, factor: 2, basedOnPreviousValue: false }
routes:
- id: pedidos
uri: lb://servicio-pedidos # lb:// = balanceo con Spring Cloud LoadBalancer
predicates:
- Path=/api/v1/pedidos/**
- Method=GET,POST,PUT,DELETE
filters:
- name: CircuitBreaker
args:
name: cbPedidos
fallbackUri: forward:/fallback/pedidos
- name: RequestRateLimiter
args:
redis-rate-limiter.replenishRate: 100 # tokens por segundo
redis-rate-limiter.burstCapacity: 200 # capacidad del cubo
redis-rate-limiter.requestedTokens: 1
key-resolver: "#{@resolutorPorUsuario}" # por usuario, NO global
- id: catalogo
uri: lb://servicio-catalogo
predicates:
- Path=/api/v1/catalogo/**
filters:
- name: LocalResponseCache # caché en el borde
args: { timeToLive: 60s, size: 50MB }
- id: legado
uri: http://monolito.interno:8080
predicates:
- Path=/api/v1/informes/**
filters:
- RewritePath=/api/v1/informes/(?<resto>.*), /legacy/reports/${resto}
httpclient:
connect-timeout: 500
response-timeout: 5s
pool: { max-connections: 500, type: elastic }
globalcors:
cors-configurations:
'[/**]':
allowedOriginPatterns: "https://*.ejemplo.com"
allowedMethods: [GET, POST, PUT, DELETE, OPTIONS]
allowedHeaders: "*"
allowCredentials: true
maxAge: 3600
// Filtro global: propagación de contexto y correlación en el borde
@Component
class ContextoGlobalFilter implements GlobalFilter, Ordered {
static final String CABECERA_PETICION = "X-Request-Id";
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
String requestId = Optional
.ofNullable(exchange.getRequest().getHeaders().getFirst(CABECERA_PETICION))
.orElseGet(() -> UUID.randomUUID().toString());
ServerHttpRequest peticion = exchange.getRequest().mutate()
.header(CABECERA_PETICION, requestId)
.header("X-Gateway-Received-At", Instant.now().toString())
.build();
exchange.getResponse().getHeaders().add(CABECERA_PETICION, requestId);
return chain.filter(exchange.mutate().request(peticion).build());
}
@Override public int getOrder() { return Ordered.HIGHEST_PRECEDENCE; }
}
/** Limitar por usuario autenticado; si es anónimo, por IP. Nunca una sola cubeta global. */
@Bean
KeyResolver resolutorPorUsuario() {
return exchange -> ReactiveSecurityContextHolder.getContext()
.map(ctx -> ctx.getAuthentication().getName())
.defaultIfEmpty(Optional.ofNullable(exchange.getRequest().getRemoteAddress())
.map(a -> a.getAddress().getHostAddress()).orElse("anonimo"));
}
@RestController
class FallbackController {
@RequestMapping("/fallback/pedidos")
ResponseEntity<ProblemDetail> pedidos() {
var pd = ProblemDetail.forStatusAndDetail(HttpStatus.SERVICE_UNAVAILABLE,
"El servicio de pedidos no está disponible. Inténtalo en unos minutos.");
return ResponseEntity.status(HttpStatus.SERVICE_UNAVAILABLE)
.header("Retry-After", "30").body(pd);
}
}
8.2 Descubrimiento de servicios: Eureka frente a DNS de Kubernetes
EUREKA (descubrimiento del lado del cliente) KUBERNETES (DNS + Service)
┌──────────┐ 1. se registra ┌────────┐ ┌──────────┐
│ Pedidos │──────────────────►│ EUREKA │ │ Pedidos │
│ │◄──2. descarga ────│ SERVER │ │ │
│ caché │ el registro └────────┘ └────┬─────┘
│ local │ │ DNS: servicio-pagos
└────┬─────┘ │ .produccion
│ 3. elige instancia y llama ▼ .svc.cluster.local
▼ (balanceo EN EL CLIENTE) ┌───────────┐
┌──────────┐ │ Service │ (IP virtual estable)
│ Pagos │ └─────┬─────┘
│ :8081 │ │ kube-proxy / iptables
└──────────┘ ┌────┴────┬────────┐
▼ ▼ ▼
Pros: sin infraestructura de red; pod1 pod2 pod3
metadatos ricos; funciona en VM.
Contras: OTRO servicio que operar (y en HA); Pros: CERO código y cero
latencia de propagación de 30-90 s (el dependencias; lo gestiona la
cliente puede llamar a instancias muertas); plataforma; healthchecks y
solo para clientes Java/Spring. readiness integrados.
Contras: atado a K8s; balanceo
L4 (problemático con gRPC/HTTP2).
http://servicio-pagos. Reserva Eureka (o Consul) para entornos híbridos,
máquinas virtuales o migraciones en las que aún no hay orquestador. Y si necesitas balanceo L7 real
(gRPC, reintentos, mTLS), la respuesta no es Eureka: es una malla de servicios.
8.3 Configuración centralizada
| Spring Cloud Config Server | ConfigMap y Secret de Kubernetes | |
|---|---|---|
| Origen | Repositorio Git (auditable, con historial y PR) | Manifiestos YAML, idealmente en Git con GitOps |
| Recarga en caliente | Sí, con @RefreshScope y /actuator/refresh o Spring Cloud Bus | Los ficheros montados se actualizan solos; las variables de entorno no (hay que reiniciar el pod) |
| Secretos | Cifrado propio o integración con Vault | Secret es solo base64: usa Sealed Secrets, External Secrets o Vault |
| Coste operativo | Un servicio más (y en alta disponibilidad) | Cero: viene con la plataforma |
| Cuándo | Fuera de Kubernetes, o si necesitas refresco sin reiniciar | Por defecto en Kubernetes |
@RefreshScope recrea los beans afectados. Si
cambias la URL de la base de datos en caliente, puedes dejar conexiones huérfanas o estados a medias. Es
excelente para feature flags, umbrales y niveles de log; es arriesgado para infraestructura. La
alternativa segura y trazable es cambiar la configuración y redesplegar: en Kubernetes
eso cuesta 30 segundos y queda registrado.
8.4 Balanceo en el cliente y malla de servicios
TRES SITIOS DONDE PUEDE VIVIR LA LÓGICA DE RED
A) EN LA APLICACIÓN (Spring Cloud LoadBalancer, Resilience4j)
┌──────────────────────────────┐
│ Tu servicio │ ✅ Control total, depuración fácil.
│ ├─ balanceo │ ❌ Cada lenguaje reimplementa lo mismo;
│ ├─ reintentos │ actualizar una política = redesplegar
│ ├─ circuit breaker │ todos los servicios.
│ └─ mTLS │
└──────────────────────────────┘
B) EN UN SIDECAR (malla de servicios: Istio, Linkerd)
┌───────────────────────────────────────────┐
│ Pod │ ✅ Políticas uniformes sin tocar
│ ┌──────────┐ ┌──────────────────┐ │ código; mTLS automático;
│ │ Tu │◄────►│ Sidecar (Envoy) │◄──┼──► métricas y trazas gratis;
│ │ servicio │ │ · mTLS │ │ canary y espejo de tráfico.
│ └──────────┘ │ · reintentos │ │ ❌ Complejidad operativa alta;
│ │ · circuit break │ │ +1-3 ms por salto; consumo de
│ │ · métricas │ │ recursos; depuración más dura.
│ └──────────────────┘ │ Linkerd es bastante más simple
└───────────────────────────────────────────┘ que Istio si no necesitas todo.
C) EN LA INFRAESTRUCTURA (Ingress, balanceador cloud)
✅ Simple. ❌ Solo en el borde: no cubre el tráfico entre servicios.
CRITERIO: menos de 10 servicios y todos en Java → (A) es suficiente.
Muchos servicios, varios lenguajes, requisito de mTLS y despliegues progresivos → (B).
8.5 Clientes declarativos: OpenFeign e interfaces HTTP
// OPCIÓN 1 — HTTP Interfaces: NATIVO en Spring Framework 6 / Boot 3, sin dependencias extra.
// Es la opción recomendada para proyectos nuevos.
public interface InventarioApi {
@GetExchange("/api/stock/{sku}")
StockDto consultar(@PathVariable String sku);
@PostExchange("/api/stock/reservas")
ReservaDto reservar(@RequestBody SolicitudReserva solicitud,
@RequestHeader("Idempotency-Key") String clave);
}
@Configuration
class ClientesHttpConfig {
@Bean
InventarioApi inventarioApi(RestClient.Builder builder,
ObservationRegistry observaciones) {
RestClient cliente = builder
.baseUrl("http://servicio-inventario") // resuelto por DNS de K8s
.requestFactory(factoriaConTimeouts())
.defaultStatusHandler(HttpStatusCode::is5xxServerError,
(req, res) -> { throw new DependenciaCaidaException("inventario"); })
.observationRegistry(observaciones) // trazas y métricas automáticas
.build();
return HttpServiceProxyFactory
.builderFor(RestClientAdapter.create(cliente))
.build()
.createClient(InventarioApi.class);
}
}
// OPCIÓN 2 — OpenFeign: útil si ya lo usas o quieres integración directa con Eureka
@FeignClient(name = "servicio-inventario", configuration = FeignConfig.class,
fallbackFactory = InventarioFallbackFactory.class)
public interface InventarioFeignClient {
@GetMapping("/api/stock/{sku}")
StockDto consultar(@PathVariable("sku") String sku);
}
@Configuration
class FeignConfig {
@Bean Request.Options opciones() {
return new Request.Options(500, TimeUnit.MILLISECONDS, // conexión
2000, TimeUnit.MILLISECONDS, // lectura
true); // seguir redirecciones
}
@Bean Logger.Level nivel() { return Logger.Level.BASIC; } // FULL loguea cuerpos: cuidado
}
8.6 Propagación de contexto entre servicios
Hay información que debe viajar con la petición en todos los saltos: el identificador de traza, el tenant, el usuario, el idioma y el deadline restante. Perderla en un solo salto rompe la observabilidad y, a veces, la seguridad.
/**
* Interceptor que propaga el contexto en llamadas salientes.
* La traza (traceparent) la propaga Micrometer Tracing automáticamente si usas
* RestClient/WebClient/RestTemplate creados desde el builder inyectado por Spring.
* Lo que NUNCA se propaga solo es el contexto de NEGOCIO: eso es cosa tuya.
*/
@Component
class PropagacionContextoInterceptor implements ClientHttpRequestInterceptor {
@Override
public ClientHttpResponse intercept(HttpRequest peticion, byte[] cuerpo,
ClientHttpRequestExecution ejecucion) throws IOException {
ContextoPeticion ctx = ContextoPeticionHolder.actual();
if (ctx != null) {
peticion.getHeaders().add("X-Tenant-Id", ctx.tenantId());
peticion.getHeaders().add("X-Request-Id", ctx.requestId());
peticion.getHeaders().add("Accept-Language", ctx.idioma());
// Deadline restante: el llamado sabe cuánto tiempo le queda de verdad
long restanteMs = ctx.deadline().toEpochMilli() - System.currentTimeMillis();
peticion.getHeaders().add("X-Deadline-Ms", String.valueOf(Math.max(0, restanteMs)));
}
return ejecucion.execute(peticion, cuerpo);
}
}
/**
* ⚠️ CON VIRTUAL THREADS Y EJECUCIÓN ASÍNCRONA, ThreadLocal NO BASTA.
* El contexto se pierde al saltar de hilo. Soluciones:
* · Usar ScopedValue (Java 21+, en preview) en lugar de ThreadLocal.
* · io.micrometer:context-propagation, que Spring Boot 3 integra.
* · Envolver los ejecutores con ContextExecutorService.
*/
@Bean
TaskDecorator decoradorDeContexto() {
return tarea -> {
ContextoPeticion ctx = ContextoPeticionHolder.actual();
Map<String, String> mdc = MDC.getCopyOfContextMap();
return () -> {
ContextoPeticionHolder.establecer(ctx);
if (mdc != null) MDC.setContextMap(mdc);
try { tarea.run(); }
finally { ContextoPeticionHolder.limpiar(); MDC.clear(); }
};
};
}
8.7 Despliegue independiente y feature flags
Si tus servicios no se pueden desplegar por separado, no tienes microservicios. La técnica que lo hace posible es separar el despliegue de la activación: despliegas código nuevo apagado y lo enciendes cuando quieres, sin volver a desplegar.
DESPLIEGUE ≠ ACTIVACIÓN (release)
Semana 1 ──► desplegar código nuevo con la bandera APAGADA (0 % usuarios)
Semana 1 ──► activar para el equipo interno (allowlist)
Semana 2 ──► activar para el 1 % → medir errores y latencia
Semana 2 ──► 10 % → 50 % → 100 %
Semana 4 ──► BORRAR la bandera y el código viejo ← el paso que todos olvidan
✅ El rollback es un cambio de configuración (segundos), no un despliegue (minutos).
✅ Permite desplegar en horario laboral, con el equipo despierto.
❌ DEUDA: cada bandera es un if permanente. 40 banderas = 2^40 combinaciones que
nadie ha probado. Pon FECHA DE CADUCIDAD a cada una y hazla fallar en CI.
// Feature flags con Spring: desde una propiedad refrescable o desde un servicio dedicado
@Service
class CalculadoraPrecios {
private final Banderas banderas;
private final MotorPreciosV1 v1;
private final MotorPreciosV2 v2;
Dinero calcular(Pedido pedido, ClienteId cliente) {
if (banderas.activa("precios.motor-v2", cliente.valor().toString())) {
Dinero nuevo = v2.calcular(pedido);
// Ejecución en la sombra: comparamos sin afectar al usuario
if (banderas.activa("precios.comparar-motores")) {
comparar(nuevo, v1.calcular(pedido), pedido);
}
return nuevo;
}
return v1.calcular(pedido);
}
}
9 · Observabilidad
Monitorización es saber si el sistema funciona; responde a preguntas que ya sabías que ibas a hacer. Observabilidad es poder averiguar por qué no funciona, incluyendo preguntas que nadie anticipó. En un monolito puedes sobrevivir con lo primero. En un sistema distribuido, sin lo segundo estás depurando a ciegas.
9.1 Los tres pilares y cómo se conectan
LOS TRES PILARES — su valor está en la CORRELACIÓN, no en cada uno por separado
┌────────────────────┬────────────────────┬────────────────────┐
│ MÉTRICAS │ TRAZAS │ LOGS │
├────────────────────┼────────────────────┼────────────────────┤
│ Números agregados │ Camino de UNA │ Detalle textual de │
│ en el tiempo │ petición por todo │ un evento concreto │
│ │ el sistema │ │
│ "¿QUÉ pasa?" │ "¿DÓNDE pasa?" │ "¿POR QUÉ pasa?" │
│ │ │ │
│ Barato, retención │ Coste medio, se │ CARO a escala, el │
│ larga (meses) │ muestrea │ mayor gasto oculto │
│ Cardinalidad LIMI- │ Contexto completo │ Cardinalidad libre │
│ TADA (¡crítico!) │ de una petición │ │
└────────────────────┴────────────────────┴────────────────────┘
FLUJO DE UNA INVESTIGACIÓN REAL (así se usan de verdad)
1. ALERTA: "el p99 de /checkout ha pasado de 200 ms a 3 s" ← MÉTRICA
│
2. PANEL: ¿desde cuándo? ¿qué servicio? ¿todos los pods? ← MÉTRICAS
│
3. TRAZA: abrir una traza lenta de ejemplo → ← TRAZA
gateway 5 ms → pedidos 2.980 ms → inventario 2.950 ms
¡el tiempo está en inventario!
│
4. LOGS: filtrar por traceId=4bf92f... en inventario → ← LOGS
"connection pool exhausted, waited 2900ms"
│
5. CAUSA: una consulta sin índice retiene las conexiones.
⚠️ Sin el traceId en los logs, el paso 4 es imposible y la investigación pasa de
5 minutos a 3 horas. Es la inversión con mejor retorno de todo el módulo.
9.2 Logs estructurados y correlación
| ❌ Mal | ✅ Bien |
|---|---|
Texto libre: log.info("Procesando pedido " + id) | Estructurado con parámetros: log.info("Pedido procesado", kv("pedidoId", id)) |
Sin traceId: imposible correlacionar | traceId y spanId en el MDC de cada línea |
Concatenación con +: construye el String aunque el nivel esté apagado | Marcadores {}: SLF4J solo formatea si va a emitir |
| Todo a nivel INFO, o todo a DEBUG en producción | Niveles con criterio; DEBUG activable por paquete sin reiniciar |
| Loguear el DNI, el email, el token o la tarjeta | Enmascarar o no loguear datos personales ni secretos |
log.error("error", e) y además relanzar | Loguear o relanzar, no las dos cosas (evita el log duplicado) |
| Loguear dentro de un bucle de 10.000 elementos | Loguear el resumen: "procesados {} elementos en {} ms" |
<!-- logback-spring.xml — JSON en producción, legible en local -->
<configuration>
<springProperty scope="context" name="app" source="spring.application.name"/>
<springProfile name="local">
<appender name="CONSOLA" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} %highlight(%-5level) [%X{traceId:-},%X{spanId:-}] %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<root level="INFO"><appender-ref ref="CONSOLA"/></root>
</springProfile>
<springProfile name="produccion">
<!-- Una línea = un objeto JSON. Nunca multilínea: los recolectores lo parten mal. -->
<appender name="JSON" class="ch.qos.logback.core.ConsoleAppender">
<encoder class="net.logstash.logback.encoder.LogstashEncoder">
<includeMdcKeyName>traceId</includeMdcKeyName>
<includeMdcKeyName>spanId</includeMdcKeyName>
<includeMdcKeyName>tenantId</includeMdcKeyName>
<includeMdcKeyName>usuarioId</includeMdcKeyName>
<customFields>{"servicio":"${app}","entorno":"produccion"}</customFields>
<fieldNames><timestamp>@timestamp</timestamp></fieldNames>
<throwableConverter class="net.logstash.logback.stacktrace.ShortenedThrowableConverter">
<maxDepthPerThrowable>30</maxDepthPerThrowable>
<exclude>^sun\.reflect\..*\.invoke</exclude>
<rootCauseFirst>true</rootCauseFirst>
</throwableConverter>
</encoder>
</appender>
<root level="INFO"><appender-ref ref="JSON"/></root>
<logger name="com.ejemplo" level="INFO"/>
<logger name="org.hibernate.SQL" level="WARN"/>
</springProfile>
</configuration>
// Enriquecer el MDC con contexto de negocio: aparece en TODAS las líneas del hilo
@Component
class MdcFilter extends OncePerRequestFilter {
@Override
protected void doFilterInternal(HttpServletRequest req, HttpServletResponse res,
FilterChain chain) throws ServletException, IOException {
try {
// traceId y spanId los pone Micrometer Tracing automáticamente
Optional.ofNullable(req.getHeader("X-Tenant-Id")).ifPresent(t -> MDC.put("tenantId", t));
Optional.ofNullable(req.getUserPrincipal())
.ifPresent(p -> MDC.put("usuarioId", p.getName()));
MDC.put("ruta", req.getRequestURI());
chain.doFilter(req, res);
} finally {
MDC.clear(); // OBLIGATORIO: los hilos se reutilizan y el contexto se filtraría
}
}
}
// Logging estructurado con pares clave-valor (logstash-logback-encoder)
import static net.logstash.logback.argument.StructuredArguments.kv;
log.info("Pedido confirmado",
kv("pedidoId", pedido.id()),
kv("clienteId", pedido.clienteId()),
kv("totalCentimos", pedido.total().importe().movePointRight(2).longValue()),
kv("numeroLineas", pedido.lineas().size()));
// Produce: {"@timestamp":"...","level":"INFO","message":"Pedido confirmado",
// "traceId":"4bf92f3577b34da6","spanId":"00f067aa0ba902b7",
// "pedidoId":"...","clienteId":"...","totalCentimos":12050,"numeroLineas":3}
// → se puede consultar con: pedidoId:"abc-123" o totalCentimos > 100000
9.3 Métricas con Micrometer: RED, USE y cardinalidad
DOS MÉTODOS COMPLEMENTARIOS
MÉTODO RED — para SERVICIOS (lo que ve el usuario)
Rate → peticiones por segundo
Errors → peticiones fallidas por segundo (y su porcentaje)
Duration → distribución de latencias (p50, p95, p99, p99.9)
MÉTODO USE — para RECURSOS (CPU, memoria, disco, pools, colas)
Utilization → % de tiempo ocupado
Saturation → cuánto trabajo hay ESPERANDO (¡la más predictiva de todas!)
Errors → errores del recurso
La SATURACIÓN es la que avisa ANTES del incidente: el pool de conexiones al 100 %
de uso con 40 hilos esperando es un incidente que ocurrirá en 5 minutos.
TIPOS DE MÉTRICA
Counter → solo sube. Peticiones, errores, eventos publicados.
Gauge → sube y baja. Conexiones activas, tamaño de cola, lag de Kafka.
Timer → duración + conteo. Latencia de endpoints y de llamadas externas.
Distribution Summary → distribución de un valor no temporal (tamaño de payload).
⚠️ CARDINALIDAD: EL ERROR QUE TUMBA PROMETHEUS
Cada combinación ÚNICA de etiquetas crea una serie temporal en memoria.
❌ Timer.builder("http.peticiones").tag("usuarioId", id) // 1.000.000 series
❌ .tag("url", urlCompleta) // infinitas (query params)
❌ .tag("pedidoId", id) // infinitas
✅ Timer.builder("http.peticiones").tag("ruta", "/api/pedidos/{id}") // plantilla
✅ .tag("metodo", "POST")
✅ .tag("estado", "200")
Regla práctica: una etiqueta no debe tener más de ~100 valores distintos, y el
producto de todas ellas debe quedarse por debajo de unos pocos miles de series.
¿Necesitas buscar por pedidoId? Eso son LOGS o TRAZAS, no métricas.
@Configuration
class MetricasConfig {
/** Etiquetas comunes a TODAS las métricas: imprescindibles para filtrar en Grafana. */
@Bean
MeterRegistryCustomizer<MeterRegistry> comunes(
@Value("${spring.application.name}") String app,
@Value("${ENTORNO:local}") String entorno,
@Value("${VERSION:dev}") String version) {
return registry -> registry.config()
.commonTags("aplicacion", app, "entorno", entorno, "version", version)
// Cortafuegos anti-cardinalidad: si alguien mete una etiqueta prohibida, se ignora
.meterFilter(MeterFilter.ignoreTags("usuarioId", "pedidoId", "email"))
.meterFilter(MeterFilter.maximumAllowableTags(
"http.server.requests", "uri", 100, MeterFilter.deny()));
}
}
@Service
class ServicioPedidosInstrumentado {
private final Counter confirmados;
private final Counter rechazados;
private final Timer tiempoConfirmacion;
private final DistributionSummary lineasPorPedido;
ServicioPedidosInstrumentado(MeterRegistry registro, ColaTrabajo cola) {
this.confirmados = Counter.builder("pedidos.confirmados")
.description("Pedidos confirmados con éxito")
.baseUnit("pedidos")
.register(registro);
this.rechazados = Counter.builder("pedidos.rechazados").register(registro);
this.tiempoConfirmacion = Timer.builder("pedidos.confirmacion.duracion")
.publishPercentiles(0.5, 0.95, 0.99) // percentiles calculados en la app
.publishPercentileHistogram() // histograma: permite agregar entre pods
.serviceLevelObjectives(Duration.ofMillis(200), Duration.ofMillis(500))
.register(registro);
this.lineasPorPedido = DistributionSummary.builder("pedidos.lineas")
.publishPercentiles(0.5, 0.95).register(registro);
// Gauge: se muestrea, no se incrementa. Ojo con las referencias fuertes.
Gauge.builder("cola.pendientes", cola, ColaTrabajo::tamano)
.description("Trabajos pendientes en la cola interna")
.register(registro);
}
ResultadoConfirmacion confirmar(ConfirmarPedidoComando comando) {
return tiempoConfirmacion.record(() -> {
try {
var resultado = casoDeUso.ejecutar(comando);
confirmados.increment();
lineasPorPedido.record(resultado.numeroLineas());
return resultado;
} catch (ReglaDeNegocioException e) {
rechazados.increment();
throw e;
}
});
}
}
// Métricas de negocio con @Timed y @Counted (requieren spring-boot-starter-aop)
@Timed(value = "facturas.generacion", percentiles = {0.5, 0.95, 0.99},
extraTags = {"tipo", "electronica"})
public Factura generar(PedidoId pedidoId) { /* ... */ }
publishPercentileHistogram(), que exporta buckets y permite calcular el percentil
en Prometheus con histogram_quantile). Si solo exportas percentiles precalculados, tus
paneles globales estarán mintiendo.
9.4 Trazas distribuidas con OpenTelemetry
ANATOMÍA DE UNA TRAZA
traceId = 4bf92f3577b34da6a3ce929d0e0e4736 (el MISMO en todos los servicios)
├─ span "POST /api/checkout" gateway [0 ────────────── 320 ms]
│ spanId=00f067aa0ba902b7 parent=null
│
├──── span "POST /pedidos" pedidos [ 8 ─────────── 310 ms]
│ spanId=a1b2c3 parent=00f067aa0ba902b7
│ atributos: pedido.id=abc, cliente.tipo=vip
│
│ ├── span "SELECT pedido" pedidos [ 12 ── 25 ms]
│ │
│ ├── span "GET /stock" inventario [ 30 ────── 95 ms]
│ │ └── span "SELECT stock" inventario [ 40 ── 88 ms] ← 48 ms aquí
│ │
│ ├── span "POST /pagos" pagos [100 ────────── 290 ms]
│ │ └── span "HTTP pasarela" pagos [110 ───────── 285 ms] ← 175 ms!
│ │ atributos: pasarela=stripe, http.status=200
│ │
│ └── span "INSERT outbox" pedidos [295 ── 305 ms]
│
└──── span "publicar evento" pedidos [312 ── 318 ms]
DIAGNÓSTICO EN 10 SEGUNDOS: el 55 % del tiempo se va en la pasarela externa.
Sin traza distribuida, esta conclusión requiere horas y varias personas.
PROPAGACIÓN DEL CONTEXTO — cabecera W3C Trace Context (estándar)
traceparent: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01
▲ ▲ ▲ ▲
versión trace-id (16 bytes) parent-id (8 B) flags
01 = muestreado
tracestate: vendor1=valor,vendor2=valor (información específica del proveedor)
⚠️ Esta cabecera debe atravesar TODO: HTTP, cabeceras de Kafka, colas, tareas
programadas y llamadas asíncronas. Un solo salto que la pierda parte la traza
en dos trazas huérfanas y la investigación se vuelve imposible.
MUESTREO (sampling) — no puedes guardar el 100 % a escala
· head-based (el más común): se decide al inicio. Simple y barato.
1 % de las trazas normales, 100 % de las que tienen error.
· tail-based: se decide al final, con la traza completa; permite quedarse con
TODAS las lentas y erróneas. Necesita un colector con memoria y estado.
· Regla de oro: SIEMPRE 100 % de errores y de trazas lentas. Lo aburrido se muestrea.
<!-- Spring Boot 3: Micrometer Tracing con puente a OpenTelemetry -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-tracing-bridge-otel</artifactId>
</dependency>
<dependency>
<groupId>io.opentelemetry</groupId>
<artifactId>opentelemetry-exporter-otlp</artifactId>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-registry-prometheus</artifactId>
</dependency>
management:
endpoints:
web:
exposure:
include: health,info,metrics,prometheus,loggers,env,threaddump,heapdump
base-path: /actuator
endpoint:
health:
show-details: when-authorized
probes:
enabled: true # /health/liveness y /health/readiness
loggers:
access: read_only # cambiar niveles de log sin reiniciar
tracing:
enabled: true
sampling:
probability: 0.1 # 10 % en producción; 1.0 en preproducción
propagation:
type: w3c # estándar; usa b3 solo por compatibilidad
otlp:
tracing:
endpoint: http://otel-collector:4318/v1/traces
timeout: 10s
metrics:
export:
enabled: false # las métricas van por Prometheus (pull)
metrics:
tags:
aplicacion: ${spring.application.name}
distribution:
percentiles-histogram:
http.server.requests: true # histograma agregable entre instancias
slo:
http.server.requests: 50ms,100ms,200ms,500ms,1s,2s
enable:
jvm: true
process: true
observations:
key-values:
entorno: ${ENTORNO:local}
# Correlación automática en los logs: Spring Boot 3 añade traceId y spanId al MDC
logging:
pattern:
correlation: "[${spring.application.name},%X{traceId:-},%X{spanId:-}] "
include-application-name: false
// ─── Instrumentación manual: @Observed crea span + métrica + log correlacionado ───
@Service
class ServicioPrecios {
@Observed(name = "precios.calculo",
contextualName = "calcular-precio-final",
lowCardinalityKeyValues = {"motor", "v2"})
public Dinero calcular(Pedido pedido) {
return motor.calcular(pedido);
}
}
@Configuration
class ObservacionConfig {
/** Necesario para que @Observed funcione (aspecto de AOP). */
@Bean
ObservedAspect observedAspect(ObservationRegistry registry) {
return new ObservedAspect(registry);
}
}
// ─── Spans manuales con la API de Tracer, para tramos internos que importan ───
@Service
class ImportadorCatalogo {
private final Tracer tracer;
void importar(List<ProductoDto> productos) {
Span span = tracer.nextSpan().name("importar-catalogo").start();
try (Tracer.SpanInScope ignored = tracer.withSpan(span)) {
// Atributos de BAJA cardinalidad como etiquetas; los identificadores, como evento
span.tag("catalogo.tamano.rango", rangoDe(productos.size()));
span.tag("catalogo.origen", "proveedor-a");
for (var lote : particionar(productos, 500)) {
Span spanLote = tracer.nextSpan().name("procesar-lote").start();
try (var s = tracer.withSpan(spanLote)) {
procesar(lote);
} catch (Exception e) {
spanLote.error(e); // marca el span como fallido
throw e;
} finally {
spanLote.end();
}
}
} catch (Exception e) {
span.error(e);
throw e;
} finally {
span.end(); // SIEMPRE en finally: si no, fuga de spans
}
}
}
// ─── Obtener el traceId para devolvérselo al usuario en un error ───
@RestControllerAdvice
class ErroresConTraza {
private final Tracer tracer;
@ExceptionHandler(Exception.class)
ProblemDetail error(Exception e) {
var pd = ProblemDetail.forStatusAndDetail(HttpStatus.INTERNAL_SERVER_ERROR,
"Se ha producido un error inesperado.");
Optional.ofNullable(tracer.currentSpan())
.ifPresent(s -> pd.setProperty("traceId", s.context().traceId()));
return pd;
// El usuario ve un id que puede dar a soporte; soporte encuentra la traza en 5 segundos.
}
}
9.5 Stacks de monitorización
| Pieza | Opción libre | Alternativas | Qué hace |
|---|---|---|---|
| Métricas | Prometheus (modelo pull) | VictoriaMetrics, Mimir, Datadog | Recolecta y almacena series temporales; lenguaje PromQL. |
| Paneles | Grafana | Datadog, New Relic, Kibana | Visualización unificada de métricas, logs y trazas. |
| Logs | Loki (indexa etiquetas, no contenido: barato) | Elasticsearch/OpenSearch, Datadog | Agregación y búsqueda de logs. |
| Trazas | Tempo o Jaeger | Zipkin, Datadog APM | Almacenamiento y consulta de trazas distribuidas. |
| Recolección | OpenTelemetry Collector | Agentes propietarios | Recibe, procesa, muestrea y reenvía telemetría. Te desacopla del proveedor. |
| Alertas | Alertmanager | PagerDuty, Opsgenie | Enrutamiento, agrupación, silenciado y escalado de alertas. |
# otel-collector-config.yaml — el punto único de control de la telemetría
receivers:
otlp:
protocols:
grpc: { endpoint: 0.0.0.0:4317 }
http: { endpoint: 0.0.0.0:4318 }
processors:
batch:
timeout: 5s
send_batch_size: 1024
memory_limiter:
check_interval: 1s
limit_mib: 1024
resource:
attributes:
- { key: deployment.environment, value: produccion, action: upsert }
attributes: # eliminar datos sensibles ANTES de almacenar
actions:
- { key: http.request.header.authorization, action: delete }
- { key: usuario.email, action: delete }
- { key: usuario.dni, action: hash }
tail_sampling: # muestreo inteligente: guarda lo que importa
decision_wait: 10s
policies:
- name: errores-siempre
type: status_code
status_code: { status_codes: [ERROR] }
- name: lentas-siempre
type: latency
latency: { threshold_ms: 1000 }
- name: resto-uno-por-ciento
type: probabilistic
probabilistic: { sampling_percentage: 1 }
exporters:
otlp/tempo:
endpoint: tempo:4317
tls: { insecure: true }
prometheus:
endpoint: 0.0.0.0:8889
service:
pipelines:
traces:
receivers: [otlp]
processors: [memory_limiter, tail_sampling, resource, attributes, batch]
exporters: [otlp/tempo]
metrics:
receivers: [otlp]
processors: [memory_limiter, resource, batch]
exporters: [prometheus]
9.6 SLI, SLO, SLA y presupuesto de error
DEFINICIONES QUE SE CONFUNDEN CONSTANTEMENTE
SLI (Indicator) = la MEDIDA. "% de peticiones con éxito y < 300 ms"
SLO (Objective) = el OBJETIVO. "99,9 % de las peticiones, en 30 días"
SLA (Agreement) = el CONTRATO con penalización económica. Siempre MÁS LAXO que
el SLO interno: si tu SLA es 99,9 %, tu SLO debe ser 99,95 %,
para enterarte antes de incumplir y pagar.
PRESUPUESTO DE ERROR (error budget) — la idea más útil de todo el SRE
SLO = 99,9 % en 30 días → presupuesto de error = 0,1 % = 43,2 minutos/mes
┌───────────────────────────────────────────────────────────────────────┐
│ Presupuesto consumido: ████████████░░░░░░░░░░░░░░░░ 40 % (17 min) │
└───────────────────────────────────────────────────────────────────────┘
Y esto se convierte en una POLÍTICA DE EQUIPO acordada de antemano:
· Queda presupuesto → se despliega rápido, se experimenta, se asumen riesgos.
· Consumido el 100 % → CONGELACIÓN de features. Todo el equipo a fiabilidad
hasta que la ventana móvil se recupere.
Convierte la discusión eterna "¿features o estabilidad?" en un DATO objetivo.
Además deja claro que el 100 % de disponibilidad NO es el objetivo: es imposible,
carísimo, y significa que estás desplegando demasiado despacio.
TABLA DE DISPONIBILIDAD (para que los números signifiquen algo)
┌──────────┬───────────────┬──────────────┬─────────────┐
│ SLO │ Caída al mes │ Caída al año │ Coste │
├──────────┼───────────────┼──────────────┼─────────────┤
│ 99 % │ 7 h 18 min │ 3,65 días │ Bajo │
│ 99,9 % │ 43,8 min │ 8,76 h │ Medio │
│ 99,95 % │ 21,9 min │ 4,38 h │ Alto │
│ 99,99 % │ 4,38 min │ 52,6 min │ Muy alto │
│ 99,999 % │ 26 segundos │ 5,26 min │ Extremo │
└──────────┴───────────────┴──────────────┴─────────────┘
Cada "nueve" adicional multiplica el coste por 3-10. Elige el SLO que el NEGOCIO
necesita y está dispuesto a pagar, no el que suena mejor en una reunión.
# Reglas de alerta en Prometheus: por SÍNTOMA (lo que sufre el usuario),
# no por causa (lo que falla por dentro). Una alerta que no requiere acción es ruido.
groups:
- name: slo-pedidos
rules:
# SLI de disponibilidad: proporción de peticiones NO 5xx
- record: sli:pedidos:disponibilidad:ratio5m
expr: |
sum(rate(http_server_requests_seconds_count{aplicacion="pedidos",status!~"5.."}[5m]))
/
sum(rate(http_server_requests_seconds_count{aplicacion="pedidos"}[5m]))
# Alerta por VELOCIDAD DE CONSUMO del presupuesto de error (multiventana).
# Salta si en 1 h consumes 14,4x el ritmo permitido: agotarías el mes en 2 días.
- alert: PresupuestoErrorConsumiendoseRapido
expr: |
(1 - sli:pedidos:disponibilidad:ratio5m) > (14.4 * 0.001)
and
(1 - sli:pedidos:disponibilidad:ratio1h) > (14.4 * 0.001)
for: 2m
labels: { severity: pagina }
annotations:
summary: "Pedidos consume el presupuesto de error 14x más rápido de lo permitido"
runbook: "https://wiki.ejemplo.com/runbooks/pedidos-errores"
- alert: LatenciaP99Degradada
expr: |
histogram_quantile(0.99,
sum by (le) (rate(http_server_requests_seconds_bucket{aplicacion="pedidos"}[5m]))
) > 1
for: 10m
labels: { severity: aviso }
annotations:
summary: "p99 de pedidos por encima de 1 s durante 10 minutos"
- alert: ConsumerLagCreciente
expr: kafka_consumergroup_lag{group="facturacion"} > 10000
for: 5m
labels: { severity: aviso }
- alert: MensajesEnDLT
expr: increase(kafka_consumer_records_consumed_total{topic=~".*\\.DLT"}[5m]) > 0
for: 1m
labels: { severity: aviso }
annotations:
summary: "Hay mensajes llegando a una dead letter topic"
- alert: CircuitBreakerAbierto
expr: resilience4j_circuitbreaker_state{state="open"} == 1
for: 1m
labels: { severity: aviso }
9.7 Health checks bien hechos
/**
* Indicador de salud PERSONALIZADO. La distinción entre liveness y readiness es
* crítica y casi siempre está mal:
*
* LIVENESS ("¿estoy vivo?") → si falla, Kubernetes REINICIA el pod.
* Solo debe fallar por estados irrecuperables
* (deadlock, memoria corrupta). NUNCA debe
* comprobar dependencias externas.
* READINESS ("¿puedo atender?")→ si falla, se retira del balanceador pero NO
* se reinicia. Aquí SÍ se comprueban las
* dependencias imprescindibles.
*
* ❌ El error clásico: poner la comprobación de la base de datos en LIVENESS.
* La BD tiene un hipo de 30 s → Kubernetes reinicia TODOS los pods a la vez →
* avalancha de reconexiones → la BD cae del todo → bucle de reinicios infinito.
*/
@Component("baseDatos")
class BaseDatosHealthIndicator implements HealthIndicator {
private final JdbcTemplate jdbc;
@Override
public Health health() {
try {
Integer uno = jdbc.queryForObject("SELECT 1", Integer.class);
var pool = (HikariDataSource) jdbc.getDataSource();
var mx = pool.getHikariPoolMXBean();
Health.Builder estado = (uno != null && uno == 1) ? Health.up() : Health.down();
return estado
.withDetail("conexionesActivas", mx.getActiveConnections())
.withDetail("conexionesInactivas", mx.getIdleConnections())
.withDetail("hilosEsperando", mx.getThreadsAwaitingConnection()) // SATURACIÓN
.build();
} catch (Exception e) {
return Health.down(e).build();
}
}
}
// Grupos de salud: separar liveness de readiness explícitamente
// management.endpoint.health.group.liveness.include=livenessState,deadlockDetector
// management.endpoint.health.group.readiness.include=readinessState,baseDatos,kafka
Checklist — Spring Cloud y observabilidad
10 · Seguridad entre servicios
Ampliación natural del módulo 10 · Seguridad, aquí desde la perspectiva de la comunicación entre servicios.
10.1 Confianza cero: la red interna no es segura
El modelo antiguo era el castillo con foso: un perímetro duro y, dentro, confianza total. Ese modelo murió con la nube. Hoy se asume que el atacante ya está dentro —por un contenedor comprometido, una dependencia maliciosa o un error de configuración— y cada llamada se autentica y autoriza como si viniera de internet.
CASTILLO Y FOSO (obsoleto) CONFIANZA CERO (zero trust)
Internet Internet
│ 🔒 firewall │ 🔒
▼ ▼
┌───────────────────────────┐ ┌───────────────────────────────┐
│ RED INTERNA "de confianza"│ │ Cada llamada: │
│ │ │ · identidad verificada │
│ A ──► B ──► C sin auth │ │ · mTLS entre servicios │
│ (si estás dentro, pasas) │ │ · autorización explícita │
│ │ │ · mínimo privilegio │
│ Un pod comprometido tiene│ │ · todo auditado │
│ acceso a TODO. │ │ │
└───────────────────────────┘ │ A ──mTLS+token──► B │
│ B ──mTLS+token──► C │
El movimiento lateral es trivial. └───────────────────────────────┘
Comprometer A no da acceso a C.
10.2 mTLS: autenticación mutua
En TLS normal solo el servidor presenta certificado. En mTLS también lo hace el cliente: ambos extremos demuestran quiénes son. Hacerlo a mano es un infierno de gestión de certificados y rotación; por eso, en la práctica, lo delega la plataforma.
| Implementación | Coste | Rotación de certificados | Cuándo |
|---|---|---|---|
| Malla de servicios (Istio, Linkerd) | Alto al montarla, nulo después | Automática, cada 24 h | Lo estándar hoy en Kubernetes. mTLS «gratis» y sin tocar código. |
| cert-manager con SPIFFE/SPIRE | Medio | Automática | Identidad de carga de trabajo sin malla completa. |
| Manual en la aplicación (keystore/truststore) | Alto y permanente | Manual: la fuente número uno de caídas por certificado caducado | Solo si no hay orquestador. Evítalo. |
10.3 Tokens de servicio y propagación de identidad
DOS IDENTIDADES DISTINTAS QUE HAY QUE DISTINGUIR SIEMPRE
1. IDENTIDAD DEL SERVICIO ("¿quién llama?")
→ mTLS o token de cliente OAuth2 (client_credentials)
→ responde a: "¿puede el servicio Pedidos llamar a Pagos?"
2. IDENTIDAD DEL USUARIO ("¿en nombre de quién?")
→ JWT del usuario propagado, o token de intercambio (token exchange, RFC 8693)
→ responde a: "¿puede María ver el pedido 42?"
❌ ANTIPATRÓN: reenviar el JWT del usuario tal cual por toda la cadena.
· Cualquier servicio de la cadena puede REUTILIZARLO para llamar a otros
en nombre del usuario, sin restricción y sin traza.
· El token dura demasiado y su alcance es demasiado amplio.
· Un servicio comprometido se convierte en el usuario.
✅ PATRÓN RECOMENDADO: intercambio de token en el gateway
Usuario ──JWT amplio──► GATEWAY
│ valida el JWT una vez
│ intercambia por un token INTERNO:
│ · audiencia = servicio destino
│ · alcance mínimo necesario
│ · vida corta (60 s)
│ · incluye "act" (quién actúa) y "sub" (usuario)
▼
Servicio ──► token acotado ──► siguiente servicio
// Servicio de recursos: valida el JWT y aplica autorización de grano fino
@Configuration
@EnableWebSecurity
@EnableMethodSecurity
class SeguridadConfig {
@Bean
SecurityFilterChain filtros(HttpSecurity http) throws Exception {
return http
.csrf(AbstractHttpConfigurer::disable) // API sin sesión: no aplica
.sessionManagement(s -> s.sessionCreationPolicy(SessionCreationPolicy.STATELESS))
.authorizeHttpRequests(a -> a
.requestMatchers("/actuator/health/**").permitAll()
.requestMatchers("/actuator/**").hasAuthority("SCOPE_actuator")
.requestMatchers(HttpMethod.GET, "/api/v1/pedidos/**").hasAuthority("SCOPE_pedidos:leer")
.requestMatchers("/api/v1/pedidos/**").hasAuthority("SCOPE_pedidos:escribir")
.anyRequest().denyAll()) // ❗ denegar por defecto
.oauth2ResourceServer(o -> o.jwt(j -> j
.jwtAuthenticationConverter(convertidor())))
.build();
}
}
// Autorización de grano fino: el alcance no basta, hay que comprobar la PROPIEDAD del dato
@Service
class ConsultarPedidoService {
@PreAuthorize("hasAuthority('SCOPE_pedidos:leer')")
@PostAuthorize("returnObject.clienteId().valor().toString() == authentication.name "
+ "or hasRole('ADMIN')")
public PedidoVista porId(PedidoId id) {
return repositorio.vistaPorId(id).orElseThrow(() -> new PedidoNoEncontradoException(id));
}
}
GET /api/v1/pedidos/{id} devuelve
cualquier pedido si conoces el UUID, porque nadie comprueba que el pedido sea tuyo. La
autenticación (quién eres) no es la autorización (qué puedes ver). Cada servicio comprueba la
propiedad del dato, siempre, aunque el gateway ya haya autenticado.
11 · Despliegue y operación
11.1 Despliegues sin caída
| Estrategia | Cómo funciona | Coste | Rollback | Riesgo |
|---|---|---|---|---|
| Rolling update (por defecto en K8s) | Sustituye réplicas poco a poco respetando maxUnavailable y maxSurge. | Bajo | Minutos (rollout deshacer) | Conviven dos versiones: los contratos deben ser compatibles. |
| Blue-green | Dos entornos completos; se conmuta el tráfico de golpe. | Alto: el doble de recursos | Instantáneo | Migraciones de BD compartidas entre azul y verde. |
| Canary | 1 % → 10 % → 50 % → 100 %, con métricas que deciden en cada paso. | Medio | Rápido | Necesita buenas métricas y automatización (Argo Rollouts, Flagger). |
| Shadow / espejo | Se duplica el tráfico real a la versión nueva sin devolver su respuesta. | Alto | N/A | Cuidado con los efectos secundarios: no dupliques cobros ni emails. |
11.2 Migraciones de esquema compatibles: expand and contract
RENOMBRAR UNA COLUMNA SIN CAÍDA (el ejemplo canónico)
❌ INGENUO: ALTER TABLE cliente RENAME telefono TO telefono_movil;
Durante el rolling update conviven v1 (lee `telefono`) y v2 (lee `telefono_movil`).
La versión antigua se rompe al instante. Caída garantizada.
✅ EXPAND AND CONTRACT — cuatro despliegues, cero caída
PASO 1 · EXPAND (solo BD, sin código nuevo)
ALTER TABLE cliente ADD COLUMN telefono_movil VARCHAR(20); -- NULLABLE
-- backfill por lotes, sin bloquear la tabla:
UPDATE cliente SET telefono_movil = telefono WHERE id BETWEEN ? AND ?;
-- trigger o código que mantenga ambas sincronizadas mientras dure la migración
PASO 2 · ESCRIBIR EN AMBAS, LEER DE LA VIEJA
v2 escribe telefono Y telefono_movil; sigue leyendo telefono.
(compatible con v1, que sigue funcionando igual)
PASO 3 · LEER DE LA NUEVA
v3 lee telefono_movil; sigue escribiendo en ambas.
← Punto de no retorno controlado: si algo falla, vuelves a v2.
PASO 4 · CONTRACT (semanas después, cuando NO queda ninguna v2 viva)
v4 solo usa telefono_movil.
ALTER TABLE cliente DROP COLUMN telefono;
REGLAS DE ORO DE LAS MIGRACIONES
1. Toda migración debe ser compatible con la versión ANTERIOR del código.
2. Nunca en el mismo despliegue que el código que la necesita.
3. Nada de DROP ni de NOT NULL en el mismo paso que se añade algo.
4. Backfill por lotes con pausas: un UPDATE de 10 millones de filas bloquea la tabla.
5. En PostgreSQL: CREATE INDEX CONCURRENTLY, y ADD COLUMN con default es barato
desde la versión 11, pero comprueba tu versión antes de fiarte.
6. Las migraciones deben ser IDEMPOTENTES y estar versionadas (Flyway, Liquibase).
11.3 Escalado horizontal y planificación de capacidad
REQUISITO PREVIO AL ESCALADO: SER STATELESS DE VERDAD
❌ Sesión HTTP en memoria → usa Redis o tokens sin estado
❌ Caché local sin coordinación → caché distribuida o TTL corto asumiendo divergencia
❌ Ficheros en disco local → almacenamiento de objetos (S3)
❌ Tareas @Scheduled en todas las réplicas → ShedLock o un único líder
❌ Contadores en variables estáticas → métricas o almacén compartido
DIMENSIONAR CON LA LEY DE LITTLE Y DATOS REALES
Datos medidos: 2.000 req/s en pico · latencia media 80 ms · 4 vCPU por pod
Concurrencia necesaria = 2.000 × 0,08 = 160 peticiones simultáneas
Con 200 hilos por pod y un objetivo del 70 % de uso → 160 / (200 × 0,7) ≈ 2 pods
Margen para fallos de zona y picos → mínimo 4 pods en 3 zonas.
⚠️ Y AHORA LA PARTE QUE SE OLVIDA: ¿aguanta la base de datos?
4 pods × 20 conexiones = 80 conexiones. PostgreSQL con max_connections=100
se queda sin margen para migraciones ni para el resto de servicios.
Escalar la aplicación SIN escalar (o poner PgBouncer delante de) la BD
no mejora nada: solo mueve el cuello de botella y lo hace más difícil de ver.
AUTOESCALADO (HPA) — qué métrica usar
· CPU: sirve para cargas ligadas a CPU. Inútil si tu servicio espera en E/S.
· Métrica personalizada (peticiones por segundo, lag de Kafka, tamaño de cola):
mucho mejor. Para consumidores de Kafka, escalar por LAG es lo correcto…
hasta el número de particiones, que es el techo real.
· Configura siempre stabilizationWindowSeconds para evitar el "flapping".
12 · Casos reales y decisiones justificadas
Caso 1 · La startup que se adelantó
Contexto: 6 desarrolladores, producto sin encaje de mercado todavía, 500 usuarios. Deciden 12 microservicios «para estar preparados para escalar».
Qué pasó: cada feature tocaba 3–4 repositorios. El entorno local necesitaba 12 contenedores y 14 GB de RAM. Nadie tenía tiempo de montar trazas distribuidas, así que depurar era adivinar. La velocidad de entrega cayó un 60 % en cuatro meses y dos personas se marcharon.
Qué se hizo: consolidar en un monolito modular con tres módulos y fronteras verificadas con ArchUnit, dejando fuera solo el procesamiento de imágenes (escalado dispar real).
Lección: los microservicios optimizan la autonomía organizativa. Sin varios equipos que coordinar, solo pagas el coste.
Caso 2 · La cascada de reintentos
Contexto: comercio electrónico, viernes de Black Friday. El servicio de precios se degrada por una consulta sin índice: p99 pasa de 50 ms a 4 s.
Qué pasó: gateway, pedidos y carrito tenían 3 reintentos cada uno, sin jitter y sin circuit breaker. Los 3 niveles multiplicaron la carga por 27 justo cuando precios estaba peor. En 90 segundos cayó todo el sitio, no solo los precios.
Qué se hizo: reintentos solo en el nivel más cercano al fallo, con jitter y presupuesto; circuit breaker con fallback a precio cacheado; y un test de carga que reproduce el escenario en CI.
Lección: los reintentos mal configurados convierten una degradación en una caída total.
Caso 3 · Los pedidos fantasma
Contexto: un servicio guardaba el pedido y después publicaba el evento en Kafka, sin outbox.
Qué pasó: durante un reinicio del clúster de Kafka, unos 400 pedidos se guardaron en la base de datos pero su evento nunca se publicó. Los clientes tenían el pedido en «Mis pedidos», pero jamás se facturaron ni se enviaron. Se detectó tres semanas después por reclamaciones, y hubo que reconstruirlos a mano cruzando tablas.
Qué se hizo: outbox transaccional con publicador por polling y
SKIP LOCKED, más una alerta sobre la antigüedad del registro pendiente más viejo.
Lección: el dual write no falla en las pruebas; falla en producción, en silencio y semanas después.
Caso 4 · La partición caliente
Contexto: plataforma SaaS multi-tenant. El topic de eventos usaba
tenantId como clave, con 24 particiones y 24 consumidores.
Qué pasó: un cliente representaba el 60 % del volumen. Su partición acumulaba millones de mensajes de lag mientras las otras 23 estaban ociosas. Añadir consumidores no servía de nada: el paralelismo lo limita la partición.
Qué se hizo: clave compuesta tenantId + entidadId para repartir el
tráfico manteniendo el orden donde importa (por entidad, no por tenant), y un topic dedicado para el
cliente grande.
Lección: la clave de partición determina si puedes escalar. Elígela mirando la distribución real de tus datos, no el modelo conceptual.
Caso 5 · La saga sin timeout
Contexto: saga coreografiada de 5 pasos para contratar un seguro.
Qué pasó: un consumidor tenía un bug que descartaba silenciosamente ciertos eventos.
Las sagas se quedaban en ESPERANDO_VALIDACION para siempre. Como nadie medía el tiempo
en cada estado, se acumularon 12.000 contrataciones a medias durante dos meses, con dinero cobrado y
pólizas sin emitir.
Qué se hizo: pasar a saga orquestada con estado persistido, un vigilante que compensa pasados 5 minutos, y una alerta sobre el número de sagas no terminales con antigüedad superior a 10 minutos.
Lección: en sistemas asíncronos, lo que no se mide no existe. Toda saga necesita timeout, compensación y una métrica de sagas atascadas.
Caso 6 · La factura de observabilidad
Contexto: migración a microservicios con logging «completo» a un SaaS de observabilidad.
Qué pasó: se logueaba cada petición HTTP entrante y saliente a nivel INFO con el
cuerpo completo, y las métricas incluían usuarioId como etiqueta. La factura mensual
pasó de 800 € a 47.000 €, y Prometheus se quedaba sin memoria por 8 millones de series.
Qué se hizo: muestreo por cola (100 % de errores y lentas, 1 % del resto), cuerpos solo en DEBUG activable bajo demanda, filtro de cardinalidad en Micrometer y retención por niveles. Factura final: 3.200 €, sin perder capacidad de diagnóstico.
Lección: la observabilidad es un producto con presupuesto. Diseña qué NO guardas con el mismo cuidado con el que decides qué guardas.
13 · Errores comunes
| # | Error | Consecuencia | Solución |
|---|---|---|---|
| 1 | Empezar con microservicios «porque es lo moderno» | Coste enorme y cero beneficio; el equipo se ahoga en operación | Monolito modular primero; dividir con criterios objetivos |
| 2 | Cortar por capas técnicas (servicio-controladores, servicio-repositorios) | Cada feature toca todos los servicios: monolito distribuido | Cortar por contextos delimitados y capacidades de negocio |
| 3 | Base de datos compartida entre servicios | Acoplamiento oculto imposible de refactorizar | Una BD por servicio; integrar por API y eventos |
| 4 | Llamadas remotas sin timeout | Agotamiento de hilos y caída en cascada | Timeout en toda llamada, decreciente hacia abajo |
| 5 | Reintentar en varios niveles a la vez | Amplificación ×27: los reintentos matan al servicio degradado | Reintentar en un solo nivel, con jitter y presupuesto |
| 6 | Reintentar operaciones no idempotentes | Cobros y pedidos duplicados | Clave de idempotencia o identificador generado por el cliente |
| 7 | Dual write: guardar en BD y publicar en Kafka | Pérdida silenciosa de eventos o eventos fantasma | Outbox transaccional (polling o CDC) |
| 8 | Consumidores no idempotentes con at-least-once | Facturas y emails duplicados | Inbox, UPSERT o restricción UNIQUE de negocio |
| 9 | Circuit breaker que cuenta los 4xx como fallos | El circuito se abre por peticiones correctas | ignore-exceptions con las excepciones de negocio |
| 10 | @TimeLimiter sobre un método bloqueante | El timeout no se aplica; falsa sensación de seguridad | Timeout en el cliente HTTP; @TimeLimiter solo con CompletableFuture |
| 11 | Cadenas síncronas largas | Latencia sumada y disponibilidad multiplicada | Núcleo síncrono mínimo; el resto por eventos |
| 12 | Clave de partición con poca variedad | Particiones calientes; imposible escalar | Clave con alta cardinalidad y distribución uniforme |
| 13 | Más consumidores que particiones | Consumidores ociosos; el escalado no hace nada | Dimensionar particiones con holgura desde el principio |
| 14 | Sin dead letter topic ni alerta | Mensajes envenenados en bucle infinito; lag creciente | DLT, ErrorHandlingDeserializer y alerta si la DLT recibe algo |
| 15 | Saga sin timeout ni compensación | Procesos de negocio colgados para siempre | Vigilante periódico, estado persistido y métricas de sagas atascadas |
| 16 | Logs sin traceId | Investigar un incidente pasa de 5 min a 3 h | Micrometer Tracing y JSON estructurado con MDC |
| 17 | Etiquetas de métricas de alta cardinalidad | Prometheus se queda sin memoria; factura disparada | MeterFilter que bloquee etiquetas prohibidas |
| 18 | Promediar percentiles entre instancias | Paneles que mienten; decisiones erróneas | Histogramas y histogram_quantile |
| 19 | Comprobar la base de datos en el probe de liveness | Reinicio masivo de pods y bucle de fallo | Dependencias solo en readiness; liveness mínimo |
| 20 | Migración de esquema destructiva en el mismo despliegue | La versión antigua se rompe durante el rolling update | Expand and contract en cuatro pasos |
| 21 | Reenviar el JWT del usuario por toda la cadena | Un servicio comprometido suplanta al usuario en todo el sistema | Intercambio de token con audiencia y alcance mínimos |
| 22 | Confiar en que «el gateway ya autorizó» | IDOR: cualquiera lee datos ajenos conociendo el UUID | Cada servicio comprueba la propiedad del dato |
| 23 | Librería común con el modelo de dominio dentro | Cambiar un campo obliga a redesplegar todo | Compartir solo utilidades técnicas; duplicar DTOs a propósito |
| 24 | Evento de dominio publicado tal cual como evento de integración | Refactorizar tu dominio rompe a otros equipos | Traducir en la frontera; contrato versionado con esquema |
| 25 | Feature flags que nunca se borran | Combinatoria inmanejable y código muerto | Fecha de caducidad por bandera y test que falla al expirar |
14 · Preguntas de entrevista
¿Cuándo NO usarías microservicios?
Cuando el equipo es pequeño (menos de 15–20 personas), cuando el dominio todavía no está claro y las fronteras van a cambiar, cuando no hay capacidad para operar la plataforma (CI/CD, observabilidad, guardias) o cuando el producto no tiene escalado dispar ni ciclos de vida distintos. Empezaría con un monolito modular con fronteras verificadas y extraería servicios cuando aparezca un criterio objetivo: equipos que se bloquean, escalado muy dispar, ciclos de vida incompatibles o necesidad de aislar fallos. Lo importante es que la decisión sea reversible mientras se pueda.
¿Cómo garantizas consistencia sin transacciones distribuidas?
Con transacciones locales encadenadas más compensaciones, es decir, el patrón saga. Cada servicio confirma en su propia base de datos y publica un evento; si un paso falla, se ejecutan compensaciones en orden inverso, que no son rollbacks sino hechos nuevos que contrarrestan a los anteriores. Para que el evento y el cambio de estado sean atómicos se usa el patrón outbox, y como la entrega es at-least-once, los consumidores deben ser idempotentes. Con más de cuatro pasos prefiero saga orquestada, porque el estado queda consultable y los timeouts son triviales. Y siempre acompaño esto de una conversación con negocio sobre cuánta incoherencia temporal es aceptable.
Explica el patrón outbox y por qué es necesario.
Resuelve el problema del dual write: escribir en la base de datos y publicar en el broker
son dos sistemas distintos, y no hay forma de hacerlo atómicamente. Puede pasar que confirmes en base
de datos y falle la publicación (el pedido existe y nadie se entera) o al revés (se factura un pedido
que no existe). Con outbox, en la misma transacción se guarda el agregado y se inserta una fila en una
tabla outbox; un publicador aparte —por polling con FOR UPDATE SKIP LOCKED o
por CDC con Debezium— la lee y la envía a Kafka. La garantía resultante es at-least-once,
porque el publicador puede morir tras enviar y antes de marcar la fila, así que el consumidor tiene que
deduplicar.
¿Cómo funciona un circuit breaker y qué debe contar como fallo?
Tiene tres estados principales. En CLOSED pasan todas las llamadas y se mide la tasa de
fallo sobre una ventana deslizante, con un mínimo de llamadas para no decidir con ruido. Al superar el
umbral pasa a OPEN y rechaza inmediatamente sin llamar, lo que protege tus hilos y da aire
al dependiente. Tras un tiempo de espera pasa a HALF_OPEN y deja pasar unas pocas llamadas
de prueba: si van bien vuelve a CLOSED, si no vuelve a OPEN. Resilience4j
añade DISABLED, FORCED_OPEN y METRICS_ONLY, este último ideal
para estrenarlo en producción sin riesgo. Como fallo deben contar timeouts, errores de red, 5xx y
llamadas lentas; nunca los 4xx de negocio, porque un 404 significa que el dependiente funciona
perfectamente.
¿Cómo evitas procesar dos veces el mismo mensaje?
Asumiendo desde el principio que va a llegar repetido. Lo primero es hacer la operación
naturalmente idempotente si se puede: un UPSERT por clave de negocio o fijar un estado
absoluto en lugar de incrementar. Si no se puede, uso una tabla inbox con clave primaria
(mensajeId, consumidor) e INSERT ... ON CONFLICT DO NOTHING, en la misma
transacción que el efecto de negocio, de modo que si el efecto falla también desaparece el registro y
el mensaje se puede reprocesar. Otra opción excelente es una restricción UNIQUE de negocio
que deje que la base de datos rechace el duplicado. Y para descartar eventos viejos que llegan
desordenados, comparo la versión del agregado.
¿Qué es realmente el «exactly-once» de Kafka?
Es exactamente-una-vez dentro de Kafka: con transacciones puedes leer de un topic,
procesar, escribir en otro topic y commitear los offsets de forma atómica, y el consumidor con
isolation.level=read_committed no ve los mensajes de transacciones abortadas. Funciona muy
bien en Kafka Streams con processing.guarantee=exactly_once_v2. Ahora bien, en cuanto tu
efecto secundario sale de Kafka —insertar en PostgreSQL, llamar a una API, enviar un correo— vuelves a
at-least-once, porque no existe transacción que abarque los dos sistemas: es el mismo problema del dual
write. Por eso en la práctica se diseña siempre para at-least-once con consumidores idempotentes.
¿Cómo depurarías una petición lenta que atraviesa cinco servicios?
Empezaría por las métricas para acotar: desde cuándo, qué endpoints, todos los pods o solo algunos,
y si coincide con un despliegue. Luego abriría una traza lenta de ejemplo, que muestra el desglose por
span y señala inmediatamente dónde está el tiempo, distinguiendo además el tiempo propio del tiempo
esperando a dependientes. Con el traceId filtraría los logs del servicio culpable para ver
el detalle: agotamiento del pool, GC, consulta lenta. Y si no hay traza que valga, miraría saturación
—hilos en espera, tamaño de colas, lag— porque suele avisar antes que la latencia. Todo esto requiere
haber invertido antes en propagación de contexto: sin ella, este trabajo pasa de minutos a horas.
¿Qué es la cardinalidad en métricas y por qué importa tanto?
Cada combinación única de etiquetas crea una serie temporal independiente que Prometheus mantiene en
memoria. Si etiquetas con usuarioId, pedidoId o la URL completa con
parámetros, generas cientos de miles o millones de series y tumbas el sistema de métricas, además de
disparar la factura si es un SaaS. La regla es usar solo etiquetas de baja cardinalidad —plantilla de
ruta, método, código de estado, servicio— y mantener el producto total en el orden de miles. Si
necesitas buscar por un identificador concreto, eso es trabajo de logs o de trazas, no de métricas.
Diferencia entre SLI, SLO y SLA, y qué es el presupuesto de error.
El SLI es la medida (por ejemplo, el porcentaje de peticiones correctas por debajo de 300 ms), el SLO es el objetivo interno sobre esa medida (99,9 % en 30 días) y el SLA es el contrato con el cliente, con penalización económica; el SLA siempre debe ser más laxo que el SLO para tener margen de reacción. El presupuesto de error es el complemento del SLO: con un 99,9 % dispones de 43 minutos de fallo al mes. Su utilidad es política además de técnica: mientras quede presupuesto se despliega rápido y se asumen riesgos; si se agota, se congelan las features y el equipo se dedica a fiabilidad. Convierte la discusión entre velocidad y estabilidad en un dato objetivo acordado de antemano.
¿Cómo versionas una API sin romper a los consumidores?
Priorizo la evolución compatible: añadir campos opcionales, añadir endpoints, relajar validaciones y
que todos los clientes ignoren los campos desconocidos. Cuando el cambio rompe de verdad —eliminar o
renombrar un campo, cambiar un tipo, cambiar la semántica— publico una versión nueva en la ruta,
mantengo la anterior en paralelo, mido con una métrica quién sigue usándola, anuncio la retirada con
cabeceras Deprecation y Sunset, y la apago solo cuando el uso llega a cero.
Con eventos hago lo mismo pero con doble publicación de v1 y v2, y apoyándome en un registro de
esquemas con compatibilidad FULL para que el propio despliegue del productor falle si rompe algo.
¿Por qué el orden en Kafka solo está garantizado por partición?
Porque un topic es un conjunto de logs independientes y cada partición es un log propio con su secuencia de offsets. Las particiones viven en brokers distintos y se consumen en paralelo, así que no existe un reloj global que ordene entre ellas. El orden se controla con la clave: mismo valor de clave significa misma partición, y por tanto orden garantizado para esa entidad. Por eso se usa el identificador del agregado como clave. Si necesitaras orden total tendrías que usar una sola partición, lo que elimina el paralelismo; en la práctica casi nunca hace falta orden global, solo orden por entidad.
¿Qué tamaño debe tener un microservicio?
El de un contexto delimitado con su propio ciclo de vida, no una cifra de líneas de código. Señales de que es demasiado pequeño: no puede hacer nada útil sin llamar a otros tres, o cualquier cambio de negocio toca varios servicios a la vez. Señales de que es demasiado grande: dos equipos se estorban al desplegar, o hay partes con necesidades de escalado muy distintas. Un buen indicador práctico es que un equipo pueda ser dueño completo del servicio, incluida su guardia, y que la media de repositorios tocados por historia de usuario se mantenga cerca de uno.
¿Qué es un despliegue sin caída y cómo lo consigues con cambios de esquema?
Es sustituir la versión en ejecución sin que el usuario perciba errores; en Kubernetes se hace con rolling update apoyado en probes de readiness y apagado elegante. La dificultad real está en la base de datos, porque durante el despliegue conviven dos versiones del código sobre el mismo esquema. La solución es expand and contract: primero se añade lo nuevo de forma no destructiva y se rellena por lotes, luego se escribe en ambos formatos, después se lee del nuevo y, semanas más tarde, cuando ya no queda ninguna instancia antigua, se elimina lo viejo. La regla es que toda migración debe ser compatible con la versión anterior del código y nunca ir en el mismo despliegue que la necesita.
¿Saga orquestada o coreografiada?
Depende del número de pasos y de la necesidad de visibilidad. Con dos o tres pasos, la coreografía es más simple y desacoplada: cada servicio reacciona a eventos y no hay componente adicional. A partir de cuatro pasos prefiero orquestación, porque la lógica del proceso queda en un sitio, el estado es una fila consultable —puedes responder «¿en qué punto está la saga del pedido 42?» con una consulta—, los timeouts por paso son triviales y se puede reintentar un paso concreto. El riesgo de la orquestación es que el coordinador acumule lógica de negocio: debe dirigir, no decidir.
¿Cómo pruebas un sistema de microservicios?
Con una pirámide adaptada. La base son tests unitarios del dominio, rapidísimos porque la arquitectura hexagonal permite ejecutarlos sin Spring. Encima, tests de caso de uso con dobles en memoria. Después, tests de integración por servicio con Testcontainers para base de datos y broker reales, y WireMock para los dependientes HTTP, incluyendo escenarios de lentitud, error y caída. La pieza que evita las sorpresas entre equipos es el contract testing, que verifica el contrato sin levantar los dos servicios. Y en la cúspide, muy pocos tests de extremo a extremo, más smoke tests en producción tras el despliegue. Todo esto está desarrollado en el módulo 07.
¿Qué harías si el consumer lag de Kafka crece sin parar?
Primero distinguiría entre pico de tráfico y consumidor degradado, mirando la tasa de producción frente a la de consumo. Si el consumidor va lento, comprobaría el tiempo de procesamiento por mensaje, si hay rebalanceos frecuentes en los logs del coordinador y si el trabajo pesado está dentro del bucle de poll. Si el problema es capacidad, añadiría consumidores, pero solo hasta el número de particiones, que es el techo real; si ya estoy en el techo, hay que aumentar particiones o procesar por lotes. También revisaría si una partición concreta concentra el lag, lo que indicaría una clave mal elegida o un mensaje envenenado bloqueando la partición.
15 · Ejercicios prácticos
Ejercicio 1 · Diseñar las fronteras (papel y lápiz, 45 min)
Un marketplace tiene: registro de vendedores, catálogo de productos, búsqueda, carrito, pedidos, pagos, comisiones, envíos, valoraciones, notificaciones y un panel de analítica.
- Identifica los contextos delimitados y clasifícalos en núcleo, soporte o genérico.
- Dibuja el mapa de contextos indicando el patrón de relación de cada arista.
- Marca qué contextos serían microservicios desde el día 1 y cuáles empezarían como módulos.
- Busca dos palabras que signifiquen cosas distintas en dos contextos y modela ambas.
- Escribe en tres líneas qué le dirías a un director técnico que pide «12 microservicios».
Ejercicio 2 · Servicio hexagonal con tests rápidos (2–3 h)
- Crea un servicio de pedidos con Spring Boot 3 y Java 21, con los paquetes
dominio,aplicacioneinfraestructura. - Implementa el agregado
Pedidocon la invariante «total igual a la suma de líneas» y los objetos de valorDinero,PedidoIdySku. - Define los puertos
RepositorioPedidosyPublicadorEventos, y sus adaptadores JPA y Kafka. - Escribe tests del dominio y del caso de uso sin Spring, con dobles en memoria.
- Añade los ocho tests de ArchUnit de la sección 3.4 y comprueba que fallan si importas Spring en el dominio.
Criterio de éxito: la suite de dominio y aplicación completa se ejecuta en menos de 2 segundos.
Ejercicio 3 · Outbox de principio a fin (2 h)
- Crea la tabla
outboxcon su índice parcial mediante una migración de Flyway. - Implementa
PublicadorEventosOutboxescribiendo dentro de la transacción del caso de uso. - Implementa el publicador por polling con
FOR UPDATE SKIP LOCKED. - Escribe un test con Testcontainers que arranque PostgreSQL y Kafka y verifique que el evento llega.
- Escribe un test que provoque un rollback y compruebe que no se publica nada.
- Levanta dos instancias del publicador y comprueba que no duplican ni se bloquean entre ellas.
- Añade una métrica con la antigüedad del registro pendiente más antiguo y una alerta si supera 60 s.
Ejercicio 4 · Resiliencia demostrable (2 h)
- Configura Resilience4j con circuit breaker, retry con jitter, bulkhead y timeout para un dependiente.
- Con WireMock, escribe los cuatro tests de la sección 5.11: lentitud, apertura del circuito, ausencia de llamadas con el circuito abierto y reintento selectivo.
- Comprueba en
/actuator/circuitbreakersy en/actuator/prometheusque las transiciones se reflejan en métricas. - Provoca a propósito el error de contar los 404 como fallo y observa cómo se abre el circuito indebidamente. Después corrígelo con
ignore-exceptions. - Documenta en cinco líneas la degradación elegida y por qué es segura para el negocio.
Ejercicio 5 · Saga con compensación y timeout (3 h)
- Implementa una saga orquestada de tres pasos: cobrar, reservar stock y crear envío.
- Persiste el estado en la tabla
saga_pedidocon bloqueo optimista. - Implementa las compensaciones en orden inverso y hazlas idempotentes.
- Añade el vigilante que compensa las sagas atascadas más de 5 minutos.
- Escribe tests para: camino feliz, fallo en el paso 2 con compensación del 1, y evento duplicado que no debe avanzar la saga dos veces.
- Añade una métrica gauge con el número de sagas no terminales por estado.
Ejercicio 6 · Observabilidad completa (3 h)
- Levanta con Docker Compose: dos servicios, Kafka, OpenTelemetry Collector, Prometheus, Tempo, Loki y Grafana.
- Configura logs JSON con
traceIdy trazas con propagación W3C, incluida la que viaja en las cabeceras de Kafka. - Comprueba que una petición que atraviesa HTTP y luego Kafka mantiene el mismo traceId.
- Instrumenta una métrica de negocio y crea un panel con el método RED.
- Define un SLO de disponibilidad, calcula el presupuesto de error y escribe la alerta multiventana.
- Introduce un fallo (una consulta lenta) y cronometra cuánto tardas en encontrar la causa usando solo tus paneles.
Criterio de éxito: localizas la causa en menos de 5 minutos partiendo únicamente de la alerta.
Ejercicio 7 · Extraer un servicio de un monolito (4 h)
- Parte de un monolito con módulos de pedidos, facturación y notificaciones.
- Elige notificaciones (periférico y sin estado) y extráelo con strangler fig.
- Interpone un gateway y enruta primero el 5 % del tráfico, con posibilidad de volver atrás.
- Sustituye la llamada en memoria por un evento y añade la capa anticorrupción necesaria.
- Escribe el guion de rollback y pruébalo de verdad.
- Documenta un ADR con la decisión, las alternativas descartadas y las consecuencias asumidas.
16 · Resumen del módulo
Las quince ideas que debes recordar
- Los microservicios resuelven un problema organizativo, no técnico. Sin varios equipos autónomos, solo pagas el coste.
- Monolito modular primero. Es la opción que conserva más futuro al menor coste presente, y se puede dividir después.
- La ley de Conway no se negocia. Reorganiza los equipos antes de reorganizar el código.
- Los contextos delimitados marcan las fronteras. Cuando una palabra significa dos cosas, has encontrado una.
- Un agregado, una transacción. Esa frase determina el tamaño de tus servicios y si vas a necesitar sagas.
- Las dependencias apuntan al dominio. Si necesitas Spring para probar tu lógica de negocio, la arquitectura está mal.
- Síncrono solo lo imprescindible. Cada salto suma latencia y multiplica la probabilidad de fallo.
- Toda llamada remota lleva timeout, y los timeouts decrecen hacia abajo en la cadena.
- Reintenta con jitter, en un solo nivel y solo si es idempotente. Si no, los reintentos causan la caída.
- Una base de datos por servicio. Compartir tablas es el acoplamiento más caro que existe.
- Outbox para publicar eventos. El dual write falla en silencio y te enteras semanas después.
- At-least-once es la realidad. Diseña consumidores idempotentes y deja de perseguir el exactly-once.
- Sin trazas distribuidas estás depurando a ciegas. El
traceIden los logs es la mejor inversión del módulo. - Cuidado con la cardinalidad y con promediar percentiles. Un panel que miente es peor que no tener panel.
- SLO y presupuesto de error convierten «features contra estabilidad» en un dato objetivo.