Plan de estudio Java 2026
Módulo 08 crítico Días 15–17 y 20 ≈ 9 h de estudio

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.

Progreso del módulo0 / 0
Cómo leer este módulo: las secciones 1–3 son de diseño y conviene leerlas seguidas, porque la decisión de dividir un sistema depende del dominio, no de la tecnología. Las secciones 4–9 son el oficio: comunicación, resiliencia, datos, Kafka, Spring Cloud y observabilidad; se pueden atacar por separado y cada una tiene código listo para ejecutar. Las 10–16 cierran con seguridad, operación, casos reales, errores, preguntas de entrevista y ejercicios. Si solo tienes una tarde, lee la 1, la 5 y la 9: son las tres que más dinero ahorran en producción.
Aviso de honestidad intelectual: la mayor parte de la literatura sobre microservicios está escrita por empresas con cientos de ingenieros y presupuestos de plataforma de siete cifras. Si tu equipo tiene menos de 20 personas, casi todo lo que hay aquí lo usarás dentro de un monolito modular, y estará bien. Aprende los patrones; sé muy conservador al desplegarlos como procesos separados.

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í
EstiloQué es exactamenteUnidad de despliegueDatos
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
Analogía: un monolito es un piso grande compartido: barato, todo a mano, pero si alguien pone música a las 3 de la mañana lo sufren todos. Los microservicios son un bloque de apartamentos independientes: cada uno reforma el suyo cuando quiere, pero ahora hace falta ascensor, portero, contadores separados, normativa de comunidad y, cuando se va la luz de un piso, alguien tiene que averiguar en qué cuadro está el diferencial. El monolito modular es un piso grande con tabiques y puertas de verdad: casi todas las ventajas de convivencia sin el coste del edificio.

1.2 Tabla honesta de ventajas y desventajas

CriterioMonolitoMonolito modularMicroserviciosServerless
Tiempo hasta el primer despliegueHorasDíasSemanas (plataforma incluida)Horas si el proveedor ya está
Coste cognitivo de una featureBajoBajoAlto: varios repos, contratos, versionesMedio
Refactor entre fronterasTrivial (lo hace el IDE)Fácil (compilador + ArchUnit)Caro: cambio de contrato y despliegue coordinadoCaro
Transacciones ACIDSí, gratisSí, gratisNo: saga y consistencia eventualNo
Depuración de un falloUn stack traceIgualTrazas distribuidas obligatoriasMuy difícil sin observabilidad
EscaladoTodo o nadaTodo o nadaPor servicio, muy finoAutomático y a cero
Aislamiento de fallosNulo: una fuga tumba todoBajoAlto si hay resiliencia; si no, peor que el monolitoAlto
Despliegue independiente por equipoNoNo (mismo binario)Sí — su razón de ser
Heterogeneidad tecnológicaNoNoSí (y suele ser una trampa)
Latencia internaNanosegundos (llamada de método)NanosegundosMilisegundos por salto, y se acumulanMilisegundos más arranques en frío
Coste de infraestructuraBajoBajoAlto: réplicas, malla, observabilidad, brokersBarato a ráfagas, caro sostenido
Personas mínimas para operarlo bien2–33–8~20+ con alguien dedicado a plataforma3–8 y dependencia del proveedor
Onboarding de un junior1 semana1 semana1–2 mesesVariable
La fila que nadie mira y decide el proyecto: «personas mínimas para operarlo bien». Los microservicios no fracasan por razones técnicas: fracasan porque un equipo de seis personas acaba manteniendo catorce despliegues, catorce pipelines, catorce paneles y catorce fuentes de guardia. La arquitectura debe caber en la organización, no al revés.

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.

Pregunta de entrevista senior: «¿cómo decidirías el número de microservicios?». La respuesta mediocre habla de líneas de código. La buena empieza así: «no lo decido desde la tecnología; lo determinan los contextos delimitados del dominio y el número de equipos que pueden ser dueños autónomos de un servicio, incluida su guardia».

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.

CosteEn qué se traduceMitigació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,9991099,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/
    }
}
El argumento decisivo: un monolito modular se puede convertir en microservicios cuando de verdad haga falta, porque las fronteras ya existen y están probadas. Un monolito con espagueti no se puede convertir en nada, y un conjunto de microservicios mal cortados no se puede volver a juntar sin un proyecto de un año. El monolito modular es la opción que conserva más opciones futuras al menor coste presente.

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.

