Mensajería, colas y pub/sub | Nicolás Garzón
Colas, pub/sub y logs distribuidos resuelven necesidades diferentes. Elegir un broker no define automáticamente entrega, orden ni procesamiento seguro: esas garantías dependen del contrato, configuración y consumidor.
La mensajería desacopla productores y consumidores en el tiempo y permite absorber diferencias de velocidad.
Texto
Copiar productor → broker → consumidorEl broker almacena o enruta mensajes; el consumidor sigue siendo responsable de interpretar y aplicar efectos correctamente.
Varios workers compiten y normalmente uno procesa cada mensaje.
Texto
Copiar jobs → [queue] → worker A | worker B | worker CÚtil para emails, procesamiento de archivos o tareas diferidas.
Cada suscriptor obtiene una copia lógica.
Texto
Copiar OrderConfirmed
├── Notifications subscription
├── Reporting subscription
└── Delivery subscriptionÚtil para propagar hechos a capacidades independientes.
Conserva registros ordenados por partición y permite que consumidores manejen offsets y replay.
Útil para streams, proyecciones y alto throughput, pero exige gestión de particiones y retención.
Productor serializa y publica.
Broker confirma según su garantía.
Mensaje se almacena o enruta.
Consumidor recibe.
Procesa en su sistema local.
Hace ack o commit de offset.
Broker retira o conserva según modelo.
El punto crítico es el intervalo entre aplicar el efecto y confirmar. Un fallo puede producir redelivery.
Se confirma antes de procesar o no se reintenta. Puede perder mensajes, evita duplicación de entrega.
Se reentrega hasta obtener confirmación. Reduce pérdida, permite duplicados.
Suele existir únicamente dentro de una frontera concreta del broker o framework. Un email, cobro o escritura externa todavía necesita idempotencia.
En una cola, un mensaje puede ocultarse durante procesamiento. Si el worker no confirma antes del visibility timeout, reaparece.
El timeout debe superar procesamiento normal, pero no ser tan largo que retrase recuperación. Para trabajos extensos pueden usarse heartbeats o dividir la tarea.
Un log suele garantizar orden dentro de una partición.
Texto
Copiar partition key = orderIdTodos los eventos de un pedido llegan al mismo orden relativo. Aumentar particiones mejora paralelismo, pero cambia distribución y puede complicar rebalances.
Orden global limita throughput y disponibilidad. Solicítalo solo si el dominio realmente lo necesita.
Consumidores de un mismo grupo comparten particiones. Si hay más consumers que particiones, algunos quedan inactivos. El escalado depende de particiones y costo por mensaje.
Un rebalance puede pausar procesamiento y causar redelivery. Los handlers deben tolerarlo.
Si la tasa de llegada supera capacidad, crece lag.
Texto
Copiar arrival rate > processing rate
→ queue depth/lag aumenta
→ latencia crece
→ retención o almacenamiento se agota
Escalar consumidores.
Limitar productores.
Priorizar trabajo.
Aplicar batching.
Degradar tareas secundarias.
Rechazar antes del colapso.
Autoscaling tardío no corrige un consumidor bloqueado por dependencia externa.
Retries inmediatos pueden bloquear throughput. Estrategias:
Backoff con jitter.
Retry topics o colas diferidas.
Máximo de intentos.
Clasificación transitorio/permanente.
DLQ después del límite.
Un error de schema no mejorará reintentando cien veces.
Una DLQ no es un cementerio olvidado. Debe conservar:
Mensaje original.
Motivo.
Intentos.
Timestamps.
Consumer y versión.
Correlation ID.
Necesita alertas, herramienta de inspección, corrección y replay controlado.
Un mensaje que siempre falla puede bloquear una partición si se procesa estrictamente en orden. Opciones:
Mover a DLQ.
Saltar con registro de hueco.
Pausar partición e intervenir.
La elección depende de si mensajes posteriores pueden procesarse sin el anterior.
Mensajes grandes aumentan latencia, costo y presión de memoria. Puede publicarse una referencia, pero eso crea dependencia de almacenamiento y permisos.
No uses el broker como transporte de archivos sin evaluar límites.
Cambios aditivos.
Valores desconocidos tolerables.
Schema registry.
Versiones.
Upcasters o adapters.
Contract tests.
La compatibilidad depende de productor y consumidores desplegados simultáneamente.
Autenticación de productores y consumidores.
Autorización por topic/queue.
Cifrado.
Datos sensibles mínimos.
Retención y borrado.
Protección contra mensajes falsificados.
Un tenant ID dentro del payload no sustituye autorización del consumidor.
Publish rate y errores.
Queue depth o consumer lag.
Processing latency.
Redelivery.
Retry y DLQ rate.
Edad del mensaje más antiguo.
Saturación del consumer.
La métrica más importante suele ser tiempo desde hecho hasta efecto, no solo mensajes por segundo.
Texto
Copiar archivo cargado
→ ImportRequested
→ worker valida lote
→ mensajes por producto
→ consumidores actualizan catálogo
→ ImportCompleted
Idempotency key por archivo.
Partición por tenant.
Batch size.
Estado de progreso.
DLQ para filas inválidas.
No bloquear todo por un producto erróneo.
Reanudación desde checkpoint.
Productor usa outbox o falla explícitamente. No pierdas el hecho después de commit.
Redelivery; inbox/idempotencia evita repetir efecto.
La deduplicación por event ID no ayuda si el productor generó dos eventos para la misma intención. Necesita idempotencia en origen.
Revisa capacidad, dependencia lenta, partición caliente y retención. Escalar puede no bastar.
Separa proyecciones reconstruibles de efectos irreversibles.
Suponer orden global.
Confundir entrega con efecto.
Retries ilimitados.
DLQ sin operación.
Particionar por clave caliente.
Payload interno completo.
No medir edad y lag.
Escalar consumers sin particiones.
Cola: distribuir trabajo.
Pub/sub: múltiples reacciones.
Log: replay, streams y alto throughput ordenado por clave.
La elección debe seguir semántica, no solo producto disponible.
At-least-once implica duplicados.
Orden normalmente es por partición.
Backpressure debe diseñarse antes de saturación.
DLQ necesita proceso operativo.
Broker no elimina acoplamiento; lo mueve a schema y temporalidad.
Procesamiento seguro exige idempotencia y observabilidad.
¿Por qué más consumers no siempre aumentan throughput?
¿Qué ocurre si el visibility timeout es demasiado corto?
¿Cuándo usarías un log en lugar de una cola?
¿Qué diferencia existe entre mensaje entregado y efecto aplicado?
¿Cómo manejarías un poison message con orden estricto?
Sagas y transacciones distribuidas coordina procesos que necesitan varias transacciones locales sin prometer atomicidad global.