Redis Stream — fiabilitatea și scalabilitatea sistemelor tale de mesagerie

Redis Stream — fiabilitatea și scalabilitatea sistemelor tale de mesagerie

Redis Stream — un nou tip abstract de date, introdus în Redis odată cu versiunea 5.0
Conceptual, Redis Stream este un List în care poți adăuga înregistrări. Fiecare înregistrare are un identificator unic. În mod implicit, identificatorul este generat automat și include un marcaj temporal. Așadar, poți solicita intervale de înregistrări în funcție de timp sau poți primi noi date pe măsură ce acestea sosesc în flux, similar cu comanda Unix „tail -f” care citește un fișier de jurnal și așteaptă noi date. Reține că fluxul poate fi ascultat simultan de mai mulți clienți, la fel cum multe procese „tail -f” pot citi simultan un fișier fără a intra în conflict între ele.

Pentru a înțelege toate avantajele noului tip de date, să ne amintim rapid structurile Redis care există de mult timp și care repetă parțial funcționalitatea Redis Stream.

Redis PUB/SUB

Redis Pub/Sub este un sistem de mesagerie simplu, deja integrat în magazinul tău de date key-value. Totuși, această simplitate are un preț:

  • Dacă editorul se întrerupe din diverse motive, acesta pierde toți abonații săi.
  • Editorul trebuie să cunoască adresa exactă a tuturor abonatilor săi.
  • Editorul poate supraîncărca abonații săi dacă datele sunt publicate mai repede decât pot fi procesate.
  • Mesajul este șters din bufferul editorului imediat după publicare, indiferent de câți abonați l-au primit și cât de repede au reușit să proceseze acest mesaj.
  • Toți abonații vor primi mesajul simultan. Abonații trebuie să se coordoneze între ei pentru a stabili ordinea procesării aceluiași mesaj.
  • Nu există un mecanism încorporat de confirmare a procesării cu succes a mesajului de către abonat. Dacă abonatul primește mesajul și cedează în timpul procesării, editorul nu va afla asta.

Redis List

Redis List este o structură de date care suportă comenzi de citire cu blocare. Poți adăuga și citi mesaje din începutul sau sfârșitul listei. Pe baza acestei structuri poți crea un stack sau o coadă decentă pentru sistemul tău distribuit, ceea ce în cele mai multe cazuri este suficient. Principalele diferențe față de Redis Pub/Sub:

  • Mesajul este livrat unui singur client. Primul client blocat prin citire va primi datele primul.
  • Clint trebuie să inițieze singur operațiunea de citire a fiecărui mesaj. List nu știe nimic despre clienți.
  • Mesajele sunt păstrate până când cineva le citește sau le șterge explicit. Dacă ați configurat serverul Redis să scrie datele pe disc, fiabilitatea sistemului crește semnificativ.

Introducere în Stream

Adăugarea unei înregistrări în flux

Comanda XADD adaugă o nouă înregistrare în flux. O înregistrare nu este doar un șir, ci constă dintr-o sau mai multe perechi cheie-valoare. Astfel, fiecare înregistrare este deja structurată și seamănă cu structura unui fișier CSV.

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

În exemplul de mai sus, adăugăm în fluxul numit (cheia) „mystream” două câmpuri: „sensor-id” și „temperature” cu valorile „1234” și „19.8” respectiv. Ca al doilea argument, comanda acceptă un identificator care va fi atribuit înregistrării - acest identificator identifică în mod unic fiecare înregistrare din flux. Totuși, în acest caz, am transmis *, deoarece dorim ca Redis să genereze un nou identificator pentru noi. Fiecare nou identificator va fi mai mare. Astfel, fiecare nouă înregistrare va avea un identificator mai mare în raport cu înregistrările anterioare.

Formatul identificatorului

Identificatorul înregistrării, returnat de comandă XADD, este format din două părți:

{millisecondsTime}-{sequenceNumber}

millisecondsTime — timpul Unix în milisecunde (timpul server Redis). Totuși, dacă timpul curent este același sau mai mic decât timpul ultimei înregistrări, se va folosi marcajul temporal al ultimei înregistrări. Prin urmare, dacă timpul serverului revine în trecut, noul identificator va continua să aibă proprietatea de creștere.

sequenceNumber este folosit pentru înregistrările create în aceeași milisecundă. sequenceNumber va fi crescut cu 1 față de înregistrarea anterioară. Deoarece sequenceNumber are dimensiunea de 64 de biți, în practică nu ar trebui să întâmpinați o limitare a numărului de înregistrări care pot fi generate într-o milisecundă.

