Redis Stream — de betrouwbaarheid en schaalbaarheid van uw berichtensystemen

Redis Stream — de betrouwbaarheid en schaalbaarheid van uw berichtensystemen

Redis Stream — een nieuwe abstracte datatype gepresenteerd in Redis met de release van versie 5.0
Conceptueel is Redis Stream een lijst waarin je records kunt toevoegen. Elk record heeft een unieke identificatie. Standaard wordt de identificatie automatisch gegenereerd en bevat het een tijdstempel. Daarom kun je tijdsintervallen van records opvragen of nieuwe gegevens ontvangen zodra ze binnenkomen in de stream, net zoals de Unix-commando «tail -f» een logbestand leest en wacht op nieuwe gegevens. Merk op dat meerdere klanten tegelijkertijd naar de stream kunnen luisteren, net zoals veel «tail -f» processen tegelijkertijd een bestand kunnen lezen zonder elkaar te storen.

Om alle voordelen van deze nieuwe datatype te begrijpen, laten we kort de al lang bestaande structuren van Redis herhalen die gedeeltelijk de functionaliteit van Redis Stream herhalen.

Redis PUB/SUB

Redis Pub/Sub is een eenvoudig berichten systeem, al ingebouwd in je key-value opslag. Maar voor de eenvoud moet je een prijs betalen:

  • Als de uitgever om welke reden dan ook uitvalt, verliest hij al zijn abonnees.
  • De uitgever moet het exacte adres van al zijn abonnees kennen.
  • De uitgever kan zijn abonnees overbelasten als gegevens sneller worden gepubliceerd dan ze kunnen worden verwerkt.
  • Berichten worden uit de buffer van de uitgever verwijderd onmiddellijk na publicatie, ongeacht aan hoeveel abonnees het is geleverd en hoe snel zij dit bericht hebben kunnen verwerken.
  • Alle abonnees ontvangen het bericht tegelijkertijd. Abonnees moeten onderling een manier vinden om de volgorde van verwerking van hetzelfde bericht af te stemmen.
  • Er is geen ingebouwd mechanisme voor bevestiging van succesvolle verwerking van het bericht door de abonnee. Als de abonnee het bericht ontvangt en tijdens de verwerking uitvalt, zal de uitgever hiervan niets vernemen.

Redis Lijst

Redis Lijst is een datastructuur die commando's voor blokkerend lezen ondersteunt. Je kunt berichten toevoegen en lezen vanaf het begin of het einde van de lijst. Op basis van deze structuur kun je een stevige stack of queue voor je gedistribueerde systeem creƫren en in de meeste gevallen zal dit voldoende zijn. De belangrijkste verschillen met Redis Pub/Sub zijn:

  • Het bericht wordt aan ƩƩn klant afgeleverd. De eerste geblokkeerde leeskand kan de gegevens als eerste ontvangen.
  • Clint moet zelf de lezing van elk bericht initiĆ«ren. List weet niets van de klanten.
  • Berichten worden bewaard totdat iemand ze leest of ze expliciet verwijdert. Als je de Redis-server hebt ingesteld om gegevens naar de schijf te schrijven, neemt de betrouwbaarheid van het systeem aanzienlijk toe.

Inleiding tot Streams

Een record toevoegen aan de stream

Opdracht XADD voegt een nieuw record toe aan de stream. Een record is niet gewoon een regel; het bestaat uit een of meer sleutel-waardeparen. Elke opname is dus al gestructureerd en lijkt op de structuur van een CSV-bestand.

> XADD mystream * sensor-id 1234 temperature 19.8
1518951480106-0

