Reestructuración de Paquetes para Microservicios Asíncronos con RocketMQ

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 contra cache.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. Utiliza MessageModel.BROADCASTING para garantizar que cada instancia del servicio reciba actualizaciones de configuración. Al capturar un mensaje, se sobrescribe el objeto singleton en cache.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 en EventoAlertaContextual.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_transaccion que 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.java se 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.java prioriza 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      │        │          │
│  └─────────────────┘                         │        │          │
└───────────────────────────────────────────────┴────────┴──────────┘
  1. Fase de Ingreso: El broker entrega lotes de métricas a suscriptores/, que los transforma en objetos ligeros de contrato.
  2. 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.
  3. 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.
  4. 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 en emisores/. 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.

Etiquetas: RocketMQ arquitectura-microservicios patrones-mensajeria comunicacion-asincrona spring-boot

Publicado el 10-10 22:52