Formatul acestor identificatori poate părea ciudat la prima vedere. Un cititor neîncrezător s-ar putea întreba de ce timpul este parte integrantă a identificatorului. Motivul este că fluxurile Redis suportă interogări de interval pe baza identificatorilor. Deoarece identificatorul este legat de momentul în care a fost creat înregistrarea, acest lucru permite interogarea intervalelor de timp. Vom analiza un exemplu concret atunci când vom trece la Studiul comenzii XRANGE.

Dacă, dintr-un anumit motiv, utilizatorului îi este necesar să-și specifice propriul identificator, care, de exemplu, este legat de o anumită sistem extern, atunci putem să-l transmitem comenzii XADD în loc de semnul * așa cum se arată mai jos:

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

Rețineți că, în acest caz, trebuie să urmăriți singuri creșterea identificatorului. În exemplul nostru, identificatorul minim este «0-1», așa că comanda nu va accepta un alt identificator care este egal sau mai mic decât «0-1».

> XADD somestream 0-1 foo bar
(error) ERR ID-ul specificat în XADD este egal sau mai mic decât elementul de vârf al fluxului țintă

Numărul de înregistrări din flux

Se poate obține numărul de înregistrări din flux, folosind pur și simplu comanda XLEN. Pentru exemplul nostru, această comandă va returna următoarea valoare:

> XLEN somestream
(integer) 2

Interogări pe interval — XRANGE și XREVRANGE

Pentru a solicita date pe baza unui interval, trebuie să specificăm două identificatori — începutul și sfârșitul intervalului. Intervalul returnat va include toate elementele, inclusiv limitele. De asemenea, există doi identificatori speciali «-» și «+», care reprezintă, respectiv, cel mai mic (prima înregistrare) și cel mai mare (ultima înregistrare) identificator din flux. Exemplul de mai jos va afișa toate înregistrările fluxului.

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

Fiecare înregistrare returnată reprezintă un array cu două elemente: identificatorul și o listă de perechi cheie-valoare. Am mai menționat că identificatorii înregistrărilor sunt legați de timp. Astfel, putem solicita un interval specific de timp. Totuși, putem specifica în cerere nu un identificator complet, ci doar timpul Unix, lăsând deoparte partea referitoare la sequenceNumberPartea omisă a identificatorului se va considera automat ca fiind zero la începutul intervalului și la maximul posibil la sfârșitul acestuia. Iată un exemplu de cum se poate solicita un interval de două milisecunde.

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

Avem doar un singur înregistrare în acest interval, totuși în seturi de date reale, rezultatul returnat poate fi gigantic. Din acest motiv XRANGE susține opțiunea COUNT. Specificând numărul, putem obține pur și simplu primele N înregistrări. Dacă avem nevoie să obținem următoarele N înregistrări (paginare), putem folosi ultimul identificator obținut, să-i incrementăm sequenceNumber cu unul și să cerem din nou. Să vedem asta în următorul exemplu. Începem să adăugăm 10 elemente folosind XADD (să presupunem că fluxul mystream a fost deja umplut cu 10 elemente). Pentru a începe iterația, obținând 2 elemente la comandă, începem cu intervalul complet, dar cu COUNT egal cu 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"

Pentru a continua iterația cu următoarele două elemente, trebuie să selectăm ultimul identificator obținut, adică 1519073279157-0, și să adăugăm 1 la sequenceNumber.
Identificatorul rezultat, în acest caz 1519073279157-1, poate fi acum utilizat ca nou argument de început al intervalului pentru următoarea apelare 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"

Și așa mai departe. Deoarece complexitatea XRANGE este O(log (N)) pentru căutare, iar apoi O(M) pentru a returna M elemente, fiecare pas de iterație este rapid. Astfel, folosind XRANGE poți itera eficient prin fluxuri.

Comanda XREVRANGE este echivalent cu XRANGE, dar returnează elementele în ordine inversă:

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

Observați că comanda XREVRANGE primește argumentele intervalului start și stop în ordine inversă.

Citirea noilor înregistrări folosind XREAD

