advanced · temario
7.4tema 4 de 7

Outbox Pattern

Transactional outbox, dual-write problem y entrega garantizada de eventos.

El problema del dual-write: cuando guardas un aggregate Y publicas un domain event al broker en dos operaciones separadas, cualquiera de las dos puede fallar de forma independiente. Un commit exitoso a la base de datos no garantiza que el evento llegó al broker — el proceso puede morir entre las dos operaciones, dejando el sistema en un estado inconsistente donde el aggregate fue guardado pero el evento nunca se publicó.

El Outbox Pattern resuelve esto: dentro de la misma transacción de base de datos, guardas el aggregate Y escribes el evento en una tabla de outbox. Un proceso relay separado lee la tabla de outbox y publica los eventos al broker de forma asíncrona. La atomicidad está garantizada porque ambas escrituras ocurren en el mismo commit de BD.

El outbox garantiza entrega at-least-once: el relay puede publicar el mismo evento más de una vez si crashea después de publicar pero antes de marcar el registro como enviado. Por esto, los consumidores del broker deben ser idempotentes — procesar el mismo mensaje dos veces debe tener el mismo efecto que procesarlo una vez. Se usa el event_id como clave de deduplicación.

El relay (Transactional Outbox Relay) puede implementarse con polling (consulta periódica a la tabla de outbox) o con CDC (Change Data Capture, por ejemplo Debezium) que detecta cambios en la tabla a nivel de log de base de datos. El CDC es más eficiente a escala porque elimina el polling constante.

Cuándo NO usar el Outbox Pattern: si tienes un event bus en proceso (sin broker externo), o si tu broker soporta envío transaccional nativo. Añade infraestructura: tabla de outbox, proceso relay, limpieza periódica de registros publicados. Solo vale la pena para integraciones externas o entre bounded contexts con brokers reales.

structure.txt
# infrastructure/outbox/outbox_writer.py    — Writes event to outbox in same tx
# infrastructure/outbox/outbox_relay.py     — Reads outbox, publishes to broker
# infrastructure/persistence/<repo>.py       — Saves aggregate + calls outbox writer

class OutboxWriter:
    def __init__(self, session: <DbSession>) -> None:
        self._session = session

    def write(self, event: <DomainEvent>) -> None:
        # Must be called within the same DB transaction as the aggregate save
        record = OutboxRecord(
            event_id=str(event.event_id),
            event_type=type(event).__name__,
            payload=event.to_dict(),
            created_at=datetime.utcnow(),
            published=False,
        )
        self._session.add(record)
        # No commit here — caller owns the transaction

class <AggregateRepository>:
    def save(self, aggregate: <Aggregate>) -> None:
        # 1. Persist aggregate state
        self._session.merge(<AggregateORM>.from_domain(aggregate))
        # 2. Write events to outbox IN THE SAME TRANSACTION
        for event in aggregate.pull_events():
            self._outbox_writer.write(event)
        # 3. Single commit — both aggregate and outbox are atomic
        self._session.commit()

class OutboxRelay:
    def __init__(self, session: <DbSession>, broker: <MessageBroker>) -> None:
        self._session = session
        self._broker = broker

    def relay(self) -> None:
        unpublished = (
            self._session.query(OutboxRecord)
            .filter_by(published=False)
            .order_by(OutboxRecord.created_at)
            .limit(100)
            .all()
        )
        for record in unpublished:
            self._broker.publish(record.event_type, record.payload)
            record.published = True
        self._session.commit()
project/
<project>
<bc>
infrastructure
outbox
domain

Debugging lab

Detecta y corrige el error o la violación de diseño.

0/5 tests passing0%
  1. 7.4.5.1

    # Bug: repository saves aggregate, then calls broker.publish() outside the DB transaction class AggregateRepository: def save(self, aggregate: Aggregate) -> None: self._session.merge(AggregateORM.from_domain(aggregate)) self._session.commit() # Transaction committed here # Bug: publishing outside transaction — if broker fails, event is lost for event in aggregate.pull_events(): self._broker.publish(type(event).__name__, event.to_dict())

  2. 7.4.5.2

    # Bug: outbox relay marks record.published = True before calling broker.publish() class OutboxRelay: def relay(self) -> None: unpublished = self._session.query(OutboxRecord).filter_by(published=False).limit(100).all() for record in unpublished: # Bug: marking as published BEFORE sending to broker record.published = True self._session.commit() self._broker.publish(record.event_type, record.payload)

  3. 7.4.5.3

    # Bug: integration event consumer processes same event_id twice, applying the same effect twice class PaymentIntegrationHandler: def handle(self, event: PaymentReceivedIntegrationEvent) -> None: # Bug: no idempotency check — same event processed twice = double effect account = self._repo.find(AccountId(event.account_id)) account.credit(Amount(event.amount)) self._repo.save(account)

  4. 7.4.5.4

    # Bug: outbox_writer.write() called with a separate session, outside the repository transaction class AggregateRepository: def __init__(self, session: DbSession) -> None: self._session = session # Bug: OutboxWriter gets its own independent session self._outbox_writer = OutboxWriter(session=DbSession()) def save(self, aggregate: Aggregate) -> None: self._session.merge(AggregateORM.from_domain(aggregate)) for event in aggregate.pull_events(): self._outbox_writer.write(event) # Different session = different transaction! self._session.commit()

  5. 7.4.5.5

    # Bug: outbox pattern used for in-process domain events within a single bounded context # These events never leave the process — they go to in-process handlers only class OrderService: def place_order(self, cmd: PlaceOrderCommand) -> None: order = Order.place(cmd.order_id, cmd.items) self._repo.save(order) # Bug: using outbox for in-process events that never reach a broker for event in order.pull_events(): self._outbox_writer.write(event) # Unnecessary outbox overhead