infrastructure · temario
6.6tema 6 de 7

Messaging y Event Bus

El adapter de mensajería publica domain events sin acoplar el dominio al broker.

Los domain events son hechos del negocio que el dominio necesita comunicar al resto del sistema: `<SomethingHappened>`, `<StateChanged>`, `<ResourceCreated>`. Para publicarlos, el sistema necesita un mecanismo de transporte —un message broker, un bus en memoria, un outbox— pero el dominio no puede conocer ese mecanismo. Si el Aggregate importa RabbitMQ o Redis directamente, el dominio queda acoplado a la infraestructura de mensajería: cambiar de broker o hacer tests sin broker se vuelve imposible.

El EventBus es un puerto driven: una abstracción que vive en la capa de aplicación o en el shared_kernel. Se define como `ABC` o `Protocol` con un método `publish(events: list[DomainEvent])`. El Aggregate no conoce el EventBus; solo acumula eventos en una lista interna. Es el Application Service, después de ejecutar el comportamiento del Aggregate y de hacer commit, quien drena esos eventos y los pasa al EventBus. Esa secuencia —ejecutar → persistir → publicar— es crítica para evitar publicaciones dobles o pérdidas.

El EventBus in-memory es el adapter de test y de desarrollo local. No conecta a ningún broker: simplemente colecciona los eventos publicados en una lista accesible desde el test. Esto permite escribir tests unitarios que afirman qué eventos se publicaron después de un caso de uso sin levantar RabbitMQ ni Redis. El adapter in-memory implementa exactamente el mismo puerto que el adapter de producción.

El adapter real de mensajería conecta al broker, serializa el evento a JSON o a un formato de mensaje, y lo publica en el exchange o topic correcto. Este adapter puede usar cualquier SDK de broker: `pika` para RabbitMQ, `redis-py` para Redis Streams, `confluent-kafka` para Kafka. Toda esa complejidad de serialización y routing queda encapsulada en el adapter; el Application Service nunca la ve. Si se cambia de RabbitMQ a Redis Streams, solo cambia el adapter.

El momento de publicación es una decisión arquitectónica con consecuencias de consistencia. Publicar antes del commit crea un dual-write: el evento puede llegar a los suscriptores antes de que la transacción confirme, o el commit puede fallar después de publicar. La regla es: publicar después de commit exitoso, desde el Application Service. Si se necesita garantía at-least-once, el patrón Outbox registra los eventos en la misma transacción y un proceso separado los publica al broker.

structure.txt
# shared_kernel/application/event_bus.py — Port (abstract)
# Lives in application layer. No broker imports.
from abc import ABC, abstractmethod
from ..domain.events.domain_event import DomainEvent

class EventBus(ABC):
    @abstractmethod
    def publish(self, events: list[DomainEvent]) -> None: ...

# tests/unit/<bounded_context>/fakes/in_memory_event_bus.py — Test Adapter
from ....shared_kernel.application.event_bus import EventBus
from ....shared_kernel.domain.events.domain_event import DomainEvent

class InMemoryEventBus(EventBus):
    def __init__(self) -> None:
        self.published: list[DomainEvent] = []

    def publish(self, events: list[DomainEvent]) -> None:
        self.published.extend(events)

# infrastructure/events/<broker>_event_bus.py — Real Adapter
# Connects to broker. Serializes events. No domain logic.
import json
from ....shared_kernel.application.event_bus import EventBus
from ....shared_kernel.domain.events.domain_event import DomainEvent

class <Broker>EventBus(EventBus):
    def __init__(self, connection: <BrokerConnection>) -> None:
        self._conn = connection

    def publish(self, events: list[DomainEvent]) -> None:
        for event in events:
            payload = json.dumps({
                'type': event.__class__.__name__,
                'aggregate_id': str(event.aggregate_id),
            })
            self._conn.publish(topic=event.__class__.__name__, message=payload)

# application/commands/<do_something>/handler.py — Publish after commit
class <DoSomething>Handler:
    def __init__(self, repo: <Aggregate>Repository, uow: UnitOfWork, event_bus: EventBus) -> None:
        self._repo = repo
        self._uow = uow
        self._event_bus = event_bus

    def handle(self, command: <DoSomething>Command) -> None:
        with self._uow:
            aggregate = self._repo.get(command.aggregate_id)
            aggregate.<perform_action>(command.<param>)
            self._repo.add(aggregate)
        # Drain and publish AFTER commit; never inside the transaction
        # BAD: event_bus.publish(...) inside 'with self._uow:' = dual-write risk
        events = aggregate.pull_events()
        self._event_bus.publish(events)
project/

Debugging lab

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

0/5 tests passing0%
  1. 6.6.5.1

    # domain/model/<aggregate>.py from shared_kernel.application.event_bus import EventBus class <Aggregate>: def __init__(self, event_bus: EventBus) -> None: self._event_bus = event_bus def <perform_action>(self, <param>: <ValueObject>) -> None: self._<field> = <param> self._event_bus.publish([<SomethingHappened>(aggregate_id=self.id)])

  2. 6.6.5.2

    # application/commands/<do_something>/handler.py class <DoSomething>Handler: def handle(self, command: <DoSomething>Command) -> None: with self._uow: aggregate = self._repo.get(command.aggregate_id) aggregate.<perform_action>(command.<param>) self._repo.add(aggregate) events = aggregate.pull_events() self._event_bus.publish(events) # published INSIDE transaction

  3. 6.6.5.3

    # infrastructure/events/<broker>_event_bus.py class <Broker>EventBus(EventBus): def publish(self, events: list[DomainEvent]) -> None: for event in events: if isinstance(event, <SomethingHappened>): if event.<field> == '<special_value>': return # business rule: skip special events payload = json.dumps({'type': event.__class__.__name__}) self._conn.publish(topic=event.__class__.__name__, message=payload)

  4. 6.6.5.4

    # application/commands/<do_something>/handler.py class <DoSomething>Handler: def handle(self, command: <DoSomething>Command) -> None: with self._uow: aggregate = self._repo.get(command.aggregate_id) aggregate.<perform_action>(command.<param>) self._repo.add(aggregate) try: events = aggregate.pull_events() self._event_bus.publish(events) except Exception: pass # silently swallowing publish failure

  5. 6.6.5.5

    # <ContextA>/application/commands/<do_something>/handler.py from contexts.<context_b>.domain.events.<event_b_happened> import <EventBHappened> class <DoSomething>Handler: def handle(self, command: <DoSomething>Command) -> None: with self._uow: aggregate = self._repo.get(command.aggregate_id) aggregate.<perform_action>(command.<param>) events = aggregate.pull_events() # Publishing an event from another bounded context self._event_bus.publish([<EventBHappened>(aggregate_id=aggregate.id)])