Adesea, apare sarcina de a te abona la un flux și de a primi doar mesaje noi. Această concept poate părea asemănătoare cu Redis Pub/Sub sau lista Redis blocantă, dar există diferențe esențiale în modul de utilizare a Redis Stream:

  1. Fiecare mesaj nou este livrat în mod implicit fiecărui abonat. Acest comportament diferă de lista Redis blocantă, unde un nou mesaj va fi citit doar de un singur abonat.
  2. În timp ce în Redis Pub/Sub toate mesajele sunt uitate și niciodată nu sunt salvate, în Stream toate mesajele sunt păstrate pentru o perioadă nedeterminată (dacă clientul nu solicită în mod explicit ștergerea acestora).
  3. Redis Stream permite delimitarea accesului la mesaje în interiorul unui singur flux. Un anumit abonat poate vedea doar istoricul său personal de mesaje.

Puteți să vă abonați la flux și să primiți mesaje noi folosind comanda XREAD. Aceasta este puțin mai complicat decât XRANGE, de aceea vom începe întâi cu exemple mai simple.

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

În exemplul de mai sus este specificată forma neblocat. XREADRețineți că opțiunea COUNT nu este obligatorie. De fapt, singura opțiune obligatorie a comenzii este opțiunea STREAMS, care definește lista fluxurilor împreună cu identificatorul maxim corespunzător. Am scris „STREAMS mystream 0” – dorim să primim toate înregistrările fluxului mystream cu un identificator mai mare decât „0-0”. Așa cum se vede din exemplu, comanda returnează numele fluxului, deoarece ne putem abona la mai multe fluxuri simultan. Am putea scrie, de exemplu, „STREAMS mystream otherstream 0 0”. Rețineți că, după opțiunea STREAMS, trebuie mai întâi să oferim numele tuturor fluxurilor dorite și doar apoi lista identificatorilor.

În această formă simplă, comanda nu face nimic special în comparație cu XRANGE. Cu toate acestea, interesant este că putem transforma ușor XREAD într-o comandă blocantă, specificând argumentul BLOCK:

> XREAD BLOCK 0 STREAMS mystream $

În exemplul de mai sus, este specificată o nouă opțiune BLOCK cu un timp de așteptare de 0 milisecunde (acest lucru înseamnă așteptare eternă). Mai mult, în loc să se transmită un identificator obișnuit pentru fluxul mystream, a fost transmis un identificator special $. Acest identificator special înseamnă că XREAD ar trebui să folosească ca identificator maximul identificatorului din fluxul mystream. Așadar, vom primi doar mesaje noi, începând cu momentul în care am început ascultarea. Într-un anumit sens, acest lucru este similar cu comanda Unix „tail -f”.

Vă rugăm să rețineți că, atunci când folosiți opțiunea BLOCK, nu este necesar să folosim un identificator special $. Putem folosi orice identificator existent în flux. Dacă comanda poate să răspundă imediat la solicitarea noastră, fără blocaj, o va face, altfel va fi blocată.

Blocant XREAD de asemenea, poate asculta simultan mai multe fluxuri, este suficient să le specificați numele. În acest caz, comanda va returna înregistrarea primului flux în care au fost primite datele. Primul abonat, blocat pentru acest flux, va primi datele primul.

Grupuri de consumatori

În anumite sarcini, dorim să delimităm accesul abonaților la mesajele dintr-un singur flux. Un exemplu în care acest lucru poate fi util este o coadă de mesaje cu lucrători care vor primi mesaje diferite din flux pentru a permite scalarea procesării mesajelor.

Dacă presupunem că avem trei abonați C1, C2, C3 și un flux care conține mesajele 1, 2, 3, 4, 5, 6, 7, serviciul mesajelor va avea loc conform diagramei de mai jos:

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

Pentru a obține acest efect, Redis Stream utilizează un concept numit Grup de consumatori. Acest concept este similar cu un pseudo-abonat care primește date din flux, dar este efectiv deservit de mai mulți abonați din cadrul grupului, oferind anumite garanții:

  1. Fiecare mesaj este livrat diferiților abonați din grup.
  2. În cadrul grupului, abonații sunt identificați printr-un nume care este un șir cu respectarea majusculelor. Dacă un anumit abonat iese temporar din grup, el se poate recupera în grup cu propriul nume unic.
  3. Fiecare Grup de consumatori urmează conceputul 'primului mesaj necitit'. Atunci când un abonat solicită mesaje noi, el poate primi doar acele mesaje care nu au fost livrate anterior niciunui abonat din grup.
  4. Există o comandă de confirmare explicită a procesării cu succes a mesajului de către abonat. Până când această comandă nu este apelată, mesajul solicitat va rămâne în stare 'în așteptare'.
  5. În cadrul Grupului de consumatori, fiecare abonat poate solicita istoricul mesajelor care i-au fost livrate, dar care nu au fost încă procesate (în stare 'în așteptare').

