pythonarchitecturehexagonalservice-busazureasyncidempotencypatterns

El consumer que perdía mensajes bajo carga — cómo lo arreglé con hexagonal y Service Bus

Tenía un adapter de integración que funcionaba. Procesaba mensajes de una fuente externa, ejecutaba lógica de negocio, y escribía resultados. El problema era que “funcionaba” significaba “funciona la mayoría de las veces” — con mensajes duplicados bajo carga, con pérdida de mensajes en desconexiones, con lógica de negocio entremezclada con código de infraestructura hasta el punto donde era difícil modificar una sin tocar la otra.

Este post es el registro del refactor: cómo pasé de un adapter monolítico a arquitectura hexagonal con Service Bus, y por qué la idempotencia importa más de lo que parece.

El problema con el adapter original

El adapter original tenía todo mezclado: conexión al broker, deserialización de mensajes, lógica de procesamiento, escritura a base de datos, manejo de errores de red. Un archivo de 600+ líneas donde cualquier cambio en el schema del mensaje requería entender todo el flujo completo.

Las fallas que aparecían en producción:

Arquitectura hexagonal: por qué importa aquí

La arquitectura hexagonal (ports & adapters) separa el dominio de negocio de las dependencias de infraestructura. El dominio no sabe que existe Service Bus, ni MongoDB, ni ningún sistema externo. Solo conoce interfaces (ports). La infraestructura implementa esas interfaces (adapters).

irt-adapter/
├── apps/               # Entry points (Service Bus consumer, CLI)
├── domain/             # Lógica de negocio pura
│   ├── models/         # Entidades de dominio
│   ├── ports/          # Interfaces que el dominio requiere
│   └── services/       # Lógica de procesamiento
└── infrastructure/     # Implementaciones concretas
    ├── service_bus/    # Adapter Service Bus
    ├── mongo/          # Adapter MongoDB
    └── http/           # Adapter HTTP clients

El dominio tiene un MergingRepository (port) que define cómo se guardan los resultados del merge. La infraestructura tiene un MongoMergingRepository que implementa ese port. Si mañana migro de MongoDB a PostgreSQL, cambio la implementación del adapter, no toco el dominio.

La prueba de que el dominio está bien aislado: puedo testearlo sin levantar ningún servicio externo, sin Service Bus, sin MongoDB. Solo instancio los adapters con mocks o in-memory implementations.

peek_lock vs receive_and_delete

Esta es la decisión más importante del diseño de un consumer Service Bus.

receive_and_delete: el mensaje se elimina del queue en el momento de la lectura. Semántica at-most-once. Si el proceso falla entre lectura y ACK, el mensaje se pierde.

peek_lock: el mensaje se lockea en el queue pero no se elimina. El consumer tiene un TTL para procesarlo y hacer ACK (complete) o NACK (abandon). Si el consumer cae sin hacer ACK, el lock expira y el mensaje vuelve al queue para ser entregado de nuevo. Semántica at-least-once.

Para cualquier procesamiento que no sea completamente idempotente de forma natural, peek_lock es la opción correcta. La contraparte es que ahora tienes mensajes que se pueden entregar más de una vez — necesitas manejar eso.

from azure.servicebus import ServiceBusClient
from azure.servicebus._common.auto_lock_renewer import AutoLockRenewer

async def consume(connection_str: str, queue_name: str):
    async with ServiceBusClient.from_connection_string(connection_str) as client:
        async with client.get_queue_receiver(queue_name) as receiver:
            renewer = AutoLockRenewer()
            async for message in receiver:
                renewer.register(receiver, message, max_lock_renewal_duration=300)
                try:
                    await process_message(message)
                    await receiver.complete_message(message)
                except IdempotentDuplicate:
                    await receiver.complete_message(message)  # ACK sin procesar
                except ProcessingError as e:
                    await receiver.abandon_message(message)   # Vuelve al queue
                    logger.error("message.processing_failed", exc=e)

