strategic · temario
2.5tema 5 de 5

Integración entre contextos

OHS, Published Language y ACL como estrategia de integración; eventos de integración para comunicación asíncrona entre bounded contexts.

Una vez que los bounded contexts están definidos y el context map establece sus relaciones, la pregunta práctica es: ¿cómo se comunican? La integración entre contextos debe ser explícita, versionada y protegida. Existen dos grandes modalidades: síncrona (llamada directa a una API) y asíncrona (eventos de integración).

La tríada OHS + Published Language + ACL es la estrategia de integración más robusta. El upstream expone un Open Host Service: una API estable, versionada y orientada a los consumidores, expresada en un Published Language (esquema JSON, Protobuf, OpenAPI). El downstream, por su parte, instala una Anti-Corruption Layer que traduce ese contrato externo al lenguaje de su propio dominio.

Los eventos de integración son la alternativa asíncrona. A diferencia de los domain events (que son internos al aggregate y expresan hechos del negocio), los integration events son contratos entre contextos: son más estables, están versionados y se publican a través de un message broker o event bus externo. Un contexto publica lo que le sucedió; otros contextos suscriben y reaccionan según su propio modelo.

La diferencia entre domain events e integration events es importante. Un domain event es parte del lenguaje ubicuo del contexto que lo emite; puede ser rico en detalles internos. Un integration event es un contrato público entre contextos: debe ser mínimo, estable y versionado, expresado en el Published Language acordado.

La integración asíncrona desacopla completamente los ciclos de despliegue y las ventanas de disponibilidad de los contextos. La contrapartida es la consistencia eventual: el downstream procesará el evento en algún momento futuro, no de forma inmediata. Diseñar los handlers de integración como idempotentes (ejecutar el mismo evento dos veces no causa daño) es esencial para la robustez del sistema.

La elección entre integración síncrona y asíncrona depende de los requisitos de latencia, tolerancia a fallos y acoplamiento aceptable. Ambas pueden coexistir: un contexto puede exponer OHS para consultas en tiempo real y publicar integration events para notificar cambios de estado a sus suscriptores.

structure.txt
# Integration strategy: OHS + PL + ACL + Integration Events

# 1. Upstream: Open Host Service with Published Language
# contexts/<context_upstream>/infrastructure/api/schemas.py
from pydantic import BaseModel

class <ResourceV1Schema>(BaseModel):  # Published Language — versioned contract
    id: str
    <field_a>: str
    <field_b>: int

# contexts/<context_upstream>/infrastructure/api/controllers.py
# router = APIRouter(prefix='/v1/<context_upstream>')
# @router.get('/<resource>/{id}')
# def get_resource(id: str) -> <ResourceV1Schema>: ...


# 2. Downstream: Anti-Corruption Layer
# contexts/<context_downstream>/infrastructure/acl/<context_upstream>_acl.py
from src.project.contexts.<context_downstream>.domain.model.<entity> import <Entity>, <EntityId>
from src.project.contexts.<context_downstream>.domain.model.<value_obj> import <ValueObject>

class <ContextUpstream>Acl:
    """Translates <ContextUpstream> Published Language into <ContextDownstream> domain model."""

    def to_<entity>(self, schema: dict) -> <Entity>:
        return <Entity>(
            id=<EntityId>(schema['id']),
            <field>=<ValueObject>(schema['<field_a>']),
        )


# 3. Integration Event (upstream publishes what happened)
# contexts/<context_upstream>/domain/events/<something>_occurred.py
from dataclasses import dataclass
from datetime import datetime
from uuid import UUID

@dataclass(frozen=True)
class <SomethingOccurred>:  # Integration Event — public contract, versioned
    event_id: UUID
    <aggregate_id>: str      # primitive or shared VO — not a domain class
    <relevant_field>: str
    occurred_at: datetime
    schema_version: str = 'v1'


# 4. Downstream integration event handler (with ACL translation)
# contexts/<context_downstream>/application/handlers/<something>_occurred_handler.py
class <SomethingOccurred>Handler:
    def __init__(self, acl: <ContextUpstream>Acl, repo: <EntityRepository>) -> None:
        self._acl = acl
        self._repo = repo

    def handle(self, event: <SomethingOccurred>) -> None:  # idempotent!
        if self._repo.exists(<EntityId>(event.<aggregate_id>)):
            return  # already processed — idempotency guard
        entity = self._acl.to_<entity>(event.__dict__)
        self._repo.save(entity)

Debugging lab

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

0/5 tests passing0%
  1. 2.5.5.1

    # A domain event is published directly to the external message broker: # contexts/<context_a>/domain/events/aggregate_created.py class AggregateCreated: # rich domain event — contains internal VOs aggregate: <Aggregate> # ← full domain object internal_state: dict # published directly: event_bus.publish(AggregateCreated(aggregate=self))

  2. 2.5.5.2

    # Integration event handler is not idempotent: class <SomethingOccurred>Handler: def handle(self, event: <SomethingOccurred>) -> None: entity = self._acl.to_entity(event.__dict__) self._repo.save(entity) # saves every time, even if already processed

  3. 2.5.5.3

    # ACL is missing — downstream uses the upstream Published Language schema directly as a domain object: # contexts/<context_downstream>/domain/model/entity.py from src.project.contexts.<context_upstream>.infrastructure.api.schemas import ResourceV1Schema class Entity: def __init__(self, schema: ResourceV1Schema): ...

  4. 2.5.5.4

    # Integration event has no version field and a breaking change was deployed: @dataclass(frozen=True) class <SomethingOccurred>: aggregate_id: str old_field: str # removed in new version; subscribers break

  5. 2.5.5.5

    # Synchronous call between contexts creates tight coupling: # contexts/<context_downstream>/application/handlers/do_something_handler.py import httpx class DoSomethingHandler: def handle(self, cmd: DoSomethingCommand) -> None: # calls upstream synchronously — if upstream is down, this fails response = httpx.get('http://<context_upstream>/v1/resource/' + cmd.id) data = response.json() # ... process data