Într-un anumit sens, starea grupului poate fi reprezentată astfel:

+----------------------------------------+
| consumer_group_name: mygroup          
| consumer_group_stream: somekey        
| last_delivered_id: 1292309234234-92    
|                                                           
| consumers:                                          
|    "consumer-1" cu mesaje în așteptare  
|       1292309234234-4                          
|       1292309234232-8                          
|    "consumer-42" cu mesaje în așteptare 
|       ... (și așa mai departe)                             
+----------------------------------------+

Acum este timpul să ne familiarizăm cu comenzile de bază pentru Consumer Group, și anume:

  • XGROUP este folosit pentru a crea, distruge și gestiona grupuri
  • XREADGROUP este folosit pentru a citi fluxul prin grup
  • XACK — acestă comandă permite abonatului să marcheze un mesaj ca procesat cu succes

Crearea unui Consumer Group

Să presupunem că fluxul mystream există deja. Atunci comanda de creare a grupului va avea următoarea formă:

> XGROUP CREATE mystream mygroup $
OK

Când creăm grupul, trebuie să trecem un ID de la care grupul va începe să primească mesaje. Dacă dorim să primim doar toate mesajele noi, putem folosi un ID special $ (precum în exemplul nostru de mai sus). Dacă în locul ID-ului special se specifică 0, grupului îi vor fi accesibile toate mesajele din flux.

Acum, când grupul a fost creat, putem începe imediat să citim mesajele cu ajutorul comenzii XREADGROUP. Această comandă este foarte similară cu XREAD și susține opțiunea opțională BLOCK. Totuși, există o opțiune obligatorie GROUP, care trebuie întotdeauna specificată cu două argumente: numele grupului și numele abonatului. Opțiunea COUNT este de asemenea susținută.

Înainte de a citi fluxul, să plasăm câteva mesaje acolo:

> 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

Și acum să încercăm să citim acest flux prin grup:

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

Comanda de mai sus spune următoarele:

"Eu, Alice-abonat, membru al grupului mygroup, vreau să citesc din fluxul mystream un mesaj care nu a fost livrat nimănui anterior."

De fiecare dată când un abonat efectuează o operațiune cu grupul, acesta trebuie să își indice numele, identificându-se în mod unic în cadrul grupului. În comanda de mai sus există un alt detaliu foarte important — identificatorul special „>”. Acest identificator special filtrează mesajele, lăsând doar acele mesaje care nu au fost livrate niciodată până acum.

De asemenea, în cazuri speciale, puteți specifica un identificator real, cum ar fi 0 sau orice alt identificator valid. În acest caz, comanda XREADGROUP îți va returna istoricul mesajelor cu statutul „pending”, care au fost livrate abonatului specificat (Alice), dar care nu au fost încă confirmate prin comanda XACK.

Putem verifica acest comportament, specificând imediat identificatorul 0, fără opțiunea COUNT. Vom vedea doar un singur mesaj în așteptare, adică mesajul cu mărul:

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

Însă, dacă confirmăm mesajul ca fiind procesat cu succes, acesta nu va mai apărea:

> XACK mystream mygroup 1526569495631-0
(integer) 1
> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
   2) (lista sau set gol)

Acum a venit rândul lui Bob să citească ceva:

> 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, membru al grupului mygroup, a cerut nu mai mult de două mesaje. Comanda raportează doar mesajele nedeliveryate din cauza identificatorului special „>”. După cum vezi, mesajul „apple” nu apare, deoarece a fost deja livrat lui Alice, așa că Bob primește „orange” și „strawberry”.

Astfel, Alice, Bob și orice alt abonat al grupului pot citi mesaje diferite din același flux. De asemenea, ei pot citi istoricul lor de mesaje neprocesate sau pot marca mesajele ca fiind procesate.

Există câteva lucruri de care trebuie să ții cont:

  • Odată ce un abonat consideră un mesaj prin comanda XREADGROUP, acest mesaj trece în starea „pending” și este atribuit acestui abonat specific. Alți abonați ai grupului nu vor putea citi acest mesaj.
  • Abonații sunt creați automat la prima mențiune, nu este necesară crearea lor explicită.
  • Folosind XREADGROUP Puteți citi mesaje din mai multe fluxuri diferite simultan, totuși, pentru ca acest lucru să funcționeze, trebuie să creați anterior grupuri cu același nume pentru fiecare flux folosind XGROUP

