Transición Arquitectónica: De REST Síncrono a Flujo Dirigido por Eventos
La integración de un broker de mensajería en un ecosistema distribuido rara vez se limita a añadir una dependencia y configurar beans de conexión. En entornos de alto rendimiento, donde se procesan miles de transacciones por segundo, la adopción de RocketMQ obliga a redefinir la topología interna del proyecto, transformando un modelo de petición-respuesta síncrono en un canal de tuberías asíncronas.
Este análisis examina cómo se reorganiza el sistema de archivos de un servicio de monitorización biomédica (monitor-biometrico) al migrar de un controlador HTTP tradicional a un modelo basado en eventos, destacando las implicaciones en la organización del código y la circulación de datos.
Evolución Estructural del Proyecto
La diferencia más visible entre un microservicio convencional y uno impulsado por mensajería radica en la ubicación de los puntos de entrada y salida. A continuación se contrastan ambos enfoques:
| Componente | Microservicio REST (Síncrono) | Microservicio RocketMQ (Asíncrono) | Propósito Arquitectónico |
|---|---|---|---|
| Entrada de Datos | controlador/ (Expone endpoints HTTP) |
suscriptores/ (Consumen mensajes del broker) |
Elimina el bloqueo del hilo principal y permite ingestión masiva sin esperar respuestas inmediatas. |
| Salida de Notificaciones | cliente/ (Invocación directa vía Feign/gRPC) |
emisores/ (Publica payloads en tópicos) |
Desacopla la disponibilidad del servicio consumidor. El fallo del destinatario no interrumpe el flujo principal. |
| Contrato de Transferencia | dto.peticion/ / dto.respuesta/ |
contratos.evento/ |
Aislamiento de la capa de transporte. Se eliminan dependencias de HTTP, headers y cookies, centrándose únicamente en la carga útil serializada. |
Responsabilidades de los Paquetes Nucleares
com.salud.telemetria
├── suscriptores/ <--- [Nuevo] Puerta de entrada principal del flujo
│ ├── ProcesadorFlujoBiometrico.java (Modo cluster: ingesta continua)
│ └── SincronizadorReglas.java (Modo broadcast: actualización de umbrales)
├── emisores/ <--- [Nuevo] Canal de publicación asíncrona
│ └── CanalNotificaciones.java (Dispatch de alertas críticas)
├── contratos/
│ └── evento/ <--- [Refactorizado] Objetos inmutables para serialización
│ ├── MensajeMetricaRaw.java (Estructura ligera para series temporales)
│ └── EventoAlertaContextual.java (Payload enriquecido con contexto)
├── cache.local/ <--- [Colaborador] Heurística en memoria de la JVM
└── motor.calculo/ <--- [Colaborador] Lógica pura sin efectos secundarios de red
1. Paquete suscriptores/: Desplazamiento del Controlador
El 90% del tráfico abandona los endpoints REST y se canaliza directamente hacia los consumidores. La organización interna se dicta por el patrón de distribución:
- Ingesta por Lotes (Clustering): Implementado en
ProcesadorFlujoBiometrico.java. Se configura un tamaño de lote de 150 registros por invocación para reducir la sobrecarga de transaccoines. Los datos se deserializan y se cruzan directamente contracache.local/, evitando consultas SQL. Los resultados se vuelcan de forma asíncrona a bases de datos de series temporales. - Sincronización Global (Broadcasting): Implementado en
SincronizadorReglas.java. UtilizaMessageModel.BROADCASTINGpara garantizar que cada instancia del servicio reciba actualizaciones de configuración. Al capturar un mensaje, se sobrescribe el objeto singleton encache.local/, permitiendo cambios en caliente sin reinicios ni pausas en la ingesta.
2. Paquete emisores/: Publicación con Idempotencia y Contexto
Cuando el motor.calculo/ detecta una anomalía, delega la notificación en CanalNotificaciones.java. La estrategia se basa en dos pilares:
- Payload Enriquecido (Fat Payload): En lugar de enviar solo el identificador del paciente, el ensamblador extrae nombre, ubicación y parámetros vitales actuales desde
cache.local/e inyecta toda la información enEventoAlertaContextual.java. Esto garantiza que el servicio receptor pueda ejecutar acciones inmediatas (SMS, UI) sin realizar llamadas bloqueantes para recuperar datos faltantes. - Token de Idempotencia: Antes de publicar, se genera un
uuid_transaccionque se adjunta a los headers del mensaje. El consumidor downstream utiliza este valor para evitar procesamientos duplicados en caso de reintentos de red.
3. Aislamiento Estricto de contratos.evento/
Reutilizar entidades de base de datos (JPA Entities) como contratos de mensajería introduce acopalmiento temporal y estructural peligroso. La separación es obligatoria porque:
- Las migraciones de esquema no deben romper consumidores activos.
MensajeMetricaRaw.javase diseñó con un footprint mínimo, eliminando metadatos redundantes para optimizar el ancho de banda y la carga de CPU durante la serialización/deserialización masiva.EventoAlertaContextual.javaprioriza la inmediatez sobre el tamaño, incrustando todo el estado necesario para que la cadena de notificaciones opere de forma autónoma.
Topología de Flujo de Datos
┌──────────────────────────────────────┐
│ Apache RocketMQ │
└──────────┬───────────────────▲───────┘
│ │
1. Ingesta │ │ 4. Publicación
Continua │ │ (Contexto Completo)
▼ │
┌───────────────── suscriptores/ ──────────────────┼──────────┐
│ │ │
│ ProcesadorFlujoBiometrico │ │
│ │ │ │
│ │ (Deserializa → contratos.evento) │ │
│ ▼ │ │
│ ┌─────────────────┐ 2. Validación Reglas ┌─────────────┐ │ │
│ │ motor.calculo/ ├──────────────────► cache.local/│ │ │
│ │ (Filtrado/Window)│ │ (Hashmaps) │ │ │
│ └───────┬─────────┘ └──────┬──────┘ │ │
│ │ 3. Evaluación │ │ │
│ ▼ │ Refresco│ │
│ ┌─────────────────┐ ◄───────────────────────┼────────┼──────────┤
│ │ emisores/ │ Canal Broadcast │ │ │
│ │ CanlNotif. │ SincronizadorReglas │ │ │
│ └─────────────────┘ │ │ │
└───────────────────────────────────────────────┴────────┴──────────┘
- Fase de Ingreso: El broker entrega lotes de métricas a
suscriptores/, que los transforma en objetos ligeros de contrato. - Fase de Procesamiento: El
motor.calculo/aplica ventanas deslizantes y filtros de ruido. Toda la validación se ejecuta contra estructuras en memoria (ConcurrentHashMap), manteniendo la latencia por debajo de 2ms sin generar I/O de red. - Fase de Actualización: Los cambios en los umbrales críticos se propagan vía broadcasting, actualizando la caché local y modificando el comportamiento del motor en tiempo real.
- Fase de Salida: Si se superan los límites permitidos,
emisores/construye un mensaje enriquecido y lo publica para su distribución a canales externos.
Criterios de Selección: RocketMQ vs Comunicación RPC Síncrona
La coexistencia de colas de mensajes y clientes HTTP síncronos (como OpenFeign) requiere una delimitación clara de responsabilidades basada en el perfil de latencia y tolerancia a fallos.
| Dimensión | RPC Síncrono (OpenFeign/gRPC) | Cola de Mensajes (RocketMQ) |
|---|---|---|
| Patrón | Punto a punto explícito | Publicación-Suscripción (Topic) |
| Bloqueo de Hilo | Síncrono (espera respuesta) | Asíncrono (fire-and-forget) |
| Tolerancia a Fallos | Acoplado (timeout/circuit breaker necesario) | Desacoplado (persistencia en broker) |
| Gestión de Picos | Sin amortiguación (satura pool de conexiones) | Amortiguador natural (acumulación en cola) |
| Extensibilidad | Baja (modificación de código fuente) | Alta (nuevos suscriptores sin tocar productor) |
Aplicación Práctica en el Servicio
- Uso Obligatorio de RPC: Inicialización del servicio (
IniciadorCache.java). Durante el arranque, el sistema debe recuperar el catálogo completo de pacientes y reglas desde el servicio central. La dependencia es estricta y el hilo no puede avanzar sin los datos, lo que justifica una llamada síncrona con reintentos controlados. - Uso Obligatorio de Mensajería: Disparo de emergencias. El
motor.calculo/detecta bradicardia severa y delega enemisores/. La entrega debe ser no bloqueante para no ralentizar el procesamiento de métricas siguientes. Además, la persistencia del mensaje garantiza que, si el sistema de alertas sufre una caída, las notificaciones se procesen automáticamente al recuperarse, cumpliendo requisitos de alta disponibilidad.
Matriz de Decisión Técnica
- Seleccionar RPC síncrono cuando: Existe dependencia inmediata de datos, se requiere validación de éxito en el acto, y el volumen de peticiones es predecible y estable.
- Seleccionar Broker de Mensajes cuando: Se prioriza el desacoplamiento, se manejan ráfagas impredecibles, se requiere replicación a múltiples consumidores futuros, o la pérdida de un hilo de ejecución no es aceptable en cascada.