Aller au contenu principal

RedisStreamsPubSub

RedisStreamsPubSub est une implémentation de PubSub fondée sur Redis Streams. Elle distribue les événements entre les processus et les hôtes, avec persistance, groupes de consommateurs et nouvelle livraison en cas d’échec. Elle implémente également LeaseProvider, ce qui permet à la couche de signaux d’élire un seul propriétaire par ressource parmi toutes les instances. Les signaux peuvent ainsi coordonner les exécutions dans les déploiements distribués et serverless.

Utilisez-la pour les déploiements distribués dans lesquels plusieurs services partagent un flux d’événements. Pour une distribution au sein d’un seul processus, utilisez EventEmitterPubSub. Pour Google Cloud, utilisez GoogleCloudPubSub.

Chaque sujet correspond à une clé de flux Redis. Les abonnements associés à un groupe utilisent un groupe de consommateurs Redis, de sorte que ses membres se répartissent le travail à tour de rôle. Les abonnements sans groupe créent un groupe de consommateurs privé, ce qui permet à chaque abonné de recevoir tous les événements.

RedisStreamsPubSub est un transport de type pull : les consommateurs lisent les événements avec XREADGROUP. Mastra exécute donc un worker d’orchestration chargé de les lire en leur nom.

Installation
Lien direct vers Installation

npm install @mastra/redis-streams

Exemple d’utilisation
Lien direct vers Exemple d’utilisation

Fournissez une URL de connexion Redis.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})

Paramètres du constructeur
Lien direct vers Paramètres du constructeur

url?:

string
= redis://localhost:6379
URL de connexion Redis. Utilise par défaut redisOptions.url.

keyPrefix?:

string
= mastra:topic
Préfixe des clés de flux. Chaque sujet correspond à <keyPrefix>:<topic>.

blockMs?:

number
= 1000
Durée, en millisecondes, pendant laquelle chaque lecture reste bloquée dans l’attente de nouveaux événements.

redisOptions?:

RedisClientOptions
Options transmises au client redis sous-jacent pour une configuration avancée.

maxStreamLength?:

number
= 10000
Nombre maximal approximatif d’entrées conservées par flux. Définissez cette valeur sur 0 pour désactiver la suppression des anciennes entrées.

streamIdleTtlMs?:

number
= 0
Expiration après inactivité, en millisecondes : TTL glissant actualisé à chaque écriture (publication, nouvelle tentative après nack, recréation d’un groupe). Chaque écriture le réinitialise ; un flux alimenté activement n’expire donc jamais en cours d’utilisation, tandis qu’un flux inactif pendant toute cette durée est automatiquement supprimé par Redis. Notez que seules les écritures actualisent le TTL : un consommateur qui traite lentement un backlog ne le fait pas. Définissez donc cette valeur bien au-dessus du plus long intervalle prévu entre deux écritures sur un sujet actif. Il s’agit d’un filet de sécurité, et non du mécanisme de nettoyage principal : clearTopic assure la suppression normale en fin de cycle de vie ; cette option limite uniquement l’utilisation de la mémoire pour les flux qui n’atteignent jamais un appel à clearTopic (par exemple, une exécution interrompue brutalement). La valeur doit être un entier positif ou nul. Valeur par défaut : 0 (désactivé).

reclaimIntervalMs?:

number
= 30000
Fréquence, en millisecondes, à laquelle un abonnement récupère les événements qu’un consommateur précédent a lus sans jamais en accuser réception. Définissez cette valeur sur 0 pour désactiver ce comportement.

reclaimIdleMs?:

number
= 60000
Durée minimale d’inactivité, en millisecondes, avant qu’un événement en attente puisse être récupéré. Définissez-la bien au-dessus de la durée de traitement habituelle afin d’éviter une double livraison.

maxDeliveryAttempts?:

number
= 5
Nombre maximal de nouvelles livraisons d’un événement au moyen de nack avant son abandon. Transmettez Infinity pour désactiver cette limite.

logger?:

{ debug?: Function; warn?: Function }
Logger facultatif pour les diagnostics. S’il est omis, les erreurs supprimées ne produisent aucun message.

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

supportedModes:

ReadonlyArray<"pull" | "push">
Renvoie ["pull"].

Méthodes
Lien direct vers Méthodes

RedisStreamsPubSub implémente le contrat PubSub. Les méthodes ci-dessous ont un comportement propre à cette implémentation.

subscribe(topic, cb, options?)
Lien direct vers subscribetopic-cb-options

S’abonne à un sujet. Avec options.group, les membres du groupe se partagent les événements au moyen d’un groupe de consommateurs Redis. Sans groupe, l’abonné reçoit tous les événements par l’intermédiaire d’un groupe de consommateurs privé.

await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})

flush()
Lien direct vers flush

Attend la fin des publications en cours.

await pubsub.flush()

clearTopic(topic)
Lien direct vers cleartopictopic

Supprime le flux d’un sujet et tous ses groupes de consommateurs, libérant ainsi la mémoire qu’un sujet terminé occuperait autrement. Les cycles de vie des exécutions de Mastra, qu’il s’agisse d’Agents durables ou du moteur de Workflow événementiel, appellent automatiquement cette méthode lorsqu’une exécution atteint un état terminal. Ne l’appelez vous-même qu’une fois que plus rien ne doit lire le sujet. Cette opération s’effectue au mieux et ne lève jamais d’erreur. Les échecs sont consignés au niveau warn. Un abonné encore connecté au moment de la suppression du flux se rétablit de lui-même, mais ne reçoit pas les entrées supprimées.

Le nettoyage automatique nécessite des versions de @mastra/core et @mastra/redis-streams qui prennent toutes deux en charge clearTopic : le runtime achemine l’appel par sa couche de cache. Mettez donc à niveau les deux packages ensemble afin de bénéficier de la suppression du flux en fin d’exécution.

await pubsub.clearTopic('workflow.events.run-123')

close()
Lien direct vers close

Ferme les connexions Redis et arrête tous les abonnements. Appelez cette méthode lors d’un arrêt progressif.

await pubsub.close()

Nouvelle livraison et récupération
Lien direct vers Nouvelle livraison et récupération

Lorsqu’un abonné appelle nack, l’événement est republié avec une valeur deliveryAttempt incrémentée et la réception de l’original est confirmée. Dès qu’un événement atteint maxDeliveryAttempts, il est abandonné au lieu d’être livré de nouveau. Par ailleurs, chaque abonnement récupère périodiquement les événements qu’un consommateur précédent du groupe a lus sans jamais en accuser réception, selon les paramètres reclaimIntervalMs et reclaimIdleMs.

Baux distribués
Lien direct vers Baux distribués

RedisStreamsPubSub implémente le contrat LeaseProvider sur la même connexion Redis. Le runtime des signaux l’utilise pour élire un propriétaire unique, généralement pour chaque clé de fil. Ainsi, parmi toutes les instances, un seul processus se réveille et exécute l’Agent, tandis que les autres acheminent le travail de suivi vers le détenteur. C’est ce qui permet aux signaux de fonctionner dans les déploiements serverless et multi-instances ; sans bail partagé, chaque instance lancerait sa propre exécution concurrente.

Les clés de bail sont placées dans le même espace de noms keyPrefix que les sujets, sous la forme <keyPrefix>:lease:<key>. Toutes les opérations sont atomiques : acquireLease utilise SET NX PX et actualise son propre TTL de manière idempotente, tandis que releaseLease, renewLease et transferLease utilisent des scripts Lua qui vérifient la propriété avant toute modification. Ainsi, un renouvellement concurrent effectué par un autre propriétaire n’est jamais écrasé.

Vous n’appelez pas ces méthodes directement. Il suffit de configurer RedisStreamsPubSub comme backend pubsub pour que le runtime détecte et utilise cette fonctionnalité. Consultez LeaseProvider pour connaître le contrat complet des méthodes.