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

text
┌─────────────────────────────────────────────────────────────┐
│                    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

  1. 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
      }
    }
    
  2. Recibe jobs directamente de la cola

    • Los jobs se añaden con appendJobToQueue()
    • BullMQ los entrega automáticamente al processor
  3. Puede crear nuevos jobs

    typescript
    // Dentro del processor
    await this.summaryGenerationJobsQueueService.appendJobToQueue(
      requestId,
      BullJobs.MEETING_SUMMARY_GENERATION,
      { meetingId, tenantId, userId }
    );
    
  4. Tiene control total sobre el procesamiento

    • Maneja errores, reintentos, logging
    • Puede hacer operaciones complejas
    • Tiene acceso directo a servicios

Ejemplos de Processors

  • MeetingTranscriptResponseProcessor → Procesa meeting_transcript_process
  • MeetingSummaryGeneratorProcessor → Procesa meeting_summary_generation
  • MeetingVideoAudioProcessor → Procesa meeting_video_audio_conversion
  • CalendarImportProcessor → Procesa calendar_import

Flujo Típico

text
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

  1. 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
      }
    }
    
  2. 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
  3. Recibe eventos estructurados

    typescript
    interface CreateEvent<T> {
      entity: string;      // 'meeting-summary'
      operation: 'create';  // 'create' | 'update' | 'delete'
      payload: T;          // Datos del evento
    }
    
  4. Puede crear jobs en colas

    typescript
    async 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-summary
  • SyncMeetingParticipantToCrmSubscriber → Se ejecuta cuando se crea/actualiza un participant
  • UpdateCompanyStatusFromSummarySubscriber → Se ejecuta cuando se crea un meeting-summary
  • DisconnectCalendarOnUserDeleteSubscriber → Se ejecuta cuando se elimina un user

Flujo Típico

text
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

text
┌──────────────────────────────────────────────────────────────┐
│                    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

AspectoProcessorSubscriber
OrigenJobs en cola de BullMQDomain Events
Decorador@Processor(BullJobs.XXX)implements DomainEventSubscriber
Método principalprocess(job: Job<T>)handle(event: DomainEvent)
RegistroAutomático por BullMQAutomático por DI (DiscoveryService)
EscuchaCola específicaEventos de entidad + operación
EjecuciónDirecta cuando hay jobIndirecta vía DomainEventProcessor
PropósitoProcesar trabajo asíncronoReaccionar a cambios de dominio
Puede crear jobs✅ Sí✅ Sí
Control de reintentosBullMQ automáticoBullMQ automático (vía processor)
EjemploMeetingTranscriptResponseProcessorSyncMeetingSummaryToCrmSubscriber

🎯 ¿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:

typescript
// 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

  1. Ambos son workers, pero:

    • Processor = Worker directo de BullMQ
    • Subscriber = Worker indirecto vía Domain Events
  2. Ambos pueden crear jobs, pero:

    • Processor crea jobs para continuar flujos de procesamiento
    • Subscriber crea jobs como side-effects de eventos
  3. Ambos procesan trabajo asíncrono, pero:

    • Processor procesa trabajo "técnico" (conversiones, procesamiento)
    • Subscriber procesa trabajo "de negocio" (sincronizaciones, notificaciones)
  4. La diferencia principal es el origen:

    • Processor: Jobs añadidos directamente a la cola
    • Subscriber: Eventos de dominio que se convierten en jobs

📚 Referencias