Recuperare după incident

Abonatul poate să se recupereze după incident și să recitească lista sa de mesaje cu statutul „pending”. Totuși, în lumea reală, abonații pot eșua complet. Ce se întâmplă cu mesajele suspendate ale abonatului, dacă acesta nu a reușit să se recupereze după incident?
Consumer Group oferă o funcție care este utilizată exact pentru astfel de situații — când trebuie să schimbi proprietarul mesajelor.

Primul lucru pe care trebuie să-l facem este să apelăm comanda XPENDING, care afișează toate mesajele grupului cu statutul „pending”. În cea mai simplă formă, comanda este apelată doar cu doi parametri: numele fluxului și numele grupului:

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

Comanda a returnat numărul de mesaje neprocesează pentru întregul grup și pentru fiecare abonat. Avem doar pe Bob cu două mesaje neprocesate, deoarece singurul mesaj solicitat de Alice a fost confirmat folosind XACK.

Putem solicita informații suplimentare, folosind mai mulți parametri:

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

{start-id} {end-id} — interval de identificatori (se pot folosi „-” și „+”)
{count} — numărul de încercări de livrare
{consumer-name} — numele grupului

> 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

Acum avem detaliile pentru fiecare mesaj: identificatorul, numele abonatului, timpul de inactivitate în milisecunde și, în cele din urmă, numărul de încercări de livrare. Avem două mesaje de la Bob, care sunt inactive de 74170458 milisecunde, aproximativ 20 de ore.

Observați că nimeni nu ne împiedică să verificăm care a fost conținutul mesajului, pur și simplu folosind XRANGE.

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

Trebuie doar să repetați același identificator de două ori în argumente. Acum, când avem o idee despre ceea ce s-a întâmplat, Alice poate decide că după 20 de ore de inactivitate, Bob nu se va recupera și este timpul să solcite aceste mesaje și să reia procesarea lor în locul lui Bob. Pentru aceasta, folosim comanda XCLAIM:

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

Folosind această comandă, putem obține un mesaj „străin” care nu a fost încă procesat, prin schimbarea proprietarului la {consumer}. Totuși, putem oferi și un timp minim de inactivitate {min-idle-time}. Acest lucru ajută la evitarea situației în care doi clienți încearcă simultan să schimbe proprietarul pentru aceleași mesaje:

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

Primul client va reseta timpul de inactivitate și va crește numărătorul livrărilor. Așa că al doilea client nu va putea să îl solicite.

> XCLAIM mystream mygroup Alice 3600000 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "message"
      2) "orange"

Mesajul a fost solicitat cu succes de Alice, care acum poate procesa mesajul și să-l confirme.

Din exemplul de mai sus, se poate observa că executarea cu succes a cererii returnează conținutul mesajului propriu-zis. Cu toate acestea, acest lucru nu este obligatoriu. Opțiunea JUSTID poate fi utilizată pentru a returna doar identificatorii mesajului. Acest lucru este util dacă nu sunteți interesat de detaliile mesajului și doriți să îmbunătățiți performanța sistemului.

Numărătorul livrărilor

Numărătorul pe care îl observați în ieșire XPENDING — este numărul livrărilor fiecărui mesaj. Acest numărător crește în două moduri: atunci când mesajul este solicitat cu succes prin XCLAIM sau atunci când se utilizează apelul XREADGROUP.

Este normal ca unele mesaje să fie livrate de mai multe ori. Principalul lucru este că, în final, toate mesajele au fost procesate. Uneori, problemele la procesarea unui mesaj apar din cauza deteriorării mesajului în sine sau procesarea mesajului cauzează o eroare în codul manipulatorului. În acest caz, ar putea fi că acest mesaj nu va putea fi procesat de nimeni. Deoarece avem un numărător de încercări de livrare, putem utiliza acest numărător pentru a detecta astfel de situații. Prin urmare, de îndată ce numărătorul livrărilor atinge un număr mare stabilit de dvs., este probabil mai bine să mutați un astfel de mesaj într-un alt flux și să trimiteți o notificare administratorului de sistem.

Starea fluxurilor

Comanda XINFO este utilizat pentru a solicita diverse informații despre flux și grupurile sale. De exemplu, forma de bază a comenzii arată astfel:

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