In het bovenstaande voorbeeld voegen we twee velden toe aan de stream met de naam (sleutel) 'mystream': 'sensor-id' en 'temperature' met de waarden '1234' en '19.8' respectievelijk. Als tweede argument accepteert het commando de identificatie die aan het record zal worden toegewezen — deze identificatie identificeert elke record in de stream uniek. In dit geval hebben we echter * doorgegeven, omdat we willen dat Redis een nieuwe identifier voor ons genereert. Elke nieuwe identifier zal toenemen. Daarom zal elke nieuwe opname een grotere identifier hebben in verhouding tot de vorige opnames.

Format van de identifier

De identificatie van het record, geretourneerd door het commando XADD, bestaat uit twee delen:

{millisecondsTime}-{sequenceNumber}

millisecondsTime — Unix tijd in milliseconden (tijd de server Redis). Echter, als de huidige tijd gelijk is aan of lager is dan de tijd van het vorige record, wordt de tijdstempel van het vorige record gebruikt. Daarom, als de server tijd terug in de tijd gaat, zal de nieuwe identifier nog steeds zijn eigenschap van toenemen behouden.

sequenceNumber wordt gebruikt voor records die in dezelfde milliseconde zijn gemaakt. sequenceNumber zal met 1 worden verhoogd ten opzichte van het vorige record. Aangezien sequenceNumber een formaat van 64 bits heeft, zou je in de praktijk niet tegen de limiet van het aantal records moeten aanlopen dat binnen ƩƩn milliseconde kan worden gegenereerd.

Het formaat van dergelijke identificatoren kan in eerste instantie vreemd lijken. Een achterdochtige lezer kan zich afvragen waarom tijd een onderdeel van de identifier is. De reden is dat Redis-stromen ondersteuning bieden voor bereikvragen op identificatoren. Aangezien de identifier is gekoppeld aan het tijdstip van aanmaak van het record, biedt dit de mogelijkheid om tijdsbereiken op te vragen. We zullen een specifiek voorbeeld bekijken bij het bestuderen van de commando. XRANGE.

Als de gebruiker om welke reden dan ook zijn eigen identifier moet opgeven, die bijvoorbeeld is gekoppeld aan een extern systeem, kunnen we die doorgeven aan het commando. XADD in plaats van het teken * zoals hieronder weergegeven:

> XADD somestream 0-1 field value
0-1
> XADD somestream 0-2 foo bar
0-2

Let op dat je in dit geval zelf de verhoging van de identifier moet bijhouden. In ons voorbeeld is de minimale identifier '0-1', dus de commando zal geen andere identifier accepteren die gelijk is aan of kleiner is dan '0-1'.

> XADD somestream 0-1 foo bar
(error) ERR De opgegeven ID in XADD is gelijk aan of kleiner dan het bovenste item in de doelstroom

Aantal records in de stroom

Je kunt het aantal records in de stroom eenvoudig verkrijgen met het commando XLEN. Voor ons voorbeeld zal dit commando de volgende waarde retourneren:

> XLEN somestream
(integer) 2

Bereikvragen — XRANGE en XREVRANGE

Om gegevens op te vragen binnen een bereik, moeten we twee identificatoren opgeven — het begin- en het eindpunt van het bereik. Het geretourneerde bereik omvat alle elementen, inclusief de grenzen. Er zijn ook twee speciale identificatoren '-' en '+', die respectievelijk de kleinste (eerste record) en de grootste (laatste record) identifier in de stroom betekenen. Het onderstaande voorbeeld zal alle records uit de stroom weergeven.

> 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"

Elke geretourneerde record bestaat uit een array van twee elementen: een identifier en een lijst van sleutel-waarde paren. We hebben al besproken dat de identificatoren van de records betrekking hebben op tijd. Daarom kunnen we een bereik voor een specifieke tijdsperiode opvragen. We kunnen echter in de aanvraag niet de volledige identifier opgeven, maar alleen de Unix-tijd, waarbij we het deel dat betrekking heeft op... sequenceNumberDe weggelaten identifier wordt automatisch gelijkgesteld aan nul aan het begin van het bereik en aan de maximaal mogelijke waarde aan het eind van het bereik. Hieronder staat een voorbeeld van hoe je een bereik van twee milliseconden kunt aanvragen.

