โ–  CultivaFlow SaaS

Arquitectura de Microservicios

Guia de referencia: migracion desde monolito Django a microservicios distribuidos con patrones Strangler Fig, Saga, Circuit Breaker y Bulkhead.

5 servicios identificados
4 patrones de resiliencia
Kafka event bus
Python / FastAPI
Diagrama de arquitectura
CultivaFlow โ€” Vista de servicios
Capa cliente
๐ŸŒ
React SPA / Mobile App
Clientes web y movil de las agencias
โ†“
HTTPS / REST
API Gateway
๐Ÿ”€
API Gateway โ€” cultivaflow-gateway:8080
Rate limiting ยท Auth middleware ยท Request routing ยท Circuit breaker global ยท Aggregation
Gateway
โ†“
HTTP interno (gRPC planificado)
Servicios de dominio
Core
๐Ÿ”
AuthService
:8001 โ€” JWT, OAuth, orgs
Core
๐Ÿ“ข
CampaignService
:8002 โ€” CRUD campanas
Scale
๐Ÿค–
AIService
:8003 โ€” LLM, prompts, tokens
New
๐Ÿ“Š
AnalyticsService
:8004 โ€” metricas, clicks
๐Ÿ””
NotificationService
:8005 โ€” email, webhook
โ†“โ†‘
eventos asincronos
โšก
Apache Kafka โ€” Event Bus
Comunicacion asincrona entre servicios โ€” persistencia de eventos โ€” replay
campaign.created ai.content.ready analytics.tracked notification.send
Bases de datos (Database-per-service)
๐Ÿ—„๏ธ
PostgreSQL
auth_db
๐Ÿ—„๏ธ
PostgreSQL
campaign_db
โšก
Redis
ai_cache + queue
๐Ÿ”ท
ClickHouse
analytics_db
๐Ÿ—„๏ธ
PostgreSQL
notif_db
Flujo Saga โ€” crear campana
๐Ÿ“‹ OrderFulfillment Pattern โ€” transaccion distribuida atomica con compensacion
Paso 1
Validar Auth
AuthService
โ†’
Paso 2
Crear Campana
CampaignService
โ†’
Paso 3
Generar Copy
AIService
โ†’
Paso 4
Registrar Evento
AnalyticsService
โ†’
Paso 5
Notificar
NotificationService
Compensacion: si AIService falla en paso 3, se ejecuta en orden inverso: cancel_campaign (paso 2), revoke_auth_token (paso 1). Garantia de consistencia eventual.
Patrones clave implementados
1
Strangler Fig Pattern
Migracion gradual desde monolito Django
Contexto CultivaFlow: El monolito Django sigue sirviendo /admin y /legacy. El proxy enruta /api/v2/campaigns al nuevo CampaignService mientras lo demas permanece en el monolito.
Python
strangler_proxy.py
class StranglerProxy: routes = { "/api/v2/campaigns": "http://campaign-svc:8002", "/api/v2/auth": "http://auth-svc:8001", "/api/v2/ai": "http://ai-svc:8003", "/legacy": "http://django-monolith:8000", } async def route(self, request): for prefix, target in self.routes.items(): if request.path.startswith(prefix): return await self.proxy(target, request) return await self.proxy("http://django-monolith:8000", request)
2
Circuit Breaker
Proteger CampaignService de fallos de AIService
Contexto CultivaFlow: Si Claude API esta caido, las campanas siguen creandose con copy pendiente (DRAFT). El Circuit Breaker abre tras 5 errores y espera 30s antes de reintentar.
Python
ai_circuit_breaker.py
ai_breaker = CircuitBreaker( failure_threshold=5, recovery_timeout=30, success_threshold=2 ) async def generate_copy(brief: str) -> str: try: return await ai_breaker.call( ai_service.generate, brief ) except CircuitBreakerOpenError: # Fallback: campana en estado DRAFT return CopyResult(status="DRAFT", copy=None)
3
Saga Pattern (Orquestado)
Transacciones distribuidas con compensacion
Contexto CultivaFlow: Crear una campana requiere coordinar 4 servicios. Si el pago del credito AI falla, se deben cancelar los pasos anteriores automaticamente.
Python
create_campaign_saga.py
saga = CreateCampaignSaga([ SagaStep("validate_auth", action=auth_svc.validate, compensation=auth_svc.revoke), SagaStep("create_campaign", action=campaign_svc.create, compensation=campaign_svc.cancel), SagaStep("deduct_ai_credits", action=billing_svc.deduct, compensation=billing_svc.refund), SagaStep("generate_copy", action=ai_svc.generate, compensation=ai_svc.discard), ]) result = await saga.execute(campaign_data)
4
Bulkhead Pattern
Aislar AIService del resto del sistema
Contexto CultivaFlow: El AIService consume mucha CPU/memoria. Con Bulkhead, los picos de uso de generacion AI no afectan a AuthService ni CampaignService.
Python
bulkhead_config.py
# Thread pools separados por servicio ai_pool = ThreadPoolExecutor(max_workers=10) auth_pool = ThreadPoolExecutor(max_workers=20) campaign_pool = ThreadPoolExecutor(max_workers=15) # Kubernetes resource limits (YAML) # ai-service: cpu: "2", memory: "4Gi" # auth-service: cpu: "500m", memory: "512Mi" # campaign-svc: cpu: "1", memory: "1Gi" @bulkhead(pool=ai_pool, max_concurrent=10) async def generate_content(brief): ...
Configuracion de resiliencia
๐Ÿ”Œ
Circuit Breaker โ€” AIService
Abre el circuito cuando Claude API devuelve errores repetidos. Modo degradado: campanas en DRAFT sin copy.
failure_threshold: 5 recovery_timeout: 30s success_threshold: 2 fallback: DRAFT_MODE
๐Ÿ”„
Retry con Backoff โ€” HTTP calls
Reintentos exponenciales para errores transientes en llamadas entre servicios via httpx.
max_attempts: 3 min_wait: 2s max_wait: 10s multiplier: exponential
๐ŸŠ
Bulkhead โ€” recursos aislados
Pools de threads separados evitan que AIService (intensivo en CPU) sature los servicios criticos.
ai_pool: 10 workers auth_pool: 20 workers campaign_pool: 15 workers queue_size: 100 tasks
Matriz de comunicacion entre servicios
Desde โ†’ Hasta Patron Protocolo Razon Estado
Gateway โ†’ AuthService Sincrono REST / HTTP Validar JWT en cada request โ€” baja latencia requerida โœ“ Activo
Gateway โ†’ CampaignService Sincrono REST / HTTP CRUD campanas โ€” respuesta inmediata al usuario โœ“ Activo
CampaignService โ†’ AIService Asincrono Kafka topic Generacion de copy puede tardar 5-30s โ€” no bloquear UI โœ“ Activo
AIService โ†’ AnalyticsService Asincrono Kafka topic Registro de tokens usados โ€” no critico para el flujo โœ“ Activo
CampaignService โ†’ NotificationService Asincrono Kafka topic Notificacion al agente cuando campana esta lista โš  Pendiente
Gateway โ†’ AnalyticsService Sincrono REST / HTTP Dashboard en tiempo real โ€” datos frescos requeridos โš  En migracion