> Discover all available pages from the documentation index: https://mastra.zisheng.pro/fr/llms.txt # RedisStreamsPubSub `RedisStreamsPubSub` est une implémentation de [`PubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/base) fondée sur [Redis Streams](https://redis.io/docs/latest/develop/data-types/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`](https://mastra.zisheng.pro/fr/reference/pubsub/lease-provider), 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`](https://mastra.zisheng.pro/fr/reference/pubsub/event-emitter). Pour Google Cloud, utilisez [`GoogleCloudPubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/google-cloud-pubsub). 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 **npm**: ```bash npm install @mastra/redis-streams ``` **pnpm**: ```bash pnpm add @mastra/redis-streams ``` **Yarn**: ```bash yarn add @mastra/redis-streams ``` **Bun**: ```bash bun add @mastra/redis-streams ``` ## Exemple d’utilisation Fournissez une URL de connexion Redis. ```typescript 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 **url** (`string`): URL de connexion Redis. Utilise par défaut redisOptions.url. (Default: `redis://localhost:6379`) **keyPrefix** (`string`): Préfixe des clés de flux. Chaque sujet correspond à \:\. (Default: `mastra:topic`) **blockMs** (`number`): Durée, en millisecondes, pendant laquelle chaque lecture reste bloquée dans l’attente de nouveaux événements. (Default: `1000`) **redisOptions** (`RedisClientOptions`): Options transmises au client redis sous-jacent pour une configuration avancée. **maxStreamLength** (`number`): 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. (Default: `10000`) **streamIdleTtlMs** (`number`): 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é). (Default: `0`) **reclaimIntervalMs** (`number`): 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. (Default: `30000`) **reclaimIdleMs** (`number`): 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. (Default: `60000`) **maxDeliveryAttempts** (`number`): Nombre maximal de nouvelles livraisons d’un événement au moyen de nack avant son abandon. Transmettez Infinity pour désactiver cette limite. (Default: `5`) **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 **supportedModes** (`ReadonlyArray<"pull" | "push">`): Renvoie \["pull"]. ## Méthodes `RedisStreamsPubSub` implémente le contrat [`PubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/base). Les méthodes ci-dessous ont un comportement propre à cette implémentation. ### `subscribe(topic, 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é. ```typescript await pubsub.subscribe('workflow.events', (event, ack, nack) => { console.log(event) }) ``` ### `flush()` Attend la fin des publications en cours. ```typescript await pubsub.flush() ``` ### `clearTopic(topic)` 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. ```typescript await pubsub.clearTopic('workflow.events.run-123') ``` ### `close()` Ferme les connexions Redis et arrête tous les abonnements. Appelez cette méthode lors d’un arrêt progressif. ```typescript await pubsub.close() ``` ## 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 `RedisStreamsPubSub` implémente le contrat [`LeaseProvider`](https://mastra.zisheng.pro/fr/reference/pubsub/lease-provider) sur la même connexion Redis. Le [runtime des signaux](https://mastra.zisheng.pro/fr/docs/long-running-agents/signals) 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 `:lease:`. 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`](https://mastra.zisheng.pro/fr/reference/pubsub/lease-provider) pour connaître le contrat complet des méthodes.