Processor vs Subscriber: Diferencias Clave
🎯 Resumen Ejecutivo
TL;DR:
.processor= Worker de BullMQ que procesa jobs directamente de una cola específica.subscriber= Handler de eventos de dominio que se ejecuta cuando ocurre un evento (ej: "meeting created")
Ambos procesan trabajo asíncrono, pero tienen orígenes y propósitos diferentes.
📊 Comparación Visual
┌─────────────────────────────────────────────────────────────┐
│ PROCESOR (Worker) │
├─────────────────────────────────────────────────────────────┤
│ Origen: Jobs directamente en cola de BullMQ │
│ Ejemplo: meeting_transcript_process │
│ │
│ [BullMQ Queue] → @Processor → process(job) → Lógica │
│ │
│ - Escucha una cola específica │
│ - Procesa jobs de forma directa │
│ - Puede crear nuevos jobs en otras colas │
└─────────────────────────────────────────────────────────────┘
┌─────────────────────────────────────────────────────────────┐
│ SUBSCRIBER (Event Handler) │
├─────────────────────────────────────────────────────────────┤
│ Origen: Domain Events (eventos de dominio) │
│ Ejemplo: sync-meeting-summary-to-crm │
│ │
│ [DB Mutation] → DomainEvent → [Queue] → Processor │
│ │
│ ↓ Fan-out │
│ │
│ [Subscriber Queue] → @Processor │
│ │
│ ↓ │
│ │
│ Subscriber.handle(event) → Lógica │
│ │
│ - Escucha eventos de dominio (create/update/delete) │
│ - Se ejecuta automáticamente cuando ocurre el evento │
│ - Puede crear jobs en colas │
└─────────────────────────────────────────────────────────────┘
🔍 Processor (Worker de BullMQ)
¿Qué es?
Un Processor es un worker de BullMQ que procesa jobs directamente de una cola específica. Es el equivalente a un "worker" o "consumer" en sistemas de colas.
Características
-
Escucha una cola específica de BullMQ
typescript@Processor(BullJobs.TRANSCRIPT_PROCESS, { ...DEFAULT_WORKER_OPTIONS, concurrency: 12, }) export class MeetingTranscriptResponseProcessor extends WorkerHostProcessor { async process(job: Job<EventMeetingTranscriptResponse>) { // Procesa el job directamente } } -
Recibe jobs directamente de la cola
- Los jobs se añaden con
appendJobToQueue() - BullMQ los entrega automáticamente al processor
- Los jobs se añaden con
-
Puede crear nuevos jobs
typescript// Dentro del processor await this.summaryGenerationJobsQueueService.appendJobToQueue( requestId, BullJobs.MEETING_SUMMARY_GENERATION, { meetingId, tenantId, userId } ); -
Tiene control total sobre el procesamiento
- Maneja errores, reintentos, logging
- Puede hacer operaciones complejas
- Tiene acceso directo a servicios
Ejemplos de Processors
MeetingTranscriptResponseProcessor→ Procesameeting_transcript_processMeetingSummaryGeneratorProcessor→ Procesameeting_summary_generationMeetingVideoAudioProcessor→ Procesameeting_video_audio_conversionCalendarImportProcessor→ Procesacalendar_import
Flujo Típico
1. Algo en el código llama: appendJobToQueue(queue, payload)
2. BullMQ añade el job a la cola
3. Processor escucha la cola y recibe el job
4. Processor ejecuta process(job)
5. Processor puede crear nuevos jobs si es necesario
📡 Subscriber (Event Handler)
¿Qué es?
Un Subscriber es un handler de eventos de dominio que se ejecuta automáticamente cuando ocurre un evento (ej: "meeting created", "user deleted"). Implementa el patrón Observer/Subscriber del DDD.
Características
-
Escucha eventos de dominio, no colas directamente
typescript@Injectable() export class SyncMeetingSummaryToCrmSubscriber extends DomainSubscriber implements DomainEventSubscriber<CreateEvent<'meeting-summary'>> { readonly entity = DomainEntities.MEETING_SUMMARY; readonly operation = DomainOperations.CREATE; async handle(event: CreateEvent<'meeting-summary'>) { // Se ejecuta cuando se crea un meeting-summary } } -
Se registra automáticamente mediante DI
- No necesitas registrarlo manualmente en una cola
- El sistema lo descubre automáticamente
- Se ejecuta cuando ocurre el evento correspondiente
-
Recibe eventos estructurados
typescriptinterface CreateEvent<T> { entity: string; // 'meeting-summary' operation: 'create'; // 'create' | 'update' | 'delete' payload: T; // Datos del evento } -
Puede crear jobs en colas
typescriptasync handle(event: CreateEvent<'meeting-summary'>) { await this.crmMeetingOutgoingJobsQueueService.appendJobToQueue( requestId, BullJobs.Crm.Outgoing.Meeting, { meetingId: event.payload.meetingId } ); }
Ejemplos de Subscribers
SyncMeetingSummaryToCrmSubscriber→ Se ejecuta cuando se crea un meeting-summarySyncMeetingParticipantToCrmSubscriber→ Se ejecuta cuando se crea/actualiza un participantUpdateCompanyStatusFromSummarySubscriber→ Se ejecuta cuando se crea un meeting-summaryDisconnectCalendarOnUserDeleteSubscriber→ Se ejecuta cuando se elimina un user
Flujo Típico
1. Algo en el código hace: domainEventBus.appendEventToQueue(event)
2. DomainEventProcessor recibe el evento
3. DomainEventProcessor hace "fan-out" creando un job por cada subscriber
4. DomainEventSubscriberProcessor ejecuta cada subscriber
5. Subscriber.handle(event) se ejecuta
6. Subscriber puede crear jobs si es necesario
🔄 Arquitectura Completa: Cómo Trabajan Juntos
┌──────────────────────────────────────────────────────────────┐
│ FLUJO COMPLETO │
└──────────────────────────────────────────────────────────────┘
1. Código de negocio:
meetingService.create() → Guarda en DB
↓
domainEventBus.appendEventToQueue({
entity: 'meeting',
operation: 'create',
payload: meeting
})
↓
2. DomainEventProcessor (PROCESSOR):
- Recibe evento de dominio
- Busca todos los subscribers registrados para 'meeting:create'
- Hace fan-out: crea un job por cada subscriber
↓
3. DomainEventSubscriberProcessor (PROCESSOR):
- Recibe job de subscriber individual
- Ejecuta: subscriber.handle(event)
↓
4. SyncMeetingSummaryToCrmSubscriber (SUBSCRIBER):
- handle(event) se ejecuta
- Crea job: appendJobToQueue('crm_meeting_outgoing', ...)
↓
5. CrmMeetingOutgoingProcessor (PROCESSOR):
- Recibe job de 'crm_meeting_outgoing'
- Procesa sincronización con CRM
📋 Tabla Comparativa
| Aspecto | Processor | Subscriber |
|---|---|---|
| Origen | Jobs en cola de BullMQ | Domain Events |
| Decorador | @Processor(BullJobs.XXX) | implements DomainEventSubscriber |
| Método principal | process(job: Job<T>) | handle(event: DomainEvent) |
| Registro | Automático por BullMQ | Automático por DI (DiscoveryService) |
| Escucha | Cola específica | Eventos de entidad + operación |
| Ejecución | Directa cuando hay job | Indirecta vía DomainEventProcessor |
| Propósito | Procesar trabajo asíncrono | Reaccionar a cambios de dominio |
| Puede crear jobs | ✅ Sí | ✅ Sí |
| Control de reintentos | BullMQ automático | BullMQ automático (vía processor) |
| Ejemplo | MeetingTranscriptResponseProcessor | SyncMeetingSummaryToCrmSubscriber |
🎯 ¿Cuándo Usar Cada Uno?
Usa Processor cuando:
- ✅ Necesitas procesar trabajo que no está relacionado con eventos de dominio
- ✅ Tienes un flujo de procesamiento específico (ej: convertir video a audio)
- ✅ El trabajo se dispara manualmente o por cron
- ✅ Necesitas control directo sobre el procesamiento
- ✅ Ejemplos:
- Procesar transcripciones
- Convertir video a audio
- Importar datos de APIs externas
- Procesar webhooks externos
Usa Subscriber cuando:
- ✅ Necesitas reaccionar a cambios en el dominio
- ✅ Quieres desacoplar lógica de side-effects
- ✅ Necesitas que múltiples cosas pasen cuando algo cambia
- ✅ Quieres seguir el patrón DDD con Domain Events
- ✅ Ejemplos:
- Sincronizar con CRM cuando se crea un meeting
- Actualizar status de compañía cuando se crea un summary
- Enviar webhooks cuando cambia algo
- Aplicar control de acceso cuando se crea un recurso
🔗 Relación entre Ambos
Los Subscribers suelen crear jobs que son procesados por Processors:
// SUBSCRIBER: Reacciona al evento
@Injectable()
export class SyncMeetingSummaryToCrmSubscriber extends DomainSubscriber {
async handle(event: CreateEvent<'meeting-summary'>) {
// Crea un job que será procesado por un Processor
await this.crmMeetingOutgoingJobsQueueService.appendJobToQueue(
requestId,
BullJobs.Crm.Outgoing.Meeting, // ← Esta cola tiene un Processor
{ meetingId: event.payload.meetingId }
);
}
}
// PROCESSOR: Procesa el job creado por el Subscriber
@Processor(BullJobs.Crm.Outgoing.Meeting)
export class CrmMeetingOutgoingProcessor extends WorkerHostProcessor {
async process(job: Job<EventCrmMeetingIntegration>) {
// Procesa la sincronización con CRM
}
}
💡 Puntos Clave
-
Ambos son workers, pero:
- Processor = Worker directo de BullMQ
- Subscriber = Worker indirecto vía Domain Events
-
Ambos pueden crear jobs, pero:
- Processor crea jobs para continuar flujos de procesamiento
- Subscriber crea jobs como side-effects de eventos
-
Ambos procesan trabajo asíncrono, pero:
- Processor procesa trabajo "técnico" (conversiones, procesamiento)
- Subscriber procesa trabajo "de negocio" (sincronizaciones, notificaciones)
-
La diferencia principal es el origen:
- Processor: Jobs añadidos directamente a la cola
- Subscriber: Eventos de dominio que se convierten en jobs