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.
# 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 dominioDebugging lab
Detecta y corrige el error o la violación de diseño.
- 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
- 5.8.5.2
import pika # RabbitMQ class PlaceOrderHandler: def __init__(self, repo, uow, connection: pika.BlockingConnection): ...
- 5.8.5.3
@dataclass(frozen=True) class OrderCreating(DomainEvent): # nombre en gerundio
- 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.8.5.5
# En el aggregate: class Order: def place(self): import smtplib smtplib.SMTP('smtp.example.com').sendmail(...) # notificar en el dominio