EventEmitterPubSub
EventEmitterPubSub est l’implémentation par défaut de PubSub. Elle distribue les événements au sein du processus au moyen d’un EventEmitter Node.js et fonctionne donc sans aucun service externe.
Utilisez-la pour les applications à processus unique. Pour une distribution entre plusieurs processus sur un même hôte, consultez UnixSocketPubSub. Pour une distribution répartie, consultez RedisStreamsPubSub ou GoogleCloudPubSub.
Comme elle fonctionne au sein du processus, les événements ne sont ni conservés ni partagés avec d’autres processus. Encapsulez-la dans CachingPubSub lorsque vous devez relire des événements pour des flux pouvant être repris.
Exemple d’utilisationLien direct vers Exemple d’utilisation
EventEmitterPubSub est utilisé automatiquement lorsque vous ne configurez pas d’option pubsub. La plupart des applications ne l’instancient donc jamais directement. Créez explicitement une instance uniquement si vous souhaitez la configurer ou la partager.
import { Mastra } from '@mastra/core'
import { EventEmitterPubSub } from '@mastra/core/events'
export const mastra = new Mastra({
pubsub: new EventEmitterPubSub(),
})
Pour partager un émetteur avec d’autres parties de votre application, transmettez un EventEmitter existant :
import EventEmitter from 'node:events'
import { EventEmitterPubSub } from '@mastra/core/events'
const emitter = new EventEmitter()
const pubsub = new EventEmitterPubSub(emitter)
Pour signaler les erreurs de distribution par lots, transmettez un logger :
import { EventEmitterPubSub } from '@mastra/core/events'
const pubsub = new EventEmitterPubSub(undefined, { logger })
Paramètres du constructeurLien direct vers Paramètres du constructeur
existingEmitter?:
options?:
PropriétésLien direct vers Propriétés
supportedModes:
["pull", "push"]. L’émetteur peut alimenter un worker de type pull ou transmettre directement les événements aux écouteurs.supportsNativeBatching:
true. Les abonnés peuvent activer la distribution par lots avec options.batch.MéthodesLien direct vers Méthodes
EventEmitterPubSub implémente le contrat PubSub. Les méthodes ci-dessous adoptent un comportement propre à cette implémentation.
subscribe(topic, cb, options?)Lien direct vers subscribetopic-cb-options
Enregistre une fonction de rappel pour un topic. Sans options.group, chaque abonné reçoit chaque événement. Avec un groupe, les événements sont distribués à tour de rôle entre les membres de ce groupe.
Transmettez options.batch pour activer la distribution par lots. Consultez la section Regroupement par lots ci-dessous.
await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})
flush()Lien direct vers flush
Attend le déclenchement de toutes les redistributions en attente provenant de nack avant d’être résolue.
await pubsub.flush()
close()Lien direct vers close
Supprime tous les écouteurs et annule les redistributions en attente. Appelez cette méthode lors d’un arrêt gracieux.
await pubsub.close()
RedistributionLien direct vers Redistribution
Lorsqu’un abonné groupé appelle nack, l’événement est distribué à nouveau au groupe après un court délai et son compteur deliveryAttempt augmente. L’appel de ack efface le suivi de cet événement. Les abonnés en diffusion générale reçoivent des fonctions ack et nack sans effet, puisque chaque événement atteint une fois chaque abonné.
Regroupement par lotsLien direct vers Regroupement par lots
EventEmitterPubSub prend en charge options.batch de façon native. Lorsqu’un abonné active cette option, les événements sont conservés dans un tampon en mémoire propre à cet abonné, puis distribués sous forme d’appels successifs de la fonction de rappel lorsqu’une condition de vidage est satisfaite. Les abonnés en diffusion générale comme les abonnés groupés peuvent utiliser des lots. Consultez SubscribeBatchOptions pour connaître la politique complète.
await pubsub.subscribe(
'workflow.events',
event => {
console.log(event)
},
{
batch: {
maxSize: 10, // flush once 10 events have queued
maxWaitMs: 500, // ...or after 500ms, whichever comes first
},
},
)
Le tampon réside en mémoire et est propre à chaque processus ; l’état des lots n’est donc pas conservé et ne survit pas à un redémarrage. Un vidage déclenché par maxWaitMs s’effectue au mieux : si une étape échoue, par exemple parce qu’un coalesce génère une erreur, celle-ci est signalée par le logger configuré au lieu d’être levée. flush() vide le tampon de chaque abonné utilisant des lots avant d’être résolue.