> XRANGE mystream 1518951480106 1518951480107
1) 1) 1518951480106-0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperature"
      4) "19.8"

We hebben maar ƩƩn record in dit bereik, maar in echte datasets kan het geretourneerde resultaat enorm zijn. Om deze reden XRANGE ondersteunt de COUNT-optie. Door het aantal op te geven, kunnen we eenvoudig de eerste N records ophalen. Als we de volgende N records nodig hebben (paginering), kunnen we de laatst ontvangen identifier nemen, deze met sequenceNumber ƩƩn verhogen en opnieuw aanvragen. Laten we dit bekijken in het volgende voorbeeld. We beginnen met het toevoegen van 10 items met behulp van XADD (stel dat de stroom mystream al was gevuld met 10 items). Om te beginnen met itereren en per commando 2 items te ontvangen, starten we met het volledige bereik, maar met COUNT gelijk aan 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"

Om de iteratie met de volgende twee elementen voort te zetten, moeten we de laatst ontvangen identifier kiezen, dat wil zeggen 1519073279157-0, en 1 toevoegen aan sequenceNumber.
de resulterende identifier, in dit geval 1519073279157-1, die nu kan worden gebruikt als nieuw argument voor het begin van het bereik voor de volgende aanroep. 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"

En zo verder. Aangezien de complexiteit XRANGE O(log (N)) is voor het zoeken en vervolgens O(M) voor het teruggeven van M elementen, is elke stap in de iteratie snel. Zo kunnen we met XRANGE efficiƫnt itereren over stromen.

Opdracht XREVRANGE is het equivalent XRANGE, maar retourneert de elementen in omgekeerde volgorde:

> XREVRANGE mystream + - COUNT 1
1) 1) 1519073287312-0
   2) 1) "foo"
      2) "value_10"

Let op dat de opdracht XREVRANGE de arguments start en stop in omgekeerde volgorde aanneemt.

Nieuwe records lezen met XREAD

Er ontstaat vaak de behoefte om je aan te melden voor een stroom en alleen nieuwe berichten te ontvangen. Dit concept lijkt misschien op Redis Pub/Sub of een blokkeren Redis Lijst, maar er zijn fundamentele verschillen in hoe Redis Streams moeten worden gebruikt:

  1. Elk nieuw bericht wordt standaard afgeleverd bij elke abonnee. Dit gedrag verschilt van de blokkering van Redis List, waar een nieuw bericht alleen door ƩƩn specifieke abonnee wordt gelezen.
  2. Terwijl in Redis Pub/Sub alle berichten worden vergeten en nooit worden opgeslagen, worden in Stream alle berichten voor onbepaalde tijd opgeslagen (tenzij de klant expliciet om verwijdering vraagt).
  3. Redis Stream maakt het mogelijk om de toegang tot berichten binnen ƩƩn stroom te scheiden. Een specifieke abonnee kan alleen zijn eigen persoonlijke geschiedenis van berichten zien.

U kunt u abonneneren op de stroom en nieuwe berichten ontvangen met het commando XREAD. Dit is iets ingewikkelder dan XRANGE, daarom beginnen we eerst met eenvoudigere voorbeelden.

> 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"

In het bovenstaande voorbeeld wordt de niet-blokkerende vorm aangegeven. XREADLet op dat de optie COUNT niet verplicht is. In feite is de enige verplichte optie van het commando de optie STREAMS, die een lijst van stromen samen met de bijbehorende maximale identificator opgeeft. We hebben "STREAMS mystream 0" geschreven — we willen alle records van de stroom mystream ontvangen met een identificator groter dan "0-0". Zoals te zien is in het voorbeeld, retourneert het commando de naam van de stroom, omdat we ons op meerdere stromen tegelijk kunnen abonneren. We zouden bijvoorbeeld kunnen schrijven "STREAMS mystream otherstream 0 0". Let op dat we na de optie STREAMS eerst de namen van alle benodigde stromen moeten geven en pas daarna de lijst van identificatoren.