CriterioSeñal objetiva y medibleEjemplo 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.
Criterios que NO justifican dividir: «así lo hace Netflix», «para poder usar Go en un trozo», «el monolito tiene 200.000 líneas», «queremos que el currículum quede bien», «lo pide la nube», «para que el build sea más rápido» (eso se arregla con módulos y caché de compilación) y «para reutilizar código» (eso es una librería, no un servicio).

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:

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:

  1. 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.
  2. Capacidades con datos propios y poco acoplamiento: catálogo, búsqueda, gestión documental.
  3. Capacidades con escalado dispar: recomendaciones, procesamiento de imágenes, importaciones masivas.
  4. 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 ubicuoCoste realCorrecció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:

Regla del contexto delimitado: cuando la misma palabra significa cosas distintas para distintas personas, no busques «el modelo correcto». Has encontrado la frontera entre dos contextos. Modela dos conceptos y traduce entre ellos.
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 subdominioQué esEstrategiaEjemplo 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 soporteNecesario pero no diferencial.Desarrollo propio sencillo, sin sobreingeniería.Gestión de devoluciones, catálogo.
GenéricoResuelto igual en toda la industria.Comprar, no construir.Facturación fiscal, envío de emails, autenticación, pasarela de pago.
El error más caro de la industria: gastar dos años construyendo un subdominio genérico (un CRM propio, un motor de facturación propio, un sistema de identidad propio) mientras el subdominio núcleo se queda sin gente. Antes de escribir la primera línea, clasifica cada contexto en núcleo, soporte o genérico, y defiende esa clasificación delante de negocio.

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ónSignificadoCuándo usarlo
PartnershipDos 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 KernelUn trozo de modelo compartido en código.Solo para tipos verdaderamente universales e inmutables (Dinero, Cif). Cada línea compartida es acoplamiento permanente.
Customer / SupplierEl de aguas abajo es cliente y puede negociar prioridades con el de arriba.La relación sana por defecto dentro de una misma empresa.
ConformistEl 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 LayerTraducción defensiva en la frontera.Sistemas legados, SaaS externos y todo lo que no controlas.
Open Host ServiceEl de arriba publica una API pensada para muchos consumidores.Contextos con muchos clientes: catálogo, identidad.
Published LanguageUn 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 WaysNo 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:

Un agregado = una transacción = un bloqueo. Todo lo que deba ser consistente de forma inmediata tiene que estar dentro del mismo agregado. Lo demás se coordina con consistencia eventual mediante eventos de dominio. Esta única frase determina el tamaño de tus servicios, tu rendimiento bajo concurrencia y si vas a necesitar sagas.

Las cuatro reglas de diseño de agregados (Vaughn Vernon):

  1. 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.
  2. Diseña agregados pequeños. Un agregado grande (un Cliente con sus 10.000 pedidos) provoca cargas enormes, bloqueos largos y conflictos de concurrencia constantes.
  3. Referencia a otros agregados solo por identidad (ClienteId, no Cliente). Esto rompe el grafo de objetos y hace posible dividir después.
  4. 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

EntidadObjeto de valor (value object)
IdentidadTiene identidad propia que persiste aunque cambien todos sus atributos.Se define solo por sus atributos: dos con los mismos valores son el mismo.
MutabilidadSu estado evoluciona en el tiempo.Inmutable siempre. «Cambiar» es crear otro.
IgualdadPor identificador.Por todos los componentes (record lo da gratis).
EjemplosPedido, Cliente, Poliza.Dinero, Email, Direccion, Periodo, PedidoId.
Dónde van las validacionesEn los métodos que cambian el estado.En el constructor: es imposible construir uno inválido.
Consejo de alto impacto y bajo coste: sustituye los tipos primitivos del dominio por objetos de valor (la obsesión por primitivos). Cambiar 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.

ComandoEventoConsulta
Intención«Haz esto»«Esto ha pasado»«Dime esto»
NombreImperativo: ConfirmarPedidoPasado: PedidoConfirmadoObtenerPedido
DestinatariosExactamente unoCero o muchosUno
¿Se puede rechazar?Sí (validación)No, ya ocurrió
AcoplamientoEl emisor conoce al receptorEl emisor no conoce a los receptoresDirecto

2.7 Servicios de dominio, de aplicación y repositorios

ElementoResponsabilidadQué 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. : 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) { }
Distinción crítica que casi todo el mundo confunde: el evento de dominio (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).
Prueba del algodón de una arquitectura hexagonal: ¿puedes borrar el paquete 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
Detalle que marca la diferencia: hay dos modelos de pedido, 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.java con 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.
Regla pragmática: aplica la hexagonal completa en los subdominios núcleo; usa CRUD directo con Spring Data en los de soporte; y compra los genéricos. Mezclar estilos dentro de un mismo sistema no es incoherencia: es economía.

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.