AutoLockRenewer: el problema del lock expiry

El Service Bus de Azure asigna un lock time por defecto de 60 segundos. Si el procesamiento de un mensaje tarda más de 60 segundos — un merge de datos pesado, una llamada HTTP lenta, una query compleja — el lock expira y el mensaje queda disponible para otro consumer.

AutoLockRenewer renueva el lock automáticamente en background mientras el mensaje está siendo procesado. El parámetro max_lock_renewal_duration define el límite absoluto de renovación.

Sin esto, en producción bajo carga alta, mensajes que tardaban más de 60 segundos generaban duplicados silenciosos. El sistema “funcionaba” pero procesaba algunos mensajes dos veces.

Idempotencia: más que un try/except

Idempotencia significa que procesar el mismo mensaje N veces produce el mismo resultado que procesarlo una vez. Hay dos formas de implementarla:

Natural: la operación es naturalmente idempotente. SET value = X es idempotente. INCREMENT value BY 1 no lo es.

Por deduplicación: trackear qué mensajes ya fueron procesados y descartar duplicados.

Para el segundo caso, la implementación mínima es una tabla de mensajes procesados con un índice único en el ID del mensaje:

async def process_message(message: ServiceBusMessage) -> None:
    message_id = str(message.message_id)

    # Check idempotency before processing
    if await repository.message_already_processed(message_id):
        raise IdempotentDuplicate(f"Message {message_id} already processed")

    # Process...
    result = await domain_service.process(parse_message(message))

    # Save result + mark as processed in the same transaction
    await repository.save_result_and_mark_processed(result, message_id)

El punto crítico: el save del resultado y el mark como procesado tienen que ocurrir en la misma transacción atómica. Si guardan el resultado pero falla el mark, el próximo reintento procesa el mensaje de nuevo. Si guardan el mark pero falla el resultado, perdiste el procesamiento.

Con MongoDB (que soporta transacciones desde v4.0):

async def save_result_and_mark_processed(
    self, result: MergeResult, message_id: str
) -> None:
    async with await self._client.start_session() as session:
        async with session.start_transaction():
            await self._results_col.insert_one(
                result.to_dict(), session=session
            )
            await self._processed_col.insert_one(
                {"_id": message_id, "processed_at": datetime.utcnow()},
                session=session
            )

Post-merging processor

El caso de uso concreto tenía un paso adicional después del merge: un post-merging processor que ejecutaba lógica de transformación sobre el resultado merged. Este paso también tenía que ser idempotente.

La solución: el post-merging processor recibe el resultado del merge como input, no el mensaje original. Tiene su propio estado de procesamiento (idempotency key = merge result ID). El pipeline es:

Service Bus message
    → Agent merging (idempotente por message_id)
    → Merge result (guardado en DB)
    → Post-merge trigger (idempotente por merge_result_id)
    → Final output

Cada etapa tiene su propio mecanismo de idempotencia. Un fallo en cualquier punto puede ser retomado desde esa etapa sin reejecutar las anteriores.

Lo que aprendí

peek_lock + AutoLockRenewer es la configuración correcta por defecto. La complejidad adicional de manejar idempotencia vale infinitamente más que los mensajes perdidos de receive_and_delete.

La arquitectura hexagonal no es burocracia para proyectos pequeños. Es la diferencia entre poder testear la lógica de negocio en aislamiento y tener que levantar todos los servicios externos para correr un test unitario.

Los mensajes duplicados son silenciosos. Sin métricas o logs explícitos de duplicados detectados, nunca sabrás cuántos mensajes se están procesando dos veces. Loggea los IdempotentDuplicate — es información valiosa sobre el comportamiento del sistema bajo carga.

El ID del mensaje es tu amigo. Service Bus asigna un message_id único por default. Úsalo como idempotency key antes de intentar construir el tuyo propio.

compartir

X LinkedIn
Volver al blog