CachingPubSub
CachingPubSub encapsule toute implémentation de PubSub et lui ajoute la mise en cache et la relecture des événements. Il enregistre dans un cache chaque événement publié, par sujet. Ainsi, un abonné qui se connecte tardivement ou se reconnecte après une interruption peut relire les événements qu'il a manqués avant de poursuivre avec les événements en direct.
Utilisez-le pour créer des flux pouvant être repris au-dessus d'un transport qui ne conserve pas d'historique, tel que EventEmitterPubSub. Les transports qui conservent déjà les événements, comme RedisStreamsPubSub, n'ont pas besoin de ce wrapper.
CachingPubSub est transparent pour le traitement par lots : subscribe() transmet options.batch au pub/sub interne, et supportsNativeBatching reflète la valeur de ce dernier. Si l'implémentation interne n'accepte pas le traitement par lots, les événements sont distribués individuellement, même lorsque options.batch est fourni.
Exemple d'utilisationLien direct vers Exemple d'utilisation
Encapsulez un pub/sub interne et fournissez un cache côté serveur pour stocker les événements.
import { Mastra } from '@mastra/core'
import { CachingPubSub, EventEmitterPubSub } from '@mastra/core/events'
import { InMemoryServerCache } from '@mastra/core/cache'
const cache = new InMemoryServerCache()
const pubsub = new CachingPubSub(new EventEmitterPubSub(), cache)
export const mastra = new Mastra({
pubsub,
})
La fonction utilitaire withCaching renvoie la même instance et améliore la lisibilité lors d'une encapsulation en ligne :
import { withCaching, EventEmitterPubSub } from '@mastra/core/events'
import { InMemoryServerCache } from '@mastra/core/cache'
const pubsub = withCaching(new EventEmitterPubSub(), new InMemoryServerCache())
Paramètres du constructeurLien direct vers Paramètres du constructeur
inner:
cache:
options?:
PropriétésLien direct vers Propriétés
supportsNativeBatching:
true uniquement lorsque l'implémentation encapsulée prend en charge le traitement par lots.MéthodesLien direct vers Méthodes
CachingPubSub implémente le contrat PubSub. Il redéfinit les méthodes de relecture afin de lire les événements mis en cache. Les méthodes ci-dessous décrivent le comportement de la mise en cache.
publish(topic, event)Lien direct vers publishtopic-event
Met l'événement en cache avec un index séquentiel, puis le publie sur le pub/sub interne.
await pubsub.publish('my-topic', {
type: 'example',
data: { value: 1 },
runId: 'run-123',
})
subscribeWithReplay(topic, cb)Lien direct vers subscribewithreplaytopic-cb
Relit tous les événements mis en cache pour le sujet, puis s'abonne aux événements en direct.
await pubsub.subscribeWithReplay('my-topic', event => {
console.log(event)
})
subscribeFromOffset(topic, offset, cb)Lien direct vers subscribefromoffsettopic-offset-cb
Relit les événements mis en cache à partir de offset, puis s'abonne aux événements en direct. Utilisez cette méthode lorsque le client connaît sa dernière position afin d'éviter de relire tout l'historique.
await pubsub.subscribeFromOffset('my-topic', 42, event => {
console.log(event)
})
getHistory(topic, offset?)Lien direct vers gethistorytopic-offset
Renvoie les événements mis en cache pour le sujet à partir de offset.
const events = await pubsub.getHistory('my-topic', 0)
Renvoie : Promise<Event[]>
clearTopic(topic)Lien direct vers cleartopictopic
Supprime les événements mis en cache et le compteur de décalage du sujet, puis transmet l'appel au pub/sub interne afin que les transports persistants (tels que RedisStreamsPubSub) puissent également supprimer l'état qu'ils conservent. Les cycles de vie des exécutions de Mastra appellent automatiquement cette méthode à la fin d'une exécution.
await pubsub.clearTopic('my-topic')
FonctionsLien direct vers Fonctions
withCaching(pubsub, cache, options?)Lien direct vers withcachingpubsub-cache-options
Wrapper pratique qui construit un CachingPubSub. Il accepte les mêmes arguments que le constructeur et renvoie la nouvelle instance.
import { withCaching, EventEmitterPubSub } from '@mastra/core/events'
import { InMemoryServerCache } from '@mastra/core/cache'
const pubsub = withCaching(new EventEmitterPubSub(), new InMemoryServerCache(), {
keyPrefix: 'events:',
})
Renvoie : CachingPubSub