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.
InstallationLien direct vers Installation
- npm
- pnpm
- Yarn
- Bun
npm install @mastra/redis-streams
pnpm add @mastra/redis-streams
yarn add @mastra/redis-streams
bun add @mastra/redis-streams
Exemple d’utilisationLien direct vers Exemple d’utilisation
Fournissez une URL de connexion Redis.
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 constructeurLien direct vers Paramètres du constructeur
url?:
redisOptions.url.keyPrefix?:
<keyPrefix>:<topic>.blockMs?:
redisOptions?:
redis sous-jacent pour une configuration avancée.maxStreamLength?:
streamIdleTtlMs?:
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?:
reclaimIdleMs?:
maxDeliveryAttempts?:
nack avant son abandon. Transmettez Infinity pour désactiver cette limite.logger?:
PropriétésLien direct vers Propriétés
supportedModes:
["pull"].MéthodesLien 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érationLien 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ésLien 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.