PubSub
Mastra utilise un système de publication/abonnement (pub/sub) comme bus d’événements interne. Les composants publient des événements dans des rubriques, auxquelles d’autres composants s’abonnent pour y réagir. Le backend que vous configurez détermine jusqu’où ces événements circulent : dans un même processus, entre des processus sur un même hôte, ou entre des instances distinctes.
Vous définissez le backend une seule fois sur l’instance Mastra, puis le reste du système l’utilise sans modification. Par défaut, Mastra utilise un backend intégré au processus qui ne nécessite aucune configuration.
Comment Mastra utilise PubSubLien direct vers Comment Mastra utilise PubSub
Plusieurs systèmes intégrés publient et s’abonnent à des événements via le même bus pub/sub :
- Exécution des Workflows : le planificateur publie un événement
workflow.start, qu’un worker de longue durée consomme pour exécuter le Workflow. Les événements d’étape et de cycle de vie circulent via pub/sub pendant l’exécution. - Workflows planifiés : le planificateur distribue les exécutions arrivées à échéance en publiant dans la rubrique du Workflow. Consultez les Workflows planifiés.
- Tâches en arrière-plan : un gestionnaire de tâches distribue le travail à un groupe de workers et diffuse les mises à jour du cycle de vie des tâches aux abonnés ; c’est ainsi que les flux de tâches en arrière-plan restent actifs. Consultez la diffusion de tâches en arrière-plan.
- Signaux d’Agent : l’envoi d’un signal à une exécution d’Agent active publie un événement dans une rubrique de thread afin qu’une exécution dans un autre processus reçoive ce signal.
- Flux reprenables : les segments de flux sont publiés par exécution, afin qu’un client qui se reconnecte puisse rejouer ce qu’il a manqué.
Comme ces systèmes utilisent un seul bus, le backend que vous choisissez s’applique à tous simultanément.
Modes de diffusionLien direct vers Modes de diffusion
Les backends diffusent les événements selon l’un des deux modes définis par le contrat PubSub :
- Pull : les consommateurs lisent eux-mêmes le backend, ce que Mastra fait au moyen d’une boucle de worker de longue durée. Les backends distribués, tels que
RedisStreamsPubSub, utilisent ce mode. - Push : les événements arrivent sans que le consommateur les demande, dans le processus ou via HTTP. Le backend
EventEmitterPubSubpar défaut assure cette diffusion dans le processus.
Les abonnés peuvent également répartir le travail avec des groupes de consommateurs. Les membres d’un même groupe se partagent les événements afin que chacun soit traité une seule fois. Un abonné sans groupe reçoit chaque événement, ce qui diffuse le flux à tous les abonnés non regroupés.
Backend par défautLien direct vers Backend par défaut
Lorsque vous ne définissez pas l’option pubsub, Mastra utilise EventEmitterPubSub. Il diffuse les événements dans le processus à l’aide d’un EventEmitter Node.js ; aucun service externe n’est donc requis.
import { Mastra } from '@mastra/core'
// No pubsub option: Mastra uses EventEmitterPubSub
export const mastra = new Mastra({})
Comme il s’exécute dans le processus, les événements ne sont pas persistés et n’atteignent pas les autres processus. La configuration par défaut convient à une instance unique, ce qui couvre la plupart des applications.
Choisir un backendLien direct vers Choisir un backend
Définissez l’option pubsub sur l’instance Mastra pour choisir un backend. Chaque backend implémente le même contrat PubSub, donc le reste de votre application ne change pas.
Le critère déterminant est l’emplacement où s’exécute l’abonné.
| Backend | Portée | Mode | Package |
|---|---|---|---|
EventEmitterPubSub | Processus unique | Pull et push | @mastra/core |
UnixSocketPubSub | Plusieurs processus, un hôte | Push | @mastra/core |
RedisStreamsPubSub | Distribué, plusieurs hôtes | Pull | @mastra/redis-streams |
GoogleCloudPubSub | Distribué, plusieurs hôtes | Pull | @mastra/google-cloud-pubsub |
Plusieurs processus sur un même hôteLien direct vers Plusieurs processus sur un même hôte
Utilisez UnixSocketPubSub lorsque plusieurs processus sur la même machine doivent partager un flux. Il diffuse les événements via un socket de domaine Unix et élit un processus comme broker. Si ce broker se termine, les processus restants en élisent un nouveau.
import { Mastra } from '@mastra/core'
import { UnixSocketPubSub } from '@mastra/core/events'
export const mastra = new Mastra({
pubsub: new UnixSocketPubSub('/tmp/mastra/events.sock'),
})
Déploiements distribuésLien direct vers Déploiements distribués
Utilisez un backend distribué lorsque vous exécutez plus d’une instance ou d’un hôte, afin que chaque instance reçoive les mêmes événements. Cela est nécessaire dès qu’une requête traitée par une instance doit atteindre un travail exécuté par une autre.
Par exemple, envoyer un signal à une exécution d’Agent exige que l’événement de signal franchisse la limite de processus jusqu’à l’instance propriétaire de l’exécution. Avec le backend intégré au processus par défaut, cette instance ne reçoit jamais l’événement.
Les deux backends ci-dessous diffusent entre processus et hôtes et conservent les événements afin de pouvoir les rediffuser.
L’exemple suivant utilise Redis Streams :
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'
export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})
L’exemple suivant utilise Google Cloud Pub/Sub :
import { Mastra } from '@mastra/core'
import { GoogleCloudPubSub } from '@mastra/google-cloud-pubsub'
export const mastra = new Mastra({
pubsub: new GoogleCloudPubSub({
projectId: 'my-project',
}),
})
Flux reprenablesLien direct vers Flux reprenables
Les flux reprenables permettent à un client de se reconnecter et de rejouer les événements qu’il a manqués ; le backend doit donc conserver un historique récent. Les backends distribués tels que RedisStreamsPubSub persistent les événements et prennent donc en charge le rejeu nativement.
La diffusion dans le processus ne conserve pas l’historique. Pour ajouter le rejeu à EventEmitterPubSub, encapsulez-le dans CachingPubSub, qui enregistre les événements publiés pour chaque rubrique afin qu’un abonné tardif ou qui se reconnecte puisse se remettre à jour avant de poursuivre avec les événements en direct.
import { Mastra } from '@mastra/core'
import { CachingPubSub, EventEmitterPubSub } from '@mastra/core/events'
import { InMemoryServerCache } from '@mastra/core/cache'
const cache = new InMemoryServerCache()
export const mastra = new Mastra({
pubsub: new CachingPubSub(new EventEmitterPubSub(), cache),
})
Consultez la référence PubSub pour connaître le contrat de diffusion complet et les options de configuration de chaque backend.
Ressources associéesLien direct vers Ressources associées
- Référence PubSub
- Classe Mastra
- Workers : exécutez l’orchestration de Workflows et les tâches en arrière-plan dans des processus dédiés avec PubSub
- Diffusion de tâches en arrière-plan
- Workflows planifiés