In deze eenvoudige vorm doet het commando niets bijzonders in vergelijking met XRANGE. Het interessante is echter dat we XREAD gemakkelijk kunnen omzetten in een blokkadercommando door de parameter BLOCK op te geven:

> XREAD BLOCK 0 STREAMS mystream $

In het bovenstaande voorbeeld is een nieuwe optie BLOCK met een time-out van 0 milliseconden opgegeven (dit betekent eindeloos wachten). Bovendien is in plaats van een gewone identificator voor de stroom mystream een speciale identificator $ doorgegeven. Deze speciale identificator betekent dat XREAD als identificator de maximale identificator in de stroom mystream moet gebruiken. Dus we zullen alleen nieuwe berichten ontvangen, te beginnen vanaf het moment dat we met luisteren zijn begonnen. In zekere zin lijkt dit op het Unix-commando "tail -f".

Houd er rekening mee dat we bij het gebruik van de BLOCK-optie geen specifieke ID $ hoeven te gebruiken. We kunnen elke bestaande ID in de stream gebruiken. Als het team ons verzoek onmiddellijk kan verwerken zonder te blokkeren, zal het dat doen, anders zal het blokkeren.

Blokkerend XREAD kan ook meerdere streams tegelijk beluisteren, je moet gewoon hun namen opgeven. In dat geval retourneert het team de gegevens van de eerste stream waarin gegevens zijn binnengekomen. De eerste abonnee die voor die stream is geblokkeerd, ontvangt als eerste de gegevens.

Abonnee Groepen

In bepaalde taken willen we de toegang van abonnees tot berichten binnen ƩƩn stream scheiden. Een voorbeeld waarin dit nuttig kan zijn, is een berichtenwachtrij met werkers die verschillende berichten uit de stream ontvangen, zodat de verwerking van berichten kan worden opgeschaald.

Als we ons voorstellen dat we drie abonnees C1, C2, C3 hebben en een stream die berichten 1, 2, 3, 4, 5, 6, 7 bevat, wordt de afhandeling van de berichten zoals in de onderstaande diagram weergegeven:

1 -> C1
2 -> C2
3 -> C3
4 -> C1
5 -> C2
6 -> C3
7 -> C1

Om dit effect te bereiken, maakt Redis Stream gebruik van een concept dat een Abonnee Groep wordt genoemd. Dit concept lijkt op een pseudo-abonnee die gegevens uit de stream ontvangt, maar feitelijk wordt bediend door meerdere abonnees binnen de groep, wat bepaalde garanties biedt:

  1. Elk bericht wordt geleverd aan verschillende abonnees binnen de groep.
  2. Binnen de groep worden abonnees geĆÆdentificeerd aan de hand van een naam, die een case-sensitive string is. Als een abonnee tijdelijk uit de groep valt, kan hij zich opnieuw bij de groep aansluiten met zijn unieke naam.
  3. Elke Abonnee Groep volgt het principe van 'eerste ongelezen bericht'. Wanneer een abonnee nieuwe berichten aanvraagt, kan hij alleen die berichten ontvangen die nog nooit eerder aan een abonnee binnen de groep zijn geleverd.
  4. Er is een commando voor expliciete bevestiging van de succesvolle verwerking van een bericht door een abonnee. Totdat dit commando is aangeroepen, blijft het opgevraagde bericht in de status 'pending'.
  5. Binnen de Abonnee Groep kan elke abonnee het berichtenverleden opvragen dat specifiek aan hem is geleverd, maar nog niet is verwerkt (in de status 'pending').

In zekere zin kan de status van de groep als volgt worden weergegeven:

+----------------------------------------+
| consumer_group_name: mygroup          
| consumer_group_stream: somekey        
| last_delivered_id: 1292309234234-92    
|                                                           
| consumers:                                          
|    "consumer-1" met uitstaande berichten  
|       1292309234234-4                          
|       1292309234232-8                          
|    "consumer-42" met uitstaande berichten 
|       ... (en zo verder)                             
+----------------------------------------+

Het is nu tijd om kennis te maken met de belangrijkste opdrachten voor de Consumer Group, namelijk:

  • XGROUP wordt gebruikt voor het aanmaken, vernietigen en beheren van groepen
  • XREADGROUP wordt gebruikt om een stroom via de groep te lezen
  • XACK — deze opdracht stelt de abonnee in staat om een bericht als succesvol verwerkt te markeren

Consumer Group aanmaken

Stel dat de stroom mystream al bestaat. Dan ziet de opdracht voor het aanmaken van de groep er als volgt uit:

> XGROUP CREATE mystream mygroup $
OK

Bij het aanmaken van de groep moeten we de identificatie doorgeven vanaf waar de groep berichten gaat ontvangen. Als we gewoon alle nieuwe berichten willen ontvangen, kunnen we de speciale identificatie $ gebruiken (zoals in ons voorbeeld hierboven). Als we in plaats van de speciale identificatie 0 aangeven, zijn alle berichten in de stroom beschikbaar voor de groep.

Nu de groep is aangemaakt, kunnen we onmiddellijk beginnen met het lezen van berichten met de opdracht XREADGROUP. Deze opdracht is zeer vergelijkbaar met XREAD en ondersteunt de optionele optie BLOCK. Er is echter een verplichte optie GROUP die altijd moet worden opgegeven met twee argumenten: de groepsnaam en de naam van de abonnee. De optie COUNT wordt ook ondersteund.

Voordat we de stroom lezen, laten we daar enkele berichten aan toevoegen:

> 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

Laten we nu deze stroom proberen te lezen via de groep:

> XREADGROUP GROUP mygroup Alice COUNT 1 STREAMS mystream >
1) 1) "mystream"
   2) 1) 1) 1526569495631-0
         2) 1) "message"
            2) "apple"

De bovenstaande opdracht zegt letterlijk het volgende:

"Ik, Alice-abonnee, lid van de groep mygroup, wil ƩƩn bericht lezen uit de stroom mystream dat nog nooit eerder is afgeleverd aan iemand."

Elke keer dat een abonnee een actie met een groep uitvoert, moet hij zijn naam opgeven, zodat hij zichzelf duidelijk identificeert binnen de groep. In de bovengenoemde opdracht is er nog een zeer belangrijk detail — een speciale identificatie ">". Deze speciale identificatie filtert de berichten, zodat alleen de berichten die nog niet zijn afgeleverd, overblijven.

In bepaalde gevallen kunt u ook een echte identificatie opgeven, zoals 0 of een andere geldige identificatie. In dat geval zal de opdracht XREADGROUP u de historie van berichten met de status "pending" geven, die naar de opgegeven abonnee (Alice) zijn gestuurd, maar nog niet zijn bevestigd met de opdracht XACK.

We kunnen dit gedrag controleren door meteen identificatie 0 op te geven, zonder de optie COUNT. We zullen slechts ƩƩn wachtend bericht zien, namelijk het bericht met de appel:

> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
   2) 1) 1) 1526569495631-0
         2) 1) "message"
            2) "apple"

Als we echter het bericht als succesvol verwerkt bevestigen, verschijnt het niet meer:

> XACK mystream mygroup 1526569495631-0
(integer) 1
> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
   2) (lege lijst of set)

Nu is het de beurt aan Bob om iets te lezen:

> 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) "strawberry"

Bob, een lid van de groep mygroup, vroeg om niet meer dan twee berichten. De opdracht rapporteert alleen over niet-afgeleverde berichten vanwege de speciale identificatie ">". Zoals u kunt zien, verschijnt het bericht "apple" niet, omdat het al aan Alice is afgeleverd, dus ontvangt Bob "orange" en "strawberry".

Zo kunnen Alice, Bob en elke andere abonnee van de groep verschillende berichten uit dezelfde stroom lezen. Ze kunnen ook hun geschiedenis van onbewerkte berichten lezen of berichten als verwerkt markeren.

Er zijn een paar dingen om in gedachten te houden:

  • Zodra een abonnee een bericht met de opdracht XREADGROUP, gaat dit bericht naar de status "pending" en wordt aan deze specifieke abonnee toegewezen. Andere abonnees van de groep zullen dit bericht niet kunnen lezen.
  • Abonnees worden automatisch aangemaakt bij de eerste vermelding, er is geen expliciete creatie vereist.
  • Met XREADGROUP Je kunt berichten uit verschillende streams tegelijkertijd lezen, maar om dit te laten werken, moet je vooraf groepen met dezelfde naam voor elke stream maken met behulp van XGROUP

Herstel na storing

Een abonnee kan zich na een fout herstellen en zijn lijst met berichten met de status 'pending' opnieuw bekijken. Maar in de echte wereld kunnen abonnees definitief falen. Wat gebeurt er met de vastgelopen berichten van de abonnee als hij zich niet kan herstellen na de fout?
Consumer Group biedt een functie die precies voor zulke gevallen bedoeld is — wanneer het nodig is om de eigenaar van berichten te wijzigen.

Als eerste moet je het commando aanroepen XPENDING, dat alle berichten van de groep met de status 'pending' weergeeft. In zijn eenvoudigste vorm wordt het commando alleen met twee argumenten aangeroepen: de naam van de stream en de naam van de groep:

> XPENDING mystream mygroup
1) (integer) 2
2) 1526569498055-0
3) 1526569506935-0
4) 1) 1) "Bob"
      2) "2"

Het commando gaf het aantal onbewerkte berichten voor de hele groep en voor elke abonnee weer. We hebben alleen Bob met twee onbewerkte berichten, omdat het enige bericht dat door Alice werd aangevraagd, is bevestigd met behulp van XACK.

We kunnen aanvullende informatie opvragen door meer argumenten te gebruiken:

XPENDING {key} {groupname} [{start-id} {end-id} {count} [{consumer-name}]]

{start-id} {end-id} — bereik van identificaties (je kunt '-' en '+' gebruiken)
{count} — aantal afleverpogingen
{consumer-name} — naam van de groep

> 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

Nu hebben we details voor elk bericht: identificatie, naam van de abonnee, stilstandtijd in milliseconden en tenslotte het aantal afleverpogingen. We hebben twee berichten van Bob, en ze staan al 74170458 milliseconden stil, ongeveer 20 uur.

Merk op dat niets ons tegenhoudt om te controleren wat de inhoud van het bericht was, gewoon door XRANGE.

> XRANGE mystream 1526569498055-0 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "message"
      2) "orange"

We moeten gewoon dezelfde identificatie twee keer in de argumenten herhalen. Nu, wanneer we een idee hebben, kan Alice besluiten dat, na 20 uur stilstand, Bob waarschijnlijk niet zal herstellen en het tijd is om deze berichten op te vragen en ze in behandeling te nemen in plaats van Bob. Hiervoor gebruiken we het commando XCLAIM:

XCLAIM {key} {group} {consumer} {min-idle-time} {ID-1} {ID-2} ... {ID-N}

Met deze opdracht kunnen we een "vreemd" bericht verkrijgen dat nog niet is verwerkt door de eigenaar te veranderen naar {consumer}. We kunnen echter ook de minimale stilstandstijd {min-idle-time} bieden. Dit helpt situaties te vermijden waarin twee klanten tegelijkertijd proberen de eigenaar van dezelfde berichten te veranderen:

Klant 1: XCLAIM mystream mygroup Alice 3600000 1526569498055-0
Klant 2: XCLAIM mystream mygroup Lora 3600000 1526569498055-0

De eerste klant zal de stilstandtijd resetten en de afleveringscounter verhogen. Dus de tweede klant kan het niet aanvragen.

> XCLAIM mystream mygroup Alice 3600000 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "bericht"
      2) "sinaasappel"

Het bericht is met succes opgeƫist door Alice, die nu het bericht kan verwerken en bevestigen.

Uit het bovenstaande voorbeeld blijkt dat een succesvolle uitvoering van de aanvraag de inhoud van het bericht zelf teruggeeft. Dit is echter niet verplicht. De JUSTID-optie kan worden gebruikt om alleen de identificaties van het bericht terug te geven. Dit is nuttig als je niet geĆÆnteresseerd bent in de details van het bericht en je de prestaties van het systeem wilt verbeteren.

Leveringscounter

De teller die je in de uitvoer ziet XPENDING — is het aantal leveringen van elk bericht. Deze teller verhoogt op twee manieren: wanneer een bericht succesvol is aangevraagd via XCLAIM of wanneer de oproep wordt gebruikt XREADGROUP.

Het is normaal dat sommige berichten meerdere keren worden geleverd. Het belangrijkste is dat uiteindelijk alle berichten zijn verwerkt. Soms kunnen er problemen optreden bij het verwerken van een bericht door dat het bericht zelf beschadigd is of omdat de verwerking van het bericht een fout in de handler-code veroorzaakt. In dat geval kan het zijn dat niemand in staat zal zijn om dit bericht te verwerken. Omdat we een afleveringspoging teller hebben, kunnen we deze teller gebruiken om dergelijke situaties te detecteren. Daarom, zodra de afleveringen de opgegeven hoge waarde bereiken, is het waarschijnlijk verstandiger om dit bericht in een andere stroom te plaatsen en de systeembeheerder te waarschuwen.

Staat van stromen

Opdracht XINFO wordt gebruikt om verschillende informatie over de stroom en zijn groepen op te vragen. Bijvoorbeeld, de basisvorm van het commando ziet er als volgt uit:

> XINFO STREAM mystream
 1) lengte
 2) (integer) 13
 3) radix-tree-keys
 4) (integer) 1
 5) radix-tree-nodes
 6) (integer) 2
 7) groepen
 8) (integer) 2
 9) eerste-entry
10) 1) 1524494395530-0
    2) 1) "a"
       2) "1"
       3) "b"
       4) "2"
11) laatste-entry
12) 1) 1526569544280-0
    2) 1) "bericht"
       2) "banaan"

Het bovenstaande commando toont algemene informatie over de opgegeven stroom. Nu een iets complexer voorbeeld:

> XINFO GROUPS mystream
1) 1) naam
   2) "mygroup"
   3) consumenten
   4) (integer) 2
   5) in behandeling
   6) (integer) 2
2) 1) naam
   2) "some-other-group"
   3) consumenten
   4) (integer) 1
   5) in behandeling
   6) (integer) 0

Het bovenstaande commando toont algemene informatie over alle groepen binnen de opgegeven stroom.

> XINFO CONSUMERS mystream mygroup
1) 1) naam
   2) "Alice"
   3) in behandeling
   4) (integer) 1
   5) inactief
   6) (integer) 9104628
2) 1) naam
   2) "Bob"
   3) in behandeling
   4) (integer) 1
   5) inactief
   6) (integer) 83841983

Het bovenstaande commando toont informatie over alle abonnees van de opgegeven stroom en groep.
Als je de syntaxis van het commando vergeet, vraag dan gewoon om hulp met het commando zelf:

> XINFO HELP
1) XINFO {subcommand} arg arg ... arg. Subcommando's zijn:
2) CONSUMERS {key} {groupname}  -- Toont consumenten groepen van groep {groupname}.
3) GROUPS {key}                 -- Toont de consumer groepen van de stroom.
4) STREAM {key}                 -- Toont informatie over de stroom.
5) HELP                         -- Print deze hulp.

Beperkingen op de grootte van de stroom

Veel applicaties willen niet eindeloos gegevens in een stroom verzamelen. Het is vaak nuttig om een maximaal aantal berichten in de stroom te hebben. In andere gevallen is het nuttig om alle berichten uit de stroom over te brengen naar een andere permanente opslag wanneer de opgegeven grootte van de stroom is bereikt. Je kunt de grootte van de stroom beperken met de MAXLEN parameter in het commando. XADD:

> XADD mystream MAXLEN 2 * waarde 1
1526654998691-0
> XADD mystream MAXLEN 2 * waarde 2
1526654999635-0
> XADD mystream MAXLEN 2 * waarde 3
1526655000369-0
> XLEN mystream
(integer) 2
> XRANGE mystream - +
1) 1) 1526654999635-0
   2) 1) "waarde"
      2) "2"
2) 1) 1526655000369-0
   2) 1) "waarde"
      2) "3"

Bij gebruik van MAXLEN worden oude records automatisch verwijderd wanneer de aangegeven lengte is bereikt, zodat de stroom een constante grootte heeft. Echter, het trimmen gebeurt in dit geval niet op de meest efficiƫnte manier in het geheugen van Redis. De situatie kan als volgt worden verbeterd:

XADD mystream MAXLEN ~ 1000 * ... invoervelden hier ...

De argument ~ in het bovenstaande voorbeeld betekent dat we niet noodzakelijkerwijs de lengte van de stroom aan een specifiek nummer hoeven te beperken. In ons voorbeeld kan dit elk getal zijn dat groter of gelijk is aan 1000 (bijvoorbeeld 1000, 1010 of 1030). We hebben simpelweg expliciet aangegeven dat we willen dat onze stroom minimaal 1000 records opslaat. Dit maakt het werken met het geheugen veel efficiƫnter binnen Redis.

Er is ook een apart commando beschikbaar, XTRIM, dat precies hetzelfde doet:

> XTRIM mystream MAXLEN 10

> XTRIM mystream MAXLEN ~ 10

Permanente opslag en replicatie

Redis Stream wordt asynchroon gerepliceerd naar slave-nodes en opgeslagen in AOF-bestanden (snapshot van alle gegevens) en RDB-bestanden (logboek van alle schrijfoperaties). Ook wordt de replicatie van de status van Consumer Groups ondersteund. Dus als een bericht de status 'pending' heeft op de master-node, zal het op de slave-nodes dezelfde status hebben.

Verwijderen van afzonderlijke elementen uit de stream

Voor het verwijderen van berichten is er een speciale commando XDEL. Het commando ontvangt de naam van de stream gevolgd door de identificatiecodes van de berichten die moeten worden verwijderd:

> 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"

Bij het gebruik van deze commando moet rekening worden gehouden met het feit dat het geheugen feitelijk niet onmiddellijk wordt vrijgegeven.

Streams van nul lengte

Het verschil tussen streams en andere datatypes in Redis is dat, wanneer andere datatypes geen elementen meer bevatten, de structuur als neveneffect uit het geheugen zal worden verwijderd. Een gesorteerde set zal bijvoorbeeld volledig worden verwijderd wanneer de aanroep ZREM het laatste element verwijdert. In tegenstelling tot dat, mogen streams in het geheugen blijven, zelfs als ze geen enkele element erin hebben.

Conclusie

Redis Stream is bij uitstek geschikt voor het creƫren van message brokers, message queues, uniforme logs en chatsystemen die geschiedenis opslaan.

Zoals eenmaal gezegd door Niklaus Wirth, zijn programma's algoritmen plus datastructuren, en Redis biedt je al beide.

Bron: habr.com

Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers šŸ”„ Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers | ProHoster