Aller au contenu principal

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’utilisation
Lien 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.

src/mastra/index.ts
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 constructeur
Lien direct vers Paramètres du constructeur

existingEmitter?:

EventEmitter
EventEmitter Node.js existant à utiliser pour la distribution. S’il est omis, un nouvel EventEmitter est créé.

options?:

EventEmitterPubSubOptions
Configuration facultative.
IMastraLogger

Propriétés
Lien direct vers Propriétés

supportedModes:

ReadonlyArray<"pull" | "push">
Renvoie ["pull", "push"]. L’émetteur peut alimenter un worker de type pull ou transmettre directement les événements aux écouteurs.

supportsNativeBatching:

boolean
Renvoie true. Les abonnés peuvent activer la distribution par lots avec options.batch.

Méthodes
Lien 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()

Redistribution
Lien 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 lots
Lien 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.