
Redis Stream — un nouveau type de donnée abstrait introduit dans Redis avec la version 5.0
Conceptuellement, Redis Stream est une liste dans laquelle vous pouvez ajouter des enregistrements. Chaque enregistrement a un identifiant unique. Par défaut, l'identifiant est généré automatiquement et comprend un horodatage. Vous pouvez donc interroger des plages d'enregistrements par temps ou recevoir de nouvelles données au fur et à mesure de leur arrivée dans le flux, tout comme la commande Unix « tail -f » lit un fichier journal et reste en attente de nouvelles données. Notez que plusieurs clients peuvent écouter le flux simultanément, tout comme plusieurs processus « tail -f » peuvent lire un fichier en même temps, sans conflit.
Pour comprendre tous les avantages du nouveau type de données, rappelons brièvement les structures Redis existantes qui répètent partiellement la fonctionnalité de Redis Stream.
Redis PUB/SUB
Redis Pub/Sub — un système de messagerie simple déjà intégré dans votre stockage key-value. Cependant, cette simplicité a un coût :
- Si l'éditeur tombe en panne pour une raison quelconque, il perd tous ses abonnés
- L'éditeur doit connaître l'adresse exacte de tous ses abonnés
- L'éditeur peut surcharger ses abonnés si les données sont publiées plus rapidement qu'elles ne peuvent être traitées
- Le message est supprimé du tampon de l'éditeur immédiatement après sa publication, peu importe combien d'abonnés l'ont reçu et à quelle vitesse ils ont pu traiter ce message.
- Tous les abonnés recevront le message en même temps. Les abonnés doivent eux-mêmes se mettre d'accord sur l'ordre de traitement du même message.
- Il n'y a pas de mécanisme intégré pour confirmer le traitement réussi d'un message par un abonné. Si un abonné reçoit un message et échoue pendant le traitement, l'éditeur ne le saura pas.
Redis List
Redis List — une structure de données qui supporte les commandes de lecture avec blocage. Vous pouvez ajouter et lire des messages depuis le début ou la fin de la liste. Sur la base de cette structure, vous pouvez créer une bonne pile ou une queue pour votre système distribué, et cela sera suffisant dans la plupart des cas. Les principales différences avec Redis Pub/Sub :
- Le message est délivré à un seul client. Le premier client en attente de lecture recevra les données en premier.
- Clint doit initier lui-même l'opération de lecture de chaque message. List ne sait rien des clients.
- Les messages sont stockés tant que quelqu'un ne les lit pas ou ne les supprime pas explicitement. Si vous avez configuré le serveur Redis pour qu'il écrase les données sur le disque, la fiabilité du système augmente considérablement.
Introduction à Stream
Ajout d'un enregistrement au flux
Commande XADD ajoute un nouvel enregistrement au flux. Un enregistrement n'est pas juste une chaîne, il se compose d'une ou plusieurs paires clé-valeur. Chaque enregistrement est donc déjà structuré et ressemble à la structure d'un fichier CSV.
> XADD mystream * sensor-id 1234 temperature 19.8
1518951480106-0
Dans l'exemple ci-dessus, nous ajoutons au flux nommé (clé) «mystream» deux champs : «sensor-id» et «temperature» avec les valeurs «1234» et «19.8» respectivement. Comme deuxième argument, la commande prend un identifiant qui sera attribué à l'enregistrement - cet identifiant identifie de manière unique chaque enregistrement dans le flux. Cependant, dans ce cas, nous avons passé *, car nous voulons que Redis génère un nouvel identifiant pour nous. Chaque nouvel identifiant sera en augmentation. Ainsi, chaque nouvel enregistrement aura un identifiant plus grand par rapport aux enregistrements précédents.
Format de l'identifiant
L'identifiant de l'enregistrement, retourné par la commande XADD, se compose de deux parties :
{millisecondsTime}-{sequenceNumber}
millisecondesTime — temps Unix en millisecondes (temps de serveurs Redis). Cependant, si l'heure actuelle est égale ou inférieure à l'heure de l'enregistrement précédent, l'horodatage de l'enregistrement précédent est utilisé. Ainsi, si l'heure du serveur revient dans le passé, le nouvel identifiant conservera toujours sa propriété d'incrémentation.
sequenceNumber est utilisé pour les enregistrements créés dans la même milliseconde. sequenceNumber sera augmenté de 1 par rapport à l'enregistrement précédent. Étant donné que sequenceNumber fait 64 bits, dans la pratique, vous ne devriez pas atteindre la limitation du nombre d'enregistrements pouvant être générés en une milliseconde.
Le format de tels identifiants peut sembler étrange au premier abord. Un lecteur sceptique peut se demander pourquoi le temps fait partie de l'identifiant. La raison en est que les flux Redis prennent en charge les requêtes de plage par identifiants. Étant donné que l'identifiant est associé à la création de l'enregistrement, cela permet de demander des plages de temps. Nous examinerons un exemple concret lorsque nous aborderons la commande. XRANGE.
Si, pour une raison quelconque, l'utilisateur doit spécifier son propre identifiant, par exemple lié à un système externe, nous pouvons le transmettre à la commande XADD au lieu du signe * comme indiqué ci-dessous :
> XADD somestream 0-1 field value
0-1
> XADD somestream 0-2 foo bar
0-2
Notez qu'à ce stade, vous devez vous assurer de l'augmentation de l'identifiant. Dans notre exemple, l'identifiant minimal est «0-1», donc la commande n'acceptera pas un identifiant qui est égal ou inférieur à «0-1».
> XADD somestream 0-1 foo bar
(error) ERR L'ID spécifié dans XADD est égal ou inférieur au premier élément du flux cible
Le nombre d'enregistrements dans le flux
Vous pouvez obtenir le nombre d'enregistrements dans le flux simplement en utilisant la commande XLEN. Pour notre exemple, cette commande renverra la valeur suivante :
> XLEN somestream
(integer) 2
Requêtes par intervalle — XRANGE et XREVRANGE
Pour demander des données par intervalle, nous devons spécifier deux identifiants — le début et la fin de l'intervalle. La plage retournée inclura tous les éléments, y compris les limites. Il existe également deux identifiants spéciaux «-» et «+», représentant respectivement le plus petit (premier enregistrement) et le plus grand (dernier enregistrement) identifiant dans le flux. L'exemple ci-dessous affichera tous les enregistrements du flux.
> XRANGE mystream - +
1) 1) 1518951480106-0
2) 1) "sensor-id"
2) "1234"
3) "temperature"
4) "19.8"
2) 1) 1518951482479-0
2) 1) "sensor-id"
2) "9999"
3) "temperature"
4) "18.2"
Chaque enregistrement retourné représente un tableau de deux éléments : l'identifiant et une liste de paires clé-valeur. Nous avons déjà mentionné que les identifiants d'enregistrement ont une relation avec le temps. Par conséquent, nous pouvons demander une plage d'un intervalle de temps spécifique. Cependant, nous pouvons spécifier dans la requête un identifiant incomplet, en omettant la partie relative à. sequenceNumberLa partie omise de l'identifiant sera automatiquement ramenée à zéro au début de la plage et au maximum à la valeur possible à la fin de la plage. Voici un exemple de la manière dont vous pouvez demander une plage de deux millisecondes.
> XRANGE mystream 1518951480106 1518951480107
1) 1) 1518951480106-0
2) 1) "sensor-id"
2) "1234"
3) "temperature"
4) "19.8"
Nous n'avons qu'un seul enregistrement dans cette plage, mais dans des ensembles de données réels, le résultat retourné peut être énorme. Pour cette raison, XRANGE la fonction COUNT est prise en charge. En spécifiant un nombre, nous pouvons simplement obtenir les premiers N enregistrements. Si nous devons obtenir les N enregistrements suivants (pagination), nous pouvons utiliser l'identifiant obtenu précédemment, l'augmenter sequenceNumber de un et demander à nouveau. Regardons cela dans l'exemple suivant. Nous commençons à ajouter 10 éléments avec XADD (supposons que le flux mystream ait déjà été rempli de 10 éléments). Pour commencer l'itération en récupérant 2 éléments par commande, nous commençons par la plage complète, mais avec COUNT égal à 2.
> XRANGE mystream - + COUNT 2
1) 1) 1519073278252-0
2) 1) "foo"
2) "value_1"
2) 1) 1519073279157-0
2) 1) "foo"
2) "value_2"
Pour continuer l'itération avec les deux éléments suivants, nous devons prendre le dernier identifiant obtenu, soit 1519073279157-0, et y ajouter 1 à sequenceNumber.
l'identifiant résultant, dans ce cas 1519073279157-1, qui peut maintenant être utilisé comme un nouvel argument de début de plage pour le prochain appel. XRANGE:
> XRANGE mystream 1519073279157-1 + COUNT 2
1) 1) 1519073280281-0
2) 1) "foo"
2) "value_3"
2) 1) 1519073281432-0
2) 1) "foo"
2) "value_4"
Et ainsi de suite. Étant donné que la complexité XRANGE est O(log (N)) pour la recherche, puis O(M) pour le retour de M éléments, chaque étape de l'itération est rapide. Ainsi, avec XRANGE il est possible d'itérer efficacement sur les flux.
Commande XREVRANGE est l'équivalent XRANGE, mais renvoie les éléments dans l'ordre inverse :
> XREVRANGE mystream + - COUNT 1
1) 1) 1519073287312-0
2) 1) "foo"
2) "value_10"
Notez que la commande XREVRANGE prend les arguments de plage start et stop dans l'ordre inverse.
Lire de nouveaux enregistrements avec XREAD
Il est souvent nécessaire de s'abonner à un flux et de ne recevoir que de nouveaux messages. Ce concept peut sembler similaire à Redis Pub/Sub ou à la liste Redis bloquante, mais il existe des différences fondamentales dans la façon d'utiliser Redis Stream :
- Chaque nouveau message est par défaut livré à chaque abonné. Ce comportement est différent d'une liste Redis bloquante, où un nouveau message n'est lu que par un seul abonné.
- Alors que dans Redis Pub/Sub tous les messages sont oubliés et jamais sauvegardés, dans Stream tous les messages sont conservés indéfiniment (sauf si le client les supprime explicitement).
- Redis Stream permet de restreindre l'accès aux messages au sein d'un même flux. Un abonné spécifique peut voir uniquement son propre historique de messages.
Vous pouvez vous abonner à un flux et recevoir de nouveaux messages en utilisant la commande XREAD. Cela est un peu plus complexe que XRANGE, donc nous allons d'abord commencer par des exemples plus simples.
> XREAD COUNT 2 STREAMS mystream 0
1) 1) "mystream"
2) 1) 1) 1519073278252-0
2) 1) "foo"
2) "value_1"
2) 1) 1519073279157-0
2) 1) "foo"
2) "value_2"
Dans l'exemple ci-dessus, une forme non bloquante est indiquée. XREADNotez que l'option COUNT n'est pas obligatoire. En fait, la seule option obligatoire de la commande est l'option STREAMS, qui spécifie la liste des flux avec l'identifiant maximal correspondant. Nous avons écrit «STREAMS mystream 0» — nous voulons recevoir toutes les entrées du flux mystream avec un identifiant supérieur à «0-0». Comme le montre l'exemple, la commande retourne le nom du flux, car nous pouvons nous abonner à plusieurs flux simultanément. Nous pourrions écrire, par exemple, «STREAMS mystream otherstream 0 0». Notez qu'après l'option STREAMS, nous devons d'abord fournir les noms de tous les flux nécessaires, puis la liste des identifiants.
Dans cette forme simple, la commande ne fait rien de particulier par rapport à XRANGE. Cependant, ce qui est intéressant, c'est que nous pouvons facilement transformer XREAD en une commande bloquante en spécifiant l'argument BLOCK :
> XREAD BLOCK 0 STREAMS mystream $
Dans l'exemple ci-dessus, une nouvelle option BLOCK avec un délai d'attente de 0 millisecondes est indiquée (ce qui signifie attente indéfinie). De plus, au lieu de transmettre un identifiant normal pour le flux mystream, un identifiant spécial $ a été transmis. Cet identifiant spécial signifie que XREAD doit utiliser l'identifiant maximal du flux mystream. Ainsi, nous ne recevrons que les nouveaux messages à partir du moment où nous avons commencé à écouter. D'une certaine manière, cela ressemble à la commande Unix «tail -f».
Veuillez noter qu'en utilisant l'option BLOCK, il n'est pas nécessaire d'utiliser un identifiant spécial $. Nous pouvons utiliser n'importe quel identifiant existant dans le flux. Si l'équipe peut traiter notre demande immédiatement, sans blocage, elle le fera, sinon elle se bloquera.
Bloquant XREAD peut également écouter plusieurs flux simultanément, il suffit de spécifier leurs noms. Dans ce cas, l'équipe renverra l'enregistrement du premier flux dans lequel des données sont arrivées. Le premier abonné bloqué pour ce flux recevra les données en premier.
Groupes de Consommateurs
Dans certaines tâches, nous souhaitons restreindre l'accès des abonnés aux messages au sein d'un même flux. Un exemple où cela peut être utile est une file de messages avec des travailleurs qui recevront différents messages du flux, permettant ainsi d'échelonner le traitement des messages.
Imaginons que nous ayons trois abonnés C1, C2, C3 et un flux contenant les messages 1, 2, 3, 4, 5, 6, 7, le traitement des messages se fera comme indiqué dans le diagramme ci-dessous :
1 -> C1
2 -> C2
3 -> C3
4 -> C1
5 -> C2
6 -> C3
7 -> C1
Pour obtenir cet effet, Redis Stream utilise un concept appelé Groupe de Consommateurs. Ce concept est similaire à un pseudo-abonné qui reçoit des données du flux, mais est en réalité servi par plusieurs abonnés au sein du groupe, fournissant certaines garanties.
- Chaque message est délivré à différents abonnés au sein du groupe.
- Au sein du groupe, les abonnés sont identifiés par un nom, qui est une chaîne sensible à la casse. Si un abonné sort temporairement du groupe, il peut y revenir avec son propre nom unique.
- Chaque Groupe de Consommateurs suit le principe du « premier message non lu ». Lorsque l'abonné demande de nouveaux messages, il ne peut recevoir que ceux qui n'ont jamais été livrés à aucun abonné au sein du groupe.
- Il existe une commande pour confirmer explicitement le traitement réussi d'un message par un abonné. Tant que cette commande n'est pas appelée, le message demandé restera dans un état « en attente ».
- Au sein du Groupe de Consommateurs, chaque abonné peut demander l'historique des messages qui lui ont été livrés mais qui n'ont pas encore été traités (dans l'état « en attente »).
Dans un certain sens, l'état du groupe peut être représenté ainsi :
+----------------------------------------+
| consumer_group_name: mygroup
| consumer_group_stream: somekey
| last_delivered_id: 1292309234234-92
|
| consumers:
| "consumer-1" avec des messages en attente
| 1292309234234-4
| 1292309234232-8
| "consumer-42" avec des messages en attente
| ... (et ainsi de suite)
+----------------------------------------+
Il est maintenant temps de se familiariser avec les principales commandes pour le groupe de consommateurs, à savoir :
- XGROUP est utilisé pour créer, détruire et gérer des groupes
- XREADGROUP est utilisé pour lire le flux à travers le groupe
- XACK — cette commande permet au souscripteur de marquer un message comme traité avec succès
Création du groupe de consommateurs
Supposons que le flux mystream existe déjà. Alors, la commande de création du groupe sera la suivante :
> CRÉER XGROUP mystream mygroup $
OK
Lors de la création du groupe, nous devons passer l'identifiant à partir duquel le groupe commencera à recevoir des messages. Si nous voulons simplement recevoir tous les nouveaux messages, nous pouvons utiliser l'identifiant spécial $ (comme dans notre exemple ci-dessus). Si, au lieu de l'identifiant spécial, nous spécifions 0, le groupe aura accès à tous les messages du flux.
Maintenant que le groupe est créé, nous pouvons immédiatement commencer à lire les messages avec la commande XREADGROUP. Cette commande est très similaire à XREAD et supporte l'option facultative BLOCK. Cependant, il y a une option obligatoire GROUP qui doit toujours être indiquée avec deux arguments : le nom du groupe et le nom du souscripteur. L'option COUNT est également supportée.
Avant de lire le flux, ajoutons quelques messages :
> XADD mystream * message apple
1526569495631-0
> XADD mystream * message orange
1526569498055-0
> XADD mystream * message strawberry
1526569506935-0
> XADD mystream * message apricot
1526569535168-0
> XADD mystream * message banana
1526569544280-0
Et maintenant, essayons de lire ce flux à travers le groupe :
> XREADGROUP GROUP mygroup Alice COUNT 1 STREAMS mystream >
1) 1) "mystream"
2) 1) 1) 1526569495631-0
2) 1) "message"
2) "apple"
La commande ci-dessus se lit littéralement comme suit :
« Moi, Alice-le-souscripteur, membre du groupe mygroup, souhaite lire dans le flux mystream un message qui n'a jamais été livré à personne auparavant.»
Chaque fois qu'un abonné effectue une opération avec le groupe, il doit indiquer son nom, s'identifiant ainsi de manière unique au sein du groupe. Dans la commande ci-dessus, il y a un autre détail très important : l'identifiant spécial « > ». Cet identifiant spécial filtre les messages, ne laissant passer que ceux qui n'ont pas encore été livrés.
De plus, dans des cas particuliers, vous pouvez spécifier un identifiant réel tel que 0 ou tout autre identifiant valide. Dans ce cas, la commande XREADGROUP vous renverra l'historique des messages avec le statut « en attente », qui ont été livrés à l'abonné spécifié (Alice), mais qui n'ont pas encore été confirmés via la commande XACK.
Nous pouvons vérifier ce comportement en spécifiant immédiatement l'identifiant 0, sans l'option COUNT. Nous verrons simplement le seul message en attente, c'est-à-dire le message avec la pomme :
> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
2) 1) 1) 1526569495631-0
2) 1) "message"
2) "pomme"
Cependant, si nous confirmons le message comme ayant été traité avec succès, il n'apparaîtra plus :
> XACK mystream mygroup 1526569495631-0
(integer) 1
> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
2) (liste ou ensemble vide)
Il est maintenant temps pour Bob de lire quelque chose :
> XREADGROUP GROUP mygroup Bob COUNT 2 STREAMS mystream >
1) 1) "mystream"
2) 1) 1) 1526569498055-0
2) 1) "message"
2) "orange"
2) 1) 1526569506935-0
2) 1) "message"
2) "fraise"
Bob, membre du groupe mygroup, a demandé pas plus de deux messages. La commande ne rapporte que les messages non livrés à cause de l'identifiant spécial « > ». Comme vous pouvez le voir, le message « pomme » n'apparaît pas, car il a déjà été livré à Alice, donc Bob reçoit « orange » et « fraise ».
Ainsi, Alice, Bob et tout autre abonné du groupe peuvent lire différents messages d'un même flux. Ils peuvent également consulter leur historique de messages non traités ou marquer des messages comme traités.
Il y a quelques points à prendre en compte :
- Dès qu'un abonné lit le message avec la commande XREADGROUP, ce message passe à l'état « en attente » et est attribué à cet abonné spécifique. Les autres abonnés du groupe ne pourront pas lire ce message.
- Les abonnés sont créés automatiquement dès la première mention, il n'est pas nécessaire de les créer explicitement.
- Avec XREADGROUP Vous pouvez lire les messages de plusieurs flux simultanément, mais pour que cela fonctionne, vous devez d'abord créer des groupes avec le même nom pour chaque flux à l'aide de XGROUP
Récupération après sinistre
Un abonné peut se rétablir après une panne et relire sa liste de messages avec le statut « en attente ». Cependant, dans le monde réel, les abonnés peuvent échouer définitivement. Que se passe-t-il avec les messages suspendus d'un abonné s'il ne parvient pas à se rétablir après la panne ?
Le Consumer Group propose une fonction qui est utilisée précisément dans de tels cas—lorsqu'il est nécessaire de changer le propriétaire des messages.
Tout d'abord, il est nécessaire d'appeler la commande XPENDING, qui affiche tous les messages du groupe avec le statut « en attente ». Dans sa forme la plus simple, la commande est appelée avec seulement deux arguments : le nom du flux et le nom du groupe :
> XPENDING mystream mygroup
1) (integer) 2
2) 1526569498055-0
3) 1526569506935-0
4) 1) 1) "Bob"
2) "2"
La commande a affiché le nombre de messages non traités pour l'ensemble du groupe et pour chaque abonné. Nous avons seulement Bob avec deux messages non traités, car le seul message demandé par Alice a été confirmé avec XACK.
Nous pouvons demander des informations supplémentaires en utilisant plus d'arguments :
XPENDING {key} {groupname} [{start-id} {end-id} {count} [{consumer-name}]]
{start-id} {end-id} — intervalle d'identifiants (vous pouvez utiliser «-» et «+»)
{count} — nombre de tentatives de livraison
{consumer-name} — nom du groupe
> XPENDING mystream mygroup - + 10
1) 1) 1526569498055-0
2) "Bob"
3) (integer) 74170458
4) (integer) 1
2) 1) 1526569506935-0
2) "Bob"
3) (integer) 74170458
4) (integer) 1
Maintenant, nous avons les détails de chaque message : identifiant, nom de l'abonné, temps d'inactivité en millisecondes et enfin, nombre de tentatives de livraison. Nous avons deux messages de Bob, et ils sont en attente depuis 74170458 millisecondes, soit environ 20 heures.
Notez que rien ne nous empêche de vérifier quel était le contenu du message en utilisant simplement XRANGE.
> XRANGE mystream 1526569498055-0 1526569498055-0
1) 1) 1526569498055-0
2) 1) "message"
2) "orange"
Nous devons simplement répéter le même identifiant deux fois dans les arguments. Maintenant que nous avons une idée, Alice peut décider qu'après 20 heures d'inactivité, Bob ne se rétablira probablement pas, et il est temps de réclamer ces messages et de reprendre leur traitement à la place de Bob. Pour cela, nous utilisons la commande XCLAIM:
XCLAIM {key} {group} {consumer} {min-idle-time} {ID-1} {ID-2} ... {ID-N}
Avec cette commande, nous pouvons obtenir un message « étranger » qui n'a pas encore été traité en changeant le propriétaire en {consumer}. Cependant, nous pouvons également fournir un temps d'inactivité minimal de {min-idle-time}. Cela permet d'éviter une situation où deux clients essaient de changer simultanément le propriétaire des mêmes messages.
Client 1 : XCLAIM mystream mygroup Alice 3600000 1526569498055-0
Client 2 : XCLAIM mystream mygroup Lora 3600000 1526569498055-0
Le premier client réinitialisera le temps d'inactivité et augmentera le compteur de livraisons. Le deuxième client ne pourra donc pas le demander.
> XCLAIM mystream mygroup Alice 3600000 1526569498055-0
1) 1) 1526569498055-0
2) 1) "message"
2) "orange"
Le message a été récupéré avec succès par Alice, qui peut désormais traiter le message et le confirmer.
De l'exemple ci-dessus, on peut voir qu'une exécution réussie de la demande renvoie le contenu même du message. Cependant, ce n'est pas obligatoire. L'option JUSTID peut être utilisée pour ne retourner que les identifiants des messages. Cela est utile si vous ne vous intéressez pas aux détails du message et souhaitez améliorer les performances du système.
Compteur de livraison
Compteur que vous observez dans la sortie XPENDING — c'est le nombre de livraisons de chaque message. Ce compteur augmente de deux manières : lorsque le message est demandé avec succès via XCLAIM ou lorsque l'appel est utilisé XREADGROUP.
Il est normal que certains messages soient livrés plusieurs fois. L'essentiel est que tous les messages soient finalement traités. Parfois, des problèmes surviennent lors du traitement du message en raison de la corruption du message lui-même ou du traitement du message provoquant une erreur dans le code du gestionnaire. Dans ce cas, il se peut que personne ne soit en mesure de traiter ce message. Comme nous avons un compteur de tentatives de livraison, nous pouvons utiliser ce compteur pour détecter de telles situations. Ainsi, dès que le compteur de livraisons atteint un nombre élevé que vous avez défini, il serait probablement plus judicieux de placer ce message dans un autre flux et d'envoyer une notification à l'administrateur système.
État des flux
Commande XINFO est utilisé pour demander diverses informations sur le flux et ses groupes. Par exemple, la forme de base de la commande ressemble à ceci :
> XINFO STREAM mystream
1) length
2) (integer) 13
3) radix-tree-keys
4) (integer) 1
5) radix-tree-nodes
6) (integer) 2
7) groups
8) (integer) 2
9) first-entry
10) 1) 1524494395530-0
2) 1) "a"
2) "1"
3) "b"
4) "2"
11) last-entry
12) 1) 1526569544280-0
2) 1) "message"
2) "banana"
La commande ci-dessus affiche des informations générales sur le flux spécifié. Voici un exemple un peu plus complexe :
> XINFO GROUPS mystream
1) 1) name
2) "mygroup"
3) consumers
4) (integer) 2
5) pending
6) (integer) 2
2) 1) name
2) "some-other-group"
3) consumers
4) (integer) 1
5) pending
6) (integer) 0
La commande ci-dessus affiche des informations générales sur tous les groupes du flux spécifié.
> XINFO CONSUMERS mystream mygroup
1) 1) name
2) "Alice"
3) pending
4) (integer) 1
5) idle
6) (integer) 9104628
2) 1) name
2) "Bob"
3) pending
4) (integer) 1
5) idle
6) (integer) 83841983
La commande ci-dessus affiche des informations sur tous les abonnés du flux et du groupe spécifiés.
Si vous oubliez la syntaxe de la commande, demandez simplement de l'aide à la commande elle-même :
> XINFO HELP
1) XINFO {subcommand} arg arg ... arg. Les sous-commandes sont :
2) CONSUMERS {key} {groupname} -- Afficher les groupes de consommateurs du groupe {groupname}.
3) GROUPS {key} -- Afficher les groupes de consommateurs du flux.
4) STREAM {key} -- Afficher des informations sur le flux.
5) HELP -- Imprimer cette aide.
Limite de taille du flux
De nombreuses applications ne souhaitent pas accumuler des données dans un flux indéfiniment. Il est souvent utile d'avoir un nombre maximum de messages dans un flux. Dans d'autres cas, il est utile de déplacer tous les messages du flux vers un autre stockage permanent une fois que la taille du flux a atteint un certain seuil. La taille du flux peut être limitée à l'aide du paramètre MAXLEN dans la commande. XADD:
> XADD mystream MAXLEN 2 * value 1
1526654998691-0
> XADD mystream MAXLEN 2 * value 2
1526654999635-0
> XADD mystream MAXLEN 2 * value 3
1526655000369-0
> XLEN mystream
(integer) 2
> XRANGE mystream - +
1) 1) 1526654999635-0
2) 1) "value"
2) "2"
2) 1) 1526655000369-0
2) 1) "value"
2) "3"
Lorsque vous utilisez MAXLEN, les anciennes entrées sont automatiquement supprimées lorsque la longueur spécifiée est atteinte, si bien que le flux conserve une taille constante. Cependant, l'élagage ne se fait pas de manière très efficace en mémoire Redis. Vous pouvez améliorer la situation de la manière suivante :
XADD mystream MAXLEN ~ 1000 * ... champs d'entrée ici ...
L'argument ~ dans l'exemple ci-dessus indique que nous n'avons pas besoin de limiter la longueur du flux à une valeur précise. Dans notre exemple, cela peut être n'importe quel nombre supérieur ou égal à 1000 (par exemple, 1000, 1010 ou 1030). Nous avons simplement précisé que nous souhaitons que notre flux conserve au moins 1000 entrées. Cela rend la gestion de la mémoire beaucoup plus efficace à l'intérieur de Redis.
Il existe également une commande distincte XTRIM, qui effectue la même tâche :
> XTRIM mystream MAXLEN 10
> XTRIM mystream MAXLEN ~ 10
Stockage permanent et réplication
Redis Stream est répliqué de manière asynchrone sur les nœuds esclaves et est sauvegardé dans des fichiers de type AOF (instantané de toutes les données) et RDB (journal de toutes les opérations d'écriture). La réplication de l'état des Consumer Groups est également prise en charge. Ainsi, si un message est en statut « en attente » sur le nœud maître, il aura le même statut sur les nœuds esclaves.
Suppression d'éléments individuels du flux
Il existe une commande spéciale pour supprimer des messages XDEL. La commande reçoit le nom du flux, suivi des identifiants des messages à supprimer :
> XRANGE mystream - + COUNT 2
1) 1) 1526654999635-0
2) 1) "value"
2) "2"
2) 1) 1526655000369-0
2) 1) "value"
2) "3"
> XDEL mystream 1526654999635-0
(integer) 1
> XRANGE mystream - + COUNT 2
1) 1) 1526655000369-0
2) 1) "value"
2) "3"
En utilisant cette commande, il faut prendre en compte que la mémoire ne sera en fait libérée que progressivement.
Flux de longueur nulle
La différence entre les flux et d'autres structures de données Redis est que lorsque d'autres structures de données n'ont plus d'éléments à l'intérieur, la structure elle-même sera supprimée de la mémoire en tant qu'effet secondaire. Par exemple, un ensemble trié sera complètement supprimé lorsque l'appel de ZREM supprimera le dernier élément. En revanche, les flux peuvent rester en mémoire, même s'ils ne contiennent aucun élément.
Conclusion
Redis Stream est idéal pour créer des courtiers de messages, des files d'attente de messages, des journaux unifiés et des systèmes de chat qui conservent l'historique.
Comme l'a dit un jour , les programmes sont des algorithmes plus des structures de données, et Redis vous donne déjà les deux.
Source : habr.com