Comanda de mai sus afișează informații generale despre fluxul specificat. Acum un exemplu puțin mai complicat:

> XINFO GROUPS mystream
1) 1) nume
   2) "mygroup"
   3) consumatori
   4) (integer) 2
   5) în așteptare
   6) (integer) 2
2) 1) nume
   2) "some-other-group"
   3) consumatori
   4) (integer) 1
   5) în așteptare
   6) (integer) 0

Comanda de mai sus afișează informații generale despre toate grupurile fluxului specificat.

> XINFO CONSUMERS mystream mygroup
1) 1) nume
   2) "Alice"
   3) în așteptare
   4) (integer) 1
   5) inactiv
   6) (integer) 9104628
2) 1) nume
   2) "Bob"
   3) în așteptare
   4) (integer) 1
   5) inactiv
   6) (integer) 83841983

Comanda de mai sus afișează informații despre toți abonații fluxului și grupului specificat.
Dacă uitați sintaxa comenzii, pur și simplu solicitați ajutor de la comanda însăși:

> XINFO HELP
1) XINFO {subcommand} arg arg ... arg. Subcomenzile sunt:
2) CONSUMERS {key} {groupname}  -- Afișează grupurile de consumatori ale grupului {groupname}.
3) GROUPS {key}                 -- Afișează grupurile de consumatori ale fluxului.
4) STREAM {key}                 -- Afișează informații despre flux.
5) HELP                         -- Tipărește acest ajutor.

Limitarea dimensiunii fluxului

Multe aplicații nu doresc să adune date în flux pentru totdeauna. De multe ori este util să avem un număr maxim permis de mesaje în flux. În alte cazuri, este util să mutăm toate mesajele din flux într-un alt stocare permanentă la atingerea dimensiunii specificate a fluxului. Limitarea dimensiunii fluxului se poate face cu parametrul MAXLEN în comandă. 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"

Când se folosește MAXLEN, înregistrările vechi sunt șterse automat la atingerea lungimii specificate, astfel încât fluxul să aibă o dimensiune constantă. Cu toate acestea, tăierea în acest caz nu se face în cel mai eficient mod în memoria Redis. Situația se poate îmbunătăți în felul următor:

XADD mystream MAXLEN ~ 1000 * ... câmpuri de intrare aici ...

Argumentul ~ din exemplul de mai sus înseamnă că nu trebuie să limităm lungimea fluxului la o valoare specifică. În exemplul nostru, aceasta poate fi orice număr mai mare sau egal cu 1000 (de exemplu, 1000, 1010 sau 1030). Pur și simplu am specificat în mod explicit că dorim ca fluxul nostru să păstreze cel puțin 1000 de înregistrări. Acest lucru face gestionarea memoriei mult mai eficientă în interiorul Redis.

Există, de asemenea, o comandă separată XTRIM, care efectuează aceeași acțiune:

> XTRIM mystream MAXLEN 10

> XTRIM mystream MAXLEN ~ 10

Stocare permanentă și replicare

Redis Stream se replică asincron pe nodurile slave și se salvează în fișiere de tip AOF (snapshotul tuturor datelor) și RDB (logul tuturor operațiunilor de scriere). De asemenea, se suportă replicarea stării Grupurilor de Consumatori. Prin urmare, dacă un mesaj se află în statutul 'pending' pe nodul master, atunci pe nodurile slave acest mesaj va avea același statut.

Ștergerea unor elemente din flux

Pentru a șterge mesaje există o comandă specială XDEL. Comanda primește numele fluxului, urmat de identificatorii mesajelor care trebuie șterse:

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

Atunci când se folosește această comandă, trebuie avut în vedere că efectiv memoria nu va fi eliberată imediat.

Fluxuri de lungime zero

Diferența dintre fluxuri și alte structuri de date Redis este că atunci când alte structuri de date nu mai au elemente în interior, ca un efect secundar, însăși structura de date va fi ștearsă din memorie. De exemplu, un set sortat va fi complet șters când apelul ZREM va elimina ultimul element. În schimb, fluxurile pot rămâne în memorie chiar și fără niciun element în interior.

Concluzie

Redis Stream este perfect pentru crearea de brokeri de mesaje, cozi de mesaje, jurnale unificate și sisteme de chat care păstrează istoria.

Așa cum a spus odată Niklaus Wirth, programele sunt algoritmi plus structuri de date, iar Redis îți oferă deja ambele.

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster