application · temario
5.8tema 8 de 8

Event Bus y Despacho de Eventos

El event bus publica domain events tras la transacción; los handlers reaccionan sin acoplar el emisor.

El event bus es el mecanismo por el que los domain events, generados dentro de un aggregate, se hacen llegar a los interesados. El aggregate registra eventos internamente durante su operación; el Command Handler los extrae y los envía al bus después de confirmar la transacción. El bus los entrega a los handlers registrados.

Distinguir entre domain events y application events es importante. Un domain event es un hecho ocurrido dentro del dominio ('OrderPlaced', 'PaymentConfirmed'); es parte del lenguaje ubicuo y vive en la capa de dominio. Un application event es más técnico o cross-cutting ('UserNotified', 'CacheInvalidated') y puede vivir en la capa de aplicación. Los domain events son inmutables y en tiempo pasado.

El EventBus es un puerto outbound (ABC). La implementación puede ser in-memory (síncrona, útil para tests y para empezar), o async (con un broker como RabbitMQ o Kafka) para sistemas distribuidos. La capa de aplicación siempre habla con el puerto, nunca con la implementación concreta.

El dispatch síncrono (in-memory) tiene la ventaja de la simplicidad: el evento se entrega en la misma llamada. El riesgo es que un fallo en un handler de evento puede propagar una excepción al handler de comando. Una implementación robusta in-memory captura las excepciones de los handlers individuales y las registra sin interrumpir el flujo principal.

En sistemas distribuidos, el patrón Transactional Outbox garantiza que los eventos se persisten en la misma transacción que el aggregate (en una tabla outbox), y un proceso separado los envía al broker. Esto evita el 'dual-write' y asegura que ningún evento se pierde si el sistema falla entre el commit y el envío. Este patrón avanzado se trata en el Módulo 7.

structure.txt
# shared_kernel/application/event_bus.py — Puerto (abstracción)
from abc import ABC, abstractmethod
from <project>.shared_kernel.domain.domain_event import DomainEvent

class EventBus(ABC):
    """
    Puerto outbound: publica domain events.
    La implementación puede ser in-memory, RabbitMQ, Kafka, etc.
    """
    @abstractmethod
    def publish(self, event: DomainEvent) -> None: ...

    @abstractmethod
    def subscribe(self, event_type: type, handler: 'EventHandler') -> None: ...

# shared_kernel/application/event_handler.py — Contrato de handler
from abc import ABC, abstractmethod

class EventHandler(ABC):
    @abstractmethod
    def handle(self, event: DomainEvent) -> None: ...

# --- Domain Event base (en shared_kernel/domain) ---
from dataclasses import dataclass, field
from datetime import datetime, timezone
import uuid

@dataclass(frozen=True)
class DomainEvent:
    event_id: str = field(default_factory=lambda: str(uuid.uuid4()))
    occurred_on: datetime = field(default_factory=lambda: datetime.now(timezone.utc))
    # Subclases añaden los datos del evento (en tiempo pasado, inmutables)

# --- Domain Event concreto (en domain/events/) ---
@dataclass(frozen=True)
class <SomethingHappened>(DomainEvent):
    aggregate_id: str = ''
    # ... datos relevantes del hecho

# --- Uso en el Command Handler ---
class <DoSomething>Handler:
    def handle(self, command: <DoSomething>Command) -> None:
        aggregate = self._repo.get_by_id(<Aggregate>Id(command.aggregate_id))
        aggregate.<do_something>()

        with self._uow:
            self._repo.save(aggregate)

        # Los eventos se publican DESPUÉS del commit
        for event in aggregate.pull_events():
            self._bus.publish(event)

# --- Implementación in-memory (para tests y sistemas simples) ---
class InMemoryEventBus(EventBus):
    def __init__(self) -> None:
        self._handlers: dict[type, list[EventHandler]] = {}

    def subscribe(self, event_type: type, handler: EventHandler) -> None:
        self._handlers.setdefault(event_type, []).append(handler)

    def publish(self, event: DomainEvent) -> None:
        for handler in self._handlers.get(type(event), []):
            handler.handle(event)  # síncrono; excepciones se propagan

# --- Event Handler de aplicación (reacción al evento) ---
# application/event_handlers/<something_happened>_handler.py
class On<SomethingHappened>(EventHandler):
    def __init__(self, notification_port: EmailNotificationPort) -> None:
        self._notification_port = notification_port

    def handle(self, event: <SomethingHappened>) -> None:
        # Reacción: notificar, actualizar proyección, disparar otro command, etc.
        self._notification_port.notify(
            recipient=event.aggregate_id,
            subject='Event occurred',
            body=f'Event {event.event_id} processed.',
        )
        # NO: disparar lógica de negocio aquí — eso es del dominio

Debugging lab

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

0/5 tests passing0%
  1. 5.8.5.1

    class PlaceOrderHandler: def handle(self, cmd): order = Order.create(...) with self._uow: self._repo.save(order) for e in order.pull_events(): self._bus.publish(e) # dentro del with

  2. 5.8.5.2

    import pika # RabbitMQ class PlaceOrderHandler: def __init__(self, repo, uow, connection: pika.BlockingConnection): ...

  3. 5.8.5.3

    @dataclass(frozen=True) class OrderCreating(DomainEvent): # nombre en gerundio

  4. 5.8.5.4

    class OnOrderPlaced(EventHandler): def handle(self, event: OrderPlaced) -> None: order = self._repo.get_by_id(OrderId(event.order_id)) order.mark_as_notified() # modificar el aggregate self._repo.save(order)

  5. 5.8.5.5

    # En el aggregate: class Order: def place(self): import smtplib smtplib.SMTP('smtp.example.com').sendmail(...) # notificar en el dominio