CriterioSíncrono (petición/respuesta)Asíncrono (mensajes y eventos)
Acoplamiento temporalAlto: los dos deben estar vivos a la vezNulo: el receptor puede estar caído y recuperar después
Disponibilidad resultanteSe multiplica: 0,999 × 0,999 × …Se aísla: el broker desacopla
Latencia percibidaLa suma de toda la cadenaRespuesta inmediata («aceptado»); el trabajo va detrás
ConsistenciaInmediata dentro del servicioEventual: hay una ventana de incoherencia
Complejidad de depuraciónMedia: la traza es una cascadaAlta: hay que correlacionar, ordenar y deduplicar
Manejo de erroresEl error vuelve al llamante al instanteReintentos, DLQ y alertas; el emisor ya se fue
ContrapresiónMala: el llamante satura al llamadoNatural: 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.
Heurística práctica: haz síncrono solo lo que el usuario necesita para ver una respuesta correcta en pantalla; todo lo demás, por eventos. Y si dudas, pregúntate: «¿qué pasa si este servicio tarda 30 segundos en enterarse?». Si la respuesta es «nada grave», es asíncrono.

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 versionadoEjemploProsContras
En la ruta (la más usada)/api/v1/pedidosExplícito, cacheable, trivial de enrutar en el gatewayDuplica rutas; los puristas REST protestan
Cabecera de tipo de medioAccept: application/vnd.ejemplo.pedido.v2+jsonURLs estables, versionado por recursoInvisible en el navegador; más difícil de probar y cachear
Parámetro de consulta/pedidos?version=2SimpleSe olvida y ensucia la caché
Sin versión, solo evolución compatible/api/pedidosLo ideal: nunca rompesExige 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";
}
AspectoREST + JSONgRPC + Protobuf
Tamaño del mensajeReferencia (texto, con nombres de campo repetidos)3–10× menor (binario, campos por número)
CPU de serializaciónAlta (parseo de texto)Mucho menor
Latencia típica intra-clúster~2–5 ms~0,5–2 ms (multiplexado HTTP/2, conexión persistente)
ContratoOpenAPI opcional, a menudo desactualizado.proto obligatorio; genera cliente y servidor
StreamingSSE o WebSocket, añadido aparteNativo en cuatro modos
Depuración manualTrivial (curl, navegador)Necesita grpcurl; el binario no se lee
NavegadoresDirectoRequiere gRPC-Web y un proxy
Balanceo de cargaSencillo (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é HTTPSí, estándarNo
Curva de aprendizajeMínima; todo el mundo lo conoceMedia; 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
El problema N+1 de GraphQL y su solución obligatoria: una consulta que pide 100 pedidos y, de cada uno, su cliente, dispara 1 + 100 llamadas al servicio de clientes si los resolvers son ingenuos. La solución es un DataLoader: acumula las peticiones de un mismo tick, las agrupa en una sola llamada por lotes y reparte los resultados. Sin DataLoader, GraphQL garantiza un incidente de rendimiento.
// 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:

4.5 WebSockets y Server-Sent Events

SSE (text/event-stream)WebSocket
DirecciónSolo servidor → clienteBidireccional
ProtocoloHTTP normal: funciona con proxies, CDN y HTTP/2Upgrade a ws://: infraestructura específica
ReconexiónAutomática en el navegador, con Last-Event-IDManual: la implementas tú
FormatoTextoTexto y binario
Casos típicosNotificaciones, progreso de una tarea, precios, tokens de un LLM, feedsChat, colaboración en vivo, juegos, trading, terminales
El problema de escalar conexiones persistentes: son con estado. Si el usuario A está conectado a la réplica 1 y el evento que le interesa lo genera la réplica 3, no se entera. Soluciones: (a) un backplane de publicación/suscripción (Redis Pub/Sub, Kafka, NATS) al que todas las réplicas se suscriben; (b) sesiones pegajosas en el balanceador, frágil ante despliegues; (c) un servicio dedicado de fan-out. Además, cada conexión abierta ocupa un descriptor de fichero y memoria: dimensiona 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

CriterioREST/JSONgRPCMensajería (Kafka/AMQP)GraphQL
ModeloPetición/respuestaPetición/respuesta y streamingPublicación/suscripción, fire-and-forgetPetición/respuesta con consulta declarativa
Acoplamiento temporalAltoAltoNingunoAlto
RendimientoMedioAltoMuy alto en volumen agregadoMedio (depende de los resolvers)
ContratoOpenAPI (opcional).proto (obligatorio)Esquema Avro/JSON con registroSDL (obligatorio)
EvoluciónDisciplina manualReglas de protobufCompatibilidad del registro de esquemasCampos @deprecated
DepuraciónMuy fácilMediaDifícil (asíncrono)Media
Entrega garantizadaNo (se reintenta)No, con persistencia y reintentosNo
Idóneo paraAPI públicas, integración, CRUDTráfico interno de alto volumen, streamingEventos de negocio, integración desacoplada, ETLBFF para móvil y web con muchas vistas
Peor usoChatty APIs con 20 llamadas por pantallaAPI pública para terceros y navegadoresCuando el usuario necesita el resultado yaComunicación servicio a servicio

4.8 Formatos y compatibilidad de esquemas

FormatoTamañoEsquemaLegibleCuándo
JSONGrandeOpcional (JSON Schema)API públicas, depuración y volúmenes moderados. Por defecto, salvo que midas un problema.
ProtobufMuy pequeñoObligatorio, en el fichero .protoNogRPC, tráfico interno intenso, contratos con generación de código.
AvroMuy pequeñoObligatorio, viaja por referencia a un registroNoKafka a gran escala y data lakes: el registro de esquemas hace cumplir la compatibilidad.
MessagePack / CBORPequeñoNoNoJSON 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
ModoQué garantizaCambios permitidosOrden 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.
NONENada: 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.

#FalaciaRealidadQué debes hacer
1La red es fiableSe pierden paquetes, se caen switches, hay particiones de red y reinicios de pods a diario.Reintentos con backoff, idempotencia, circuit breakers.
2La latencia es cero1 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.
3El ancho de banda es infinitoDevolver 50 MB de JSON satura el enlace y a los clientes móviles.Paginación, compresión, proyecciones, formatos binarios.
4La red es seguraDentro del clúster también hay atacantes y errores de configuración.mTLS, confianza cero, cifrado en tránsito y en reposo.
5La topología no cambiaLos 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.
6Hay un solo administradorCinco equipos, tres proveedores cloud y un SaaS que actualiza sin avisarte.Contratos explícitos, versionado, ACL, degradación.
7El coste de transporte es ceroSerializar y deserializar consume CPU; el tráfico entre zonas se factura.Medir el coste por salto, colocalidad de zona, menos conversación innecesaria.
8La red es homogéneaConviven 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

Regla absoluta: toda llamada remota lleva timeout. Sin excepciones. Una llamada sin timeout es una fuga de hilos esperando a suceder: cuando el dependiente se cuelga (no cae — se cuelga, que es peor), tus hilos se agotan uno a uno y tu servicio muere aunque su código esté perfecto. Y ojo: los valores por defecto de muchos clientes HTTP son «infinito».
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

Antes de reintentar, responde a esto: ¿la operación es idempotente? Si no lo es, reintentar un 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

AlgoritmoCómo funcionaRáfagasUso 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. , 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 fijaN peticiones por minuto natural.Sí, mal: 2N en la frontera entre ventanasSimple, aceptable para cuotas groseras.
Ventana deslizanteCuenta las peticiones de los últimos 60 s reales.ControladasMá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.
Respuesta correcta cuando limitas: 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:

  1. Valor cacheado, aunque esté algo obsoleto (stale-while-revalidate). Casi siempre es mejor un precio de hace 5 minutos que un error.
  2. Valor por defecto seguro: sin recomendaciones, muestra los más vendidos; sin cálculo de gastos de envío, aplica la tarifa estándar.
  3. Funcionalidad reducida: oculta la sección, muestra un aviso discreto, deshabilita el botón.
  4. Diferir: acepta la petición (202 Accepted) y procésala cuando el dependiente vuelva.
  5. Error explícito y honesto, con traceId y sin exponer detalles internos. Siempre mejor que colgarse 60 segundos.
Fallbacks que NO debes hacer: devolver 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/42Solo lee.
PUT /pedidos/42 {estado:"ENVIADO"}Fija un estado absoluto.
DELETE /pedidos/42La segunda vez ya no está: devuelve 204 o 404, pero el efecto es el mismo.
POST /pedidosNoCrea un recurso nuevo cada vez → pedidos duplicados.
PATCH /cuenta {saldo: -50}NoEs un incremento relativo: dos veces resta 100.
POST /pagosNo, y aquí dueleDoble 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;
        }
    }
}
Idempotencia sin tabla adicional (cuando encaja): a veces basta con una restricción 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();
Trampa clásica con @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.
La objeción habitual y su respuesta: «pero entonces no puedo hacer un 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.

Compensar no es hacer rollback: el rollback borra la historia; la compensación añade un hecho nuevo que contrarresta al anterior. No «deshaces» un cobro: emites un reembolso, y ambos quedan en el extracto. Esto tiene consecuencias de negocio y contables que hay que hablar con negocio, no decidir en el código.
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).
CriterioCoreografiadaOrquestada
Nº de pasos recomendado2–44 o más
AcoplamientoMínimoEl orquestador conoce a todos
Visibilidad del estadoBaja: hay que reconstruirla con trazasAlta: una fila en una tabla
RiesgoCiclos de eventos, lógica difusaOrquestador anémico o «dios»
Timeouts por pasoComplicadosTriviales
HerramientasKafka y listenersSpring 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;
    }
}
Alternativas al inbox, cuando encajan (y son más baratas): (1) hacer la operación naturalmente idempotente: un 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.
Cuándo NO hacer event sourcing: en un CRUD, en un subdominio de soporte, si el negocio no pide auditoría ni historia, si el equipo no lo ha usado antes, o si lo aplicas «a todo el sistema». La regla sana es: event sourcing solo en el agregado núcleo que lo justifica (la cuenta, la póliza, el expediente) y modelo clásico en todo lo demás. Aplicarlo en todas partes es la vía más rápida a un sistema que nadie sabe operar.

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

PropiedadValor recomendadoPor qué
acksall0 = «dispara y olvida» (pérdida garantizada); 1 = solo el líder (pierdes si cae antes de replicar); all = confirmado por los ISR.
enable.idempotencetrue (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.connection5 o menosCon idempotencia activada, hasta 5 mantiene el orden. Sin idempotencia, más de 1 puede reordenar mensajes al reintentar.
delivery.timeout.ms120000Es el que manda de verdad: tiempo total para dar un mensaje por entregado, reintentos incluidos.
linger.ms5–50Esperar unos milisegundos para agrupar mensajes en lotes multiplica el rendimiento. 0 minimiza latencia y desperdicia red.
batch.size32768–131072Tamaño del lote por partición.
compression.typezstd (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 instanciaSolo si necesitas transacciones: escritura atómica en varios topics más el commit de offsets.
La decisión más importante del productor es la CLAVE. Determina la partición y, por tanto, el orden y el reparto de carga. Usa el identificador del agregado (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íaCómo se consigueRiesgoCuá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.
El matiz de «exactly-once» que hay que saber en una entrevista: Kafka ofrece semántica exactly-once dentro de Kafka: leer de un topic, procesar, escribir en otro topic y commitear offsets, todo atómicamente. En el momento en que tu efecto secundario sale de Kafka —insertar en PostgreSQL, llamar a una API, enviar un email— vuelves a estar en at-least-once, porque no hay transacción que abarque los dos sistemas (es el mismo problema del dual write de la sección 6.5). La conclusión práctica es siempre la misma: diseña consumidores idempotentes y deja de perseguir el exactly-once.

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

CriterioKafkaRabbitMQSQS / SNSPulsar
ModeloLog distribuido particionadoBroker de colas con enrutamiento (AMQP)Cola y pub-sub gestionadosLog y colas, con almacenamiento separado (BookKeeper)
RendimientoMillones de mensajes/sDecenas de miles/s por colaAlto, elástico y gestionadoComparable a Kafka
Retención y relectura: días o años; se puede reprocesar todoNo: al consumir, desapareceSQS 14 días máx.; sin relecturaSí, con almacenamiento por niveles (S3)
OrdenPor particiónPor cola (con un consumidor)Solo en colas FIFO, con menor caudalPor partición o por clave
EnrutamientoSimple: topic y claveMuy rico: direct, topic, fanout, headersBásico (SNS a varios destinos)Rico
Entrega retrasadaNo nativa (se emula con retry topics)Sí: plugin de delayed exchange, TTL con DLXSí, hasta 15 minSí, nativa
Multi-tenencia y geoCon MirrorMakerFederación y shovelPor regiónNativa: tenants, namespaces, geo-replicación
OperaciónCompleja (mejor con KRaft o gestionado)SencillaCero: es un servicioLa más compleja (broker más BookKeeper)
Úsalo para…Eventos de negocio, streaming, integración a gran escala, reprocesado históricoColas de trabajo, RPC asíncrono, enrutamiento complejo, prioridadesTodo en AWS cuando no quieres operar nadaMulti-tenencia y geo fuertes; menos ecosistema
Evítalo si…Solo necesitas una cola de tareas: es un cañón para una moscaNecesitas relectura o volúmenes enormesNecesitas orden global o latencia muy bajaNo 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).
Recomendación 2026: si estás en Kubernetes, no montes Eureka. Usa el DNS del clúster y llama a 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 ServerConfigMap y Secret de Kubernetes
OrigenRepositorio Git (auditable, con historial y PR)Manifiestos YAML, idealmente en Git con GitOps
Recarga en calienteSí, con @RefreshScope y /actuator/refresh o Spring Cloud BusLos ficheros montados se actualizan solos; las variables de entorno no (hay que reiniciar el pod)
SecretosCifrado propio o integración con VaultSecret es solo base64: usa Sealed Secrets, External Secrets o Vault
Coste operativoUn servicio más (y en alta disponibilidad)Cero: viene con la plataforma
CuándoFuera de Kubernetes, o si necesitas refresco sin reiniciarPor defecto en Kubernetes
Peligro del refresco en caliente: @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 correlacionartraceId y spanId en el MDC de cada línea
Concatenación con +: construye el String aunque el nivel esté apagadoMarcadores {}: SLF4J solo formatea si va a emitir
Todo a nivel INFO, o todo a DEBUG en producciónNiveles con criterio; DEBUG activable por paquete sin reiniciar
Loguear el DNI, el email, el token o la tarjetaEnmascarar o no loguear datos personales ni secretos
log.error("error", e) y además relanzarLoguear o relanzar, no las dos cosas (evita el log duplicado)
Loguear dentro de un bucle de 10.000 elementosLoguear 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) { /* ... */ }
Percentiles: el error estadístico más caro. No se pueden promediar. Si tienes 3 pods con p99 de 100, 200 y 900 ms, el p99 del servicio no es 400 ms: puede ser cualquier cosa. Para agregar correctamente entre instancias necesitas histogramas (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

PiezaOpción libreAlternativasQué hace
MétricasPrometheus (modelo pull)VictoriaMetrics, Mimir, DatadogRecolecta y almacena series temporales; lenguaje PromQL.
PanelesGrafanaDatadog, New Relic, KibanaVisualización unificada de métricas, logs y trazas.
LogsLoki (indexa etiquetas, no contenido: barato)Elasticsearch/OpenSearch, DatadogAgregación y búsqueda de logs.
TrazasTempo o JaegerZipkin, Datadog APMAlmacenamiento y consulta de trazas distribuidas.
RecolecciónOpenTelemetry CollectorAgentes propietariosRecibe, procesa, muestrea y reenvía telemetría. Te desacopla del proveedor.
AlertasAlertmanagerPagerDuty, OpsgenieEnrutamiento, agrupación, silenciado y escalado de alertas.
La decisión estratégica: instrumenta con OpenTelemetry, que es un estándar neutral, y envía todo al Collector. A partir de ahí, cambiar de Datadog a Grafana (o al revés) es modificar la configuración del Collector, no reinstrumentar 40 servicios. Es la decisión que más dinero ahorra a tres años vista, y cuesta lo mismo tomarla bien que mal.
# 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 }
Fatiga de alertas: el fallo cultural que anula toda la inversión. Si el canal de alertas recibe 50 mensajes al día, el equipo deja de mirarlo y el día que suena la buena nadie reacciona. Reglas: (1) toda alerta que despierta a alguien debe ser accionable y tener un runbook; (2) alerta por síntomas del usuario, no por CPU al 80 %; (3) si una alerta se ignora tres veces seguidas, se borra o se corrige; (4) usa ventanas múltiples de consumo de presupuesto en vez de umbrales instantáneos.

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ónCosteRotación de certificadosCuándo
Malla de servicios (Istio, Linkerd)Alto al montarla, nulo despuésAutomática, cada 24 hLo estándar hoy en Kubernetes. mTLS «gratis» y sin tocar código.
cert-manager con SPIFFE/SPIREMedioAutomáticaIdentidad de carga de trabajo sin malla completa.
Manual en la aplicación (keystore/truststore)Alto y permanenteManual: la fuente número uno de caídas por certificado caducadoSolo 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));
    }
}
Vulnerabilidad clásica de microservicios (IDOR): el gateway valida el token y los servicios internos «confían». Entonces alguien descubre que 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

EstrategiaCómo funcionaCosteRollbackRiesgo
Rolling update (por defecto en K8s)Sustituye réplicas poco a poco respetando maxUnavailable y maxSurge.BajoMinutos (rollout deshacer)Conviven dos versiones: los contratos deben ser compatibles.
Blue-greenDos entornos completos; se conmuta el tráfico de golpe.Alto: el doble de recursosInstantáneoMigraciones de BD compartidas entre azul y verde.
Canary1 % → 10 % → 50 % → 100 %, con métricas que deciden en cada paso.MedioRápidoNecesita buenas métricas y automatización (Argo Rollouts, Flagger).
Shadow / espejoSe duplica el tráfico real a la versión nueva sin devolver su respuesta.AltoN/ACuidado 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

#ErrorConsecuenciaSolución
1Empezar con microservicios «porque es lo moderno»Coste enorme y cero beneficio; el equipo se ahoga en operaciónMonolito modular primero; dividir con criterios objetivos
2Cortar por capas técnicas (servicio-controladores, servicio-repositorios)Cada feature toca todos los servicios: monolito distribuidoCortar por contextos delimitados y capacidades de negocio
3Base de datos compartida entre serviciosAcoplamiento oculto imposible de refactorizarUna BD por servicio; integrar por API y eventos
4Llamadas remotas sin timeoutAgotamiento de hilos y caída en cascadaTimeout en toda llamada, decreciente hacia abajo
5Reintentar en varios niveles a la vezAmplificación ×27: los reintentos matan al servicio degradadoReintentar en un solo nivel, con jitter y presupuesto
6Reintentar operaciones no idempotentesCobros y pedidos duplicadosClave de idempotencia o identificador generado por el cliente
7Dual write: guardar en BD y publicar en KafkaPérdida silenciosa de eventos o eventos fantasmaOutbox transaccional (polling o CDC)
8Consumidores no idempotentes con at-least-onceFacturas y emails duplicadosInbox, UPSERT o restricción UNIQUE de negocio
9Circuit breaker que cuenta los 4xx como fallosEl circuito se abre por peticiones correctasignore-exceptions con las excepciones de negocio
10@TimeLimiter sobre un método bloqueanteEl timeout no se aplica; falsa sensación de seguridadTimeout en el cliente HTTP; @TimeLimiter solo con CompletableFuture
11Cadenas síncronas largasLatencia sumada y disponibilidad multiplicadaNúcleo síncrono mínimo; el resto por eventos
12Clave de partición con poca variedadParticiones calientes; imposible escalarClave con alta cardinalidad y distribución uniforme
13Más consumidores que particionesConsumidores ociosos; el escalado no hace nadaDimensionar particiones con holgura desde el principio
14Sin dead letter topic ni alertaMensajes envenenados en bucle infinito; lag crecienteDLT, ErrorHandlingDeserializer y alerta si la DLT recibe algo
15Saga sin timeout ni compensaciónProcesos de negocio colgados para siempreVigilante periódico, estado persistido y métricas de sagas atascadas
16Logs sin traceIdInvestigar un incidente pasa de 5 min a 3 hMicrometer Tracing y JSON estructurado con MDC
17Etiquetas de métricas de alta cardinalidadPrometheus se queda sin memoria; factura disparadaMeterFilter que bloquee etiquetas prohibidas
18Promediar percentiles entre instanciasPaneles que mienten; decisiones erróneasHistogramas y histogram_quantile
19Comprobar la base de datos en el probe de livenessReinicio masivo de pods y bucle de falloDependencias solo en readiness; liveness mínimo
20Migración de esquema destructiva en el mismo despliegueLa versión antigua se rompe durante el rolling updateExpand and contract en cuatro pasos
21Reenviar el JWT del usuario por toda la cadenaUn servicio comprometido suplanta al usuario en todo el sistemaIntercambio de token con audiencia y alcance mínimos
22Confiar en que «el gateway ya autorizó»IDOR: cualquiera lee datos ajenos conociendo el UUIDCada servicio comprueba la propiedad del dato
23Librería común con el modelo de dominio dentroCambiar un campo obliga a redesplegar todoCompartir solo utilidades técnicas; duplicar DTOs a propósito
24Evento de dominio publicado tal cual como evento de integraciónRefactorizar tu dominio rompe a otros equiposTraducir en la frontera; contrato versionado con esquema
25Feature flags que nunca se borranCombinatoria inmanejable y código muertoFecha 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.

  1. Identifica los contextos delimitados y clasifícalos en núcleo, soporte o genérico.
  2. Dibuja el mapa de contextos indicando el patrón de relación de cada arista.
  3. Marca qué contextos serían microservicios desde el día 1 y cuáles empezarían como módulos.
  4. Busca dos palabras que signifiquen cosas distintas en dos contextos y modela ambas.
  5. 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)

  1. Crea un servicio de pedidos con Spring Boot 3 y Java 21, con los paquetes dominio, aplicacion e infraestructura.
  2. Implementa el agregado Pedido con la invariante «total igual a la suma de líneas» y los objetos de valor Dinero, PedidoId y Sku.
  3. Define los puertos RepositorioPedidos y PublicadorEventos, y sus adaptadores JPA y Kafka.
  4. Escribe tests del dominio y del caso de uso sin Spring, con dobles en memoria.
  5. 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)

  1. Crea la tabla outbox con su índice parcial mediante una migración de Flyway.
  2. Implementa PublicadorEventosOutbox escribiendo dentro de la transacción del caso de uso.
  3. Implementa el publicador por polling con FOR UPDATE SKIP LOCKED.
  4. Escribe un test con Testcontainers que arranque PostgreSQL y Kafka y verifique que el evento llega.
  5. Escribe un test que provoque un rollback y compruebe que no se publica nada.
  6. Levanta dos instancias del publicador y comprueba que no duplican ni se bloquean entre ellas.
  7. 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)

  1. Configura Resilience4j con circuit breaker, retry con jitter, bulkhead y timeout para un dependiente.
  2. 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.
  3. Comprueba en /actuator/circuitbreakers y en /actuator/prometheus que las transiciones se reflejan en métricas.
  4. 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.
  5. 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)

  1. Implementa una saga orquestada de tres pasos: cobrar, reservar stock y crear envío.
  2. Persiste el estado en la tabla saga_pedido con bloqueo optimista.
  3. Implementa las compensaciones en orden inverso y hazlas idempotentes.
  4. Añade el vigilante que compensa las sagas atascadas más de 5 minutos.
  5. 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.
  6. Añade una métrica gauge con el número de sagas no terminales por estado.

Ejercicio 6 · Observabilidad completa (3 h)

  1. Levanta con Docker Compose: dos servicios, Kafka, OpenTelemetry Collector, Prometheus, Tempo, Loki y Grafana.
  2. Configura logs JSON con traceId y trazas con propagación W3C, incluida la que viaja en las cabeceras de Kafka.
  3. Comprueba que una petición que atraviesa HTTP y luego Kafka mantiene el mismo traceId.
  4. Instrumenta una métrica de negocio y crea un panel con el método RED.
  5. Define un SLO de disponibilidad, calcula el presupuesto de error y escribe la alerta multiventana.
  6. 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)

  1. Parte de un monolito con módulos de pedidos, facturación y notificaciones.
  2. Elige notificaciones (periférico y sin estado) y extráelo con strangler fig.
  3. Interpone un gateway y enruta primero el 5 % del tráfico, con posibilidad de volver atrás.
  4. Sustituye la llamada en memoria por un evento y añade la capa anticorrupción necesaria.
  5. Escribe el guion de rollback y pruébalo de verdad.
  6. 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

  1. Los microservicios resuelven un problema organizativo, no técnico. Sin varios equipos autónomos, solo pagas el coste.
  2. Monolito modular primero. Es la opción que conserva más futuro al menor coste presente, y se puede dividir después.
  3. La ley de Conway no se negocia. Reorganiza los equipos antes de reorganizar el código.
  4. Los contextos delimitados marcan las fronteras. Cuando una palabra significa dos cosas, has encontrado una.
  5. Un agregado, una transacción. Esa frase determina el tamaño de tus servicios y si vas a necesitar sagas.
  6. Las dependencias apuntan al dominio. Si necesitas Spring para probar tu lógica de negocio, la arquitectura está mal.
  7. Síncrono solo lo imprescindible. Cada salto suma latencia y multiplica la probabilidad de fallo.
  8. Toda llamada remota lleva timeout, y los timeouts decrecen hacia abajo en la cadena.
  9. Reintenta con jitter, en un solo nivel y solo si es idempotente. Si no, los reintentos causan la caída.
  10. Una base de datos por servicio. Compartir tablas es el acoplamiento más caro que existe.
  11. Outbox para publicar eventos. El dual write falla en silencio y te enteras semanas después.
  12. At-least-once es la realidad. Diseña consumidores idempotentes y deja de perseguir el exactly-once.
  13. Sin trazas distribuidas estás depurando a ciegas. El traceId en los logs es la mejor inversión del módulo.
  14. Cuidado con la cardinalidad y con promediar percentiles. Un panel que miente es peor que no tener panel.
  15. SLO y presupuesto de error convierten «features contra estabilidad» en un dato objetivo.
Si solo te llevas una frase: los microservicios no son una meta, son una herramienta cara para comprar autonomía organizativa. Todo lo que has aprendido aquí —DDD, hexagonal, resiliencia, outbox, idempotencia, observabilidad— te hará mejor ingeniero aunque nunca despliegues un microservicio, porque son técnicas para dominar la complejidad, y esa la tienes en cualquier sistema.

Checklist final del módulo 08

Continúa por aquí: el módulo 09 cubre cómo se despliega y opera todo esto (Docker, Kubernetes, CI/CD y cloud); el módulo 07 desarrolla el contract testing y Testcontainers; el módulo 10 profundiza en OAuth2, JWT y mTLS; y el módulo 03 explica los virtual threads, que cambian el dimensionamiento de todo lo que has visto en la sección 5.