Redis Stream — besueshmĂ«ri dhe shkallĂ«zueshmĂ«ri pĂ«r sistemet tuaja tĂ« mesazheve

Redis Stream — besueshmĂ«ri dhe shkallĂ«zueshmĂ«ri pĂ«r sistemet tuaja tĂ« mesazheve

Redis Stream — njĂ« tip i ri abstarct tĂ« tĂ« dhĂ«nave, tĂ« paraqitur nĂ« Redis me daljen e versionit 5.0
Koncepcionalisht, Redis Stream Ă«shtĂ« njĂ« ListĂ«, ku mund tĂ« shtoni regjistrime. Çdo regjistrim ka njĂ« identifikues unik. NĂ« mĂ«nyrĂ« tĂ« paracaktuar, identifikuesi gjenerohet automatikisht dhe pĂ«rmban njĂ« stampĂ« kohe. Prandaj, mund tĂ« kĂ«rkoni range regjistrimesh sipas kohĂ«s ose tĂ« merrni tĂ« dhĂ«na tĂ« reja siç ato mbĂ«rrijnĂ« nĂ« rrjedhĂ«, njĂ«soj si komanda Unix „tail -f“ lexon njĂ« skedar log dhe qĂ«ndron nĂ« pritje pĂ«r tĂ« dhĂ«na tĂ« reja. Vini re qĂ« rrjedha mund tĂ« dĂ«gjohet njĂ«kohĂ«sisht nga disa klientĂ«, po ashtu si shumĂ« procese „tail -f“ mund tĂ« lexojnĂ« njĂ« skedar pa konfliktuar me njĂ«ri-tjetrin.

Për të kuptuar të gjitha përfitimet e këtij tipi të ri të dhënash, le të rikujtojmë për shkurt disa struktura ekzistuese të Redis, të cilat pjesërisht përsërisin funksionalitetin e Redis Stream.

Redis PUB/SUB

Redis Pub/Sub është një sistem i thjeshtë mesazhi, i integruar tashmë në magazinën tuaj key-value. Ndërkohë, për thjeshtësinë duhet të paguani:

  • NĂ«se botuesi ndalon pĂ«r çdo arsye, ai humbet tĂ« gjithĂ« abonentĂ«t e tij
  • Botuesi duhet tĂ« dijĂ« adresĂ«n e saktĂ« tĂ« tĂ« gjithĂ« abonentĂ«ve tĂ« tij
  • Botuesi mund tĂ« ngarkojĂ« punĂ«n e abonentĂ«ve tĂ« tij, nĂ«se tĂ« dhĂ«nat publikohen mĂ« shpejt se sa ato pĂ«rpunohen
  • Mesazhi largohet nga buferi i botuesit menjĂ«herĂ« pas publikimit, pa marrĂ« parasysh se sa abonentĂ«ve iu Ă«shtĂ« dorĂ«zuar ky mesazh dhe sa shpejt ata kanĂ« arritur ta pĂ«rpunojnĂ« kĂ«tĂ« mesazh.
  • TĂ« gjithĂ« abonentĂ«t do tĂ« marrin mesazhin njĂ«kohĂ«sisht. AbonentĂ«t duhet ndonjĂ«herĂ« tĂ« bien dakord mes tyre pĂ«r rendin e pĂ«rpunimit tĂ« tĂ« njĂ«jtit mesazh.
  • Nuk ka njĂ« mekanizĂ«m tĂ« integruar pĂ«r konfirmimin e pĂ«rpunimit tĂ« suksesshĂ«m tĂ« mesazhit nga abonenti. NĂ«se abonenti merr mesazhin dhe bie gjatĂ« pĂ«rpunimit, atĂ«herĂ« botuesi nuk do ta dijĂ« kĂ«tĂ«.

Redis List

Redis List është një strukturë të dhënash që mbështet komandat e leximit me bllokim. Ju mund të shtoni dhe lexoni mesazhe nga fillimi ose fundi i listës. Në bazë të kësaj strukture mund të bëni një STEK të mirë ose një radhë për sistemin tuaj të shpërndarë dhe në shumicën e rasteve do të jetë e mjaftueshme. Dallimet kryesore nga Redis Pub/Sub janë:

  • Mesazhi dorĂ«zohet njĂ« klienti. Klienti i parĂ« qĂ« Ă«shtĂ« bllokuar pĂ«r lexim do tĂ« marrĂ« tĂ« dhĂ«nat i pari.
  • Klienti duhet tĂ« initijrojĂ« vetĂ« operacionin e leximit tĂ« çdo mesazhi. Lista nuk di asgjĂ« pĂ«r klientĂ«t.
  • Mesazhet ruhen deri sa dikush t'i lexojĂ« ose t'i fshijĂ« ato shprehimisht. NĂ«se keni konfiguruar serverin Redis pĂ«r tĂ« ruajtur tĂ« dhĂ«nat nĂ« disk, atĂ«herĂ« qĂ«ndrueshmĂ«ria e sistemit rritet ndjeshĂ«m.

Hyrje në Stream

Shtimi i një shënimi në stream

Ekipa XADD sht-on njĂ« shĂ«nim tĂ« ri nĂ« stream. NjĂ« shĂ«nim — nuk Ă«shtĂ« thjesht njĂ« varg, ai pĂ«rbĂ«het nga njĂ« ose disa çifte kyç-vlerĂ«. Prandaj, çdo shĂ«nim tashmĂ« Ă«shtĂ« i strukturuar dhe ngjan me strukturĂ«n e njĂ« skedari CSV.

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

NĂ« shembullin mĂ« sipĂ«r, po shtojmĂ« nĂ« streamin me emrin (kluçin) "mystream" dy fusha: "sensor-id" dhe "temperature" me vlerat "1234" dhe "19.8" pĂ«rkatĂ«sisht. Si argumentin e dytĂ«, komanda merr njĂ« identifikues qĂ« do t'i jepet shĂ«nimit — ky identifikues identifikon njĂ«herĂ« e pĂ«rgjithmonĂ« çdo shĂ«nim nĂ« stream. MegjithatĂ«, nĂ« kĂ«tĂ« rast e kemi kaluar *, sepse duam qĂ« Redis tĂ« na gjenerojĂ« njĂ« identifikues tĂ« ri. Çdo identifikues i ri do tĂ« rritet. Prandaj, çdo shĂ«nim i ri do tĂ« ketĂ« njĂ« identifikues mĂ« tĂ« madh nĂ« krahasim me shĂ«nimet e mĂ«parshme.

Formati i identifikuesit

Identifikuesi i shënimit, i kthyer nga komanda XADD, përbëhet nga dy pjesë:

{millisecondsTime}-{sequenceNumber}

millisecondsTime — Koha Unix nĂ« milisekonda (koha serverĂ« Redis). MegjithatĂ«, nĂ«se koha aktuale del tĂ« jetĂ« e njĂ«jtĂ« ose mĂ« e vogĂ«l se koha e shĂ«nimit tĂ« mĂ«parshĂ«m, atĂ«herĂ« pĂ«rdoret vula e kohĂ«s sĂ« shĂ«nimit tĂ« mĂ«parshĂ«m. Prandaj, nĂ«se koha e serverit kthehet nĂ« tĂ« kaluarĂ«n, identifikuesi i ri do tĂ« ruajĂ« akoma pronĂ«n e rritjes.

sequenceNumber përdoret për shënime të krijuara në të njëjtën milisekondë. sequenceNumber do të rritet me 1 në krahasim me shënimin e mëparshëm. Duke qenë se sequenceNumber ka madhësi 64 bit, në praktikë nuk duhet të hasni kufizime në numrin e shënimeve që mund të gjenerohen në një milisekondë.

Formati i këtyre identifikatorëve në shikim të parë mund të duket i çuditshëm. Një lexues i dyshimtë mund të pyesë se pse koha është pjesë e identifikatorit. Arsyetimi është se rrjedhat Redis mbështetin kërkesat për intervale sipas identifikatorëve. Duke qenë se identifikatori është i lidhur me kohën e krijimit të regjistrimit, kjo jep mundësinë për të kërkuar intervale kohore. Ne do të shqyrtojmë një shembull konkret kur të kalojmë në studimin e komandës. XRANGE.

Nëse për ndonjë arsye përdoruesi ndjen nevojën të tregojë identifikatorin e tij të vet, i cili, për shembull, lidhet me ndonjë sistem të jashtëm, atëherë ne mund ta kalojmë atë në komandë XADD në vend të shenjës * siç tregohet më poshtë:

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

Kujdes, në këtë rast ju duhet të kujdeseni vetë për rritjen e identifikatorit. Në shembullin tonë, identifikatori minimal është "0-1", kështu që komanda nuk do të pranojë një identifikator tjetër që është i barabartë ose më i vogël se "0-1".

> XADD somestream 0-1 foo bar
(error) ERR Identifikatori i specifikuar në XADD është i barabartë ose më i vogël se artikulli i sipërm i rrjedhës të targetuar

Numri i regjistrimeve në rrjedhë

Mund të merrni numrin e regjistrimeve në rrjedhë duke përdorur thjesht komandën XLEN. Për shembullin tonë, kjo komandë do të kthejë vlerën e mëposhtme:

> XLEN somestream
(integer) 2

KĂ«rkesat pĂ«r intervale — XRANGE dhe XREVRANGE

PĂ«r tĂ« kĂ«rkuar tĂ« dhĂ«na pĂ«r interval, na nevojitet tĂ« tregojmĂ« dy identifikatorĂ« — fillimin dhe fundin e intervalit. Rregulli i kthyer do tĂ« pĂ«rfshijĂ« tĂ« gjithĂ« elementĂ«t, duke pĂ«rfshirĂ« kufijtĂ«. Gjithashtu ekzistojnĂ« dy identifikatorĂ« tĂ« veçantĂ« "-" dhe "+", qĂ« pĂ«rkatĂ«sisht tregojnĂ« identifikatorin mĂ« tĂ« vogĂ«l (regjistri i parĂ«) dhe mĂ« tĂ« madhin (regjistri i fundit) nĂ« rrjedhĂ«. Shembulli mĂ« poshtĂ« do tĂ« nxjerrĂ« tĂ« gjitha regjistrimet e rrjedhĂ«s.

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

Çdo regjistrim i kthyer Ă«shtĂ« njĂ« array me dy elementĂ«: identifikatori dhe lista e çiftĂ«ve çelĂ«s-vlerĂ«. Ne tashmĂ« kemi folur se identifikatorĂ«t e regjistrimeve kanĂ« lidhje me kohĂ«n. Prandaj, mund tĂ« kĂ«rkojmĂ« intervalin e njĂ« periudhe tĂ« caktuar kohore. SidoqoftĂ«, mund tĂ« tregojmĂ« nĂ« kĂ«rkesĂ« njĂ« identifikator tĂ« plotĂ«, por vetĂ«m njĂ« kohĂ« Unix, duke lĂ«nĂ« jashtĂ« pjesĂ«n qĂ« i pĂ«rket. sequenceNumberPjesa e hequr e identifikatorit automatikisht do tĂ« barazohet me zero nĂ« fillim tĂ« intervalit dhe me vlerĂ«n maksimale tĂ« mundshme nĂ« fund tĂ« intervalit. MĂ« poshtĂ« Ă«shtĂ« njĂ« shembull se si mund tĂ« kĂ«rkoni njĂ« interval prej dy milisekondash.

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

Kemi vetëm një regjistrim në këtë interval, megjithatë në grupe të mëdha të dhënash, rezultati i kthyer mund të jetë shumë i madh. Për këtë arsye XRANGE mbështet opsionin COUNT. Duke specifikuar numrin, ne mund të marrim thjesht regjistrimet e para N. Nëse na nevojitet të marrim regjistrimet N pasuese (paginim), mund të përdorim identifikatorin e fundit të marrë, ta rrisim sequenceNumber me një dhe të kërkojmë përsëri. Le të shohim këtë në shembullin e mëposhtëm. Ne fillojmë të shtojmë 10 elemente duke përdorur XADD (supozoni se stream-i mystream ka qenë tashmë i mbushur me 10 elemente). Për të filluar iterimin, duke marrë 2 elemente në çdo komandë, fillojmë me intervalin e plotë, por me COUNT të barabartë me 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"

Për të vazhduar me iterimin e dy elementeve të ardhshëm, na nevojitet të zgjedhim identifikatorin e fundit të marrë, dmth 1519073279157-0, dhe të shtojmë 1 te sequenceNumber.
identifikatori i rezultuar, në këtë rast 1519073279157-1, tani mund të përdoret si një argument të ri fillimi për thirrjen tjetër 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"

Dhe kështu me radhë. Duke qenë se kompleksiteti XRANGE është O(log (N)) për kërkimin dhe pastaj O(M) për kthimin e M elementeve, çdo hap i iterimit është i shpejtë. Kështu që, me XRANGE mund të iteroni efektivisht stream-et.

Ekipa XREVRANGE është ekuivalente XRANGE, por kthen elementet në rendin e anasjelltë:

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

Vini re se komanda XREVRANGE merr argumentet e intervalit start dhe stop në rendin e anasjelltë.

Leximi i regjistrimeve të reja me XREAD

Shpesh lind nevoja për të u regjistruar në stream dhe për të marrë vetëm mesazhe të reja. Kjo koncept mund të duket e ngjashme me Redis Pub/Sub ose listën bllokuese Redis, por ka dallime thelbësore në mënyrën se si përdoret Redis Stream:

  1. Çdo mesazh i ri dĂ«rgohet automatikisht çdo abonenti. Ky sjellje ndryshon nga lista bllokuese Redis, ku njĂ« mesazh i ri do tĂ« lexohet vetĂ«m nga njĂ« abonent i vetĂ«m.
  2. Ndërsa në Redis Pub/Sub të gjitha mesazhet harrohen dhe kurrë nuk ruhen, në Stream të gjitha mesazhet ruajnë për një periudhë të paqartë (nëse klienti nuk kërkon qartë fshirjen).
  3. Redis Stream lejon ndarjen e aksesit në mesazhet brenda një stream. Një abonent specifik mund të shohë vetëm historinë e tij personale të mesazheve.

Ju mund të abonoheni në stream dhe të merrni mesazhe të reja duke përdorur komandën XREAD. Kjo është pak më e komplikuar se XRANGE, prandaj fillimisht do të fillojmë me disa shembuj më të thjeshtë.

> 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Ă« shembullin e mĂ«sipĂ«rm Ă«shtĂ« caktuara forma jobllokuese XREAD. Shtojini vĂ«mendjen, opsioni COUNT nuk Ă«shtĂ« i detyrueshĂ«m. NĂ« fakt, opsioni i vetĂ«m i detyrueshĂ«m i komandĂ«s Ă«shtĂ« opsioni STREAMS, i cili pĂ«rcakton listĂ«n e stream-eve sĂ« bashku me identifikuesin maksimal pĂ«rkatĂ«s. Ne shkruam "STREAMS mystream 0" — duam tĂ« marrim tĂ« gjitha shĂ«nimet e stream-it mystream me identifikues mĂ« tĂ« madh se "0-0". Siç shihet nga shembulli, komanda kthen emrin e stream-it, sepse ne mund tĂ« abonoheni nĂ« disa stream-e njĂ«herĂ«sh. Ne mund tĂ« shkruajmĂ«, pĂ«r shembull, "STREAMS mystream otherstream 0 0". Kjo Ă«shtĂ« njĂ«ri nga tĂ« dhĂ«nat, opsioni STREAMS duhet t'i jepni sĂ« pari emrat e tĂ« gjithĂ« stream-eve tĂ« nevojshme dhe vetĂ«m pastaj listĂ«n e identifikuesve.

Në këtë formë të thjeshtë komanda nuk bën asgjë të veçantë krahasuar me XRANGE. Megjithatë, interesante është se ne mund ta kthejmë lehtësisht XREAD në një komandë bllokuese, duke specifikuar argumentin BLOCK:

> XREAD BLOCK 0 STREAMS mystream $

Në shembullin e mësipërm, është specifikuar një opsion i ri BLOCK me një kohëprerje prej 0 milisekondash (kjo do të thotë pritje të pakufizuar). Për më tepër, në vend që të kalojmë identifikuesin e zakonshëm për stream-in mystream, u kalua një identifikues special $. Ky identifikues special do të thotë që XREAD duhet të përdorë si identifikues identifikuesin maksimal në stream-in mystream. Pra, do të marrim vetëm mesazhe të reja, duke filluar që nga momenti kur filluam të dëgjonim. Në njëfarë mënyre, kjo është si komanda Unix "tail -f".

Vini re përshtypje se kur përdorim opsionin BLOCK, nuk është e nevojshme të përdorim një identifikues të veçantë $. Ne mund të përdorim çdo identifikues ekzistues në rrjedhë. Nëse ekipi është në gjendje të shërbejë kërkesën tonë menjëherë, pa bllokim, ata do ta bëjnë këtë, ndryshe do të bllokohen.

Bllokues XREAD mund gjithashtu të dëgjojë disa rrjedha në të njëjtën kohë, thjesht duhet të përcaktojë emrat e tyre. Në këtë rast, ekipi do të kthejë regjistrimin e rrjedhës së parë ku të dhënat hynë. Abonent i parë, i bllokuar për këtë rrjedhë, do të marrë të dhënat i pari.

Grupet e Konsumatorëve

NĂ« disa detyra, ne duam tĂ« kufizojmĂ« qasjen e abonentĂ«ve nĂ« mesazhe brenda njĂ« rrjedhe. NjĂ« shembull se kur kjo mund tĂ« jetĂ« e dobishme — njĂ« radhĂ« mesazhesh me punonjĂ«s qĂ« do tĂ« marrin mesazhe tĂ« ndryshme nga rrjedha, duke lejuar shkallĂ«zimin e pĂ«rpunimit tĂ« mesazheve.

Nëse ne e imagjinojmë se kemi tre abonentë C1, C2, C3 dhe një rrjedhë që përmban mesazhet 1, 2, 3, 4, 5, 6, 7, atëherë shërbimi i mesazheve do të ndodhi si në diagramin e mëposhtëm:

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

Për të arritur këtë efekt, Redis Stream përdor konceptin e quajtur Grupi i Konsumatorëve. Ky koncept është si një pseudo-abonent që merr të dhëna nga rrjedha, por në fakt shërbehet nga disa abonentë brenda grupit, duke ofruar disa garanci:

  1. Çdo mesazh dĂ«rgohet abonentĂ«ve tĂ« ndryshĂ«m brenda grupit.
  2. Brenda grupit, abonentët identifikohen me emrin e tyre, i cili përbën një varg që merr parasysh regjistrin. Nëse ndonjë abonent del përkohësisht nga grupi, ai mund të rikthehet në grup me emrin e tij unik.
  3. Çdo Grup KonsumatorĂ«sh i nĂ«nshtrohet konceptit "mesazhi i parĂ« i pa lexuar". Kur njĂ« abonent kĂ«rkon mesazhe tĂ« reja, ai mund tĂ« marrĂ« vetĂ«m ato mesazhe qĂ« kurrĂ« mĂ« parĂ« nuk janĂ« dĂ«rguar asnjĂ« abonenti brenda grupit.
  4. Ekziston një komandë për konfirmimin e qartë të përpunimit të suksesshëm të mesazhit nga abonenti. Deri sa ajo komandë të thirret, mesazhi i kërkuar do të mbetet në statusin "ndërprerë".
  5. Brenda Grupit të Konsumatorëve, çdo abonent mund të kërkojë historinë e mesazheve që i janë dërguar atij, por ende nuk janë përpunuar (në statusin "ndërprerë")

Në njëfarë mënyre, gjendja e grupit mund të përshkruhet kështu:

+----------------------------------------+
| consumer_group_name: mygroup          
| consumer_group_stream: somekey        
| last_delivered_id: 1292309234234-92    
|                                                           
| konsumatorë:                                          
|    "consumer-1" me mesazhe në pritje  
|       1292309234234-4                          
|       1292309234232-8                          
|    "consumer-42" me mesazhe në pritje 
|       ... (dhe kështu me radhë)                             
+----------------------------------------+

Tani ka ardhur koha për t'u njohur me komandat kryesore për Grupin e Konsumatorëve, konkretisht:

  • XGROUP pĂ«rdoret pĂ«r tĂ« krijuar, shkatĂ«rruar dhe menaxhuar grupe
  • XREADGROUP pĂ«rdoret pĂ«r tĂ« lexuar rrjedhĂ«n pĂ«rmes grupit
  • XACK — kjo komandĂ« lejon qĂ« abonenti tĂ« shĂ«nojĂ« mesazhin si tĂ« pĂ«rpunuar me sukses

Krijimi i Grupit të Konsumatorëve

Supozoni se rrjedha mystream tashmë ekziston. Atëherë komandë për krijimin e grupit do të jetë:

> KREO XGROUP mystream mygroup $
OK

Kur krijojmĂ« grupin, duhet tĂ« kalojmĂ« identifikuesin nga i cili grupi do tĂ« fillojĂ« tĂ« marrĂ« mesazhe. NĂ«se duam tĂ« marrĂ«mi thjesht tĂ« gjitha mesazhet e reja, mund tĂ« pĂ«rdorim identifikuesin special $ (si nĂ« shembullin tonĂ« mĂ« sipĂ«r). NĂ«se vendosim zero nĂ« vend tĂ« identifikuesit tĂ« veçantĂ«, grupit do t’i jenĂ« tĂ« aksesueshme tĂ« gjitha mesazhet e rrjedhĂ«s.

Tani, kur grupi është krijuar, mund të fillojmë menjëherë të lexojmë mesazhet duke përdorur komandën XREADGROUP. Kjo komandë është shumë e ngjashme me XREAD dhe mbështet një opsion të jashtëzakonshëm BLOCK. Megjithatë, ka një opsion të detyrueshëm GROUP, i cili duhet të shpallet gjithmonë me dy argumente: emrin e grupit dhe emrin e abonentit. Opsioni COUNT mbështetet gjithashtu.

Para se të lexojmë rrjedhën, le të vendosim disa mesazhe aty:

> 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

Tani le të përpiqemi të lexojmë këtë rrjedhë përmes grupit:

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

Komanda e lartpërmendur thotë në mënyrë të thjeshtë:

«Unë, Alice-abonenti, anëtar i grupit mygroup, dua të lexoj nga rrjedha mystream një mesazh që nuk ka qenë kurrë i dorëzuar askujt më parë.»

Çdo herĂ« qĂ« njĂ« abonent kryen njĂ« operacion me grupin, ai duhet tĂ« japĂ« emrin e tij, duke u identifikuar qartĂ« brenda grupit. NĂ« komandĂ«n e mĂ«sipĂ«rme ka njĂ« detaj shumĂ« tĂ« rĂ«ndĂ«sishĂ«m — identifikuesin e veçantĂ« „>“. Ky identifikues i veçantĂ« filtroi mesazhet, duke lĂ«nĂ« vetĂ«m ato qĂ« nuk janĂ« dorĂ«zuar asnjĂ«herĂ« mĂ« parĂ«.

Po ashtu, nĂ« raste tĂ« veçanta, ju mund tĂ« specifikoni njĂ« identifikues real, si 0 ose ndonjĂ« identifikues tjetĂ«r tĂ« vlefshĂ«m. NĂ« kĂ«tĂ« rast, komanda XREADGROUP do t'ju kthejĂ« historinĂ« e mesazheve me statusin „pending“, tĂ« cilat janĂ« dorĂ«zuar abonentit tĂ« specified (Alice), por ende nuk janĂ« konfirmuar me anĂ« tĂ« komandĂ«s XACK.

Ne mund ta verifikojmë këtë sjellje duke specifikuar menjëherë identifikuesin 0, pa opsionin NUMRI. Ne do të shohim vetëm një mesazh të vetëm në pritje, pra mesazhin me mollë:

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

Megjithatë, nëse ne e konfirmojmë mesazhin si të përpunuar me sukses, atëherë ai nuk do të shfaqet më:

> XACK mystream mygroup 1526569495631-0
(integer) 1
> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
   2) (lista ose grup i zbrazët)

Tani është radha e Bobit të lexojë diçka:

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

Bobi, anĂ«tar i grupit mygroup, kĂ«rkoi jo mĂ« shumĂ« se dy mesazhe. Komanda raporton vetĂ«m pĂ«r mesazhet e pa dorĂ«zuara pĂ«r shkak tĂ« identifikuesit tĂ« veçantĂ« „>“. Siç e shihni, mesazhi „apple“ nuk shfaqet, pasi Ă«shtĂ« dorĂ«zuar tashmĂ« Alice-s, kĂ«shtu qĂ« Bobi merr „orange“ dhe „strawberry“.

Pra, Alice, Bob dhe çdo abonent tjetër i grupit mund të lexojnë mesazhe të ndryshme nga e njëjta stream. Ata gjithashtu mund të lexojnë historinë e tyre të mesazheve të papërpunuara ose të etiketojnë mesazhet si të përpunuara.

Ka disa gjëra që duhet të mbani mend:

  • Sapo njĂ« abonent e lexon mesazhin me anĂ« tĂ« komandĂ«s XREADGROUP, ky mesazh kalon nĂ« gjendjen „pending“ dhe i caktoset kĂ«tij abonenti tĂ« veçantĂ«. AbonentĂ«t e tjerĂ« tĂ« grupit nuk do tĂ« jenĂ« nĂ« gjendje ta lexojnĂ« kĂ«tĂ« mesazh.
  • AbonentĂ«t krijohen automatikisht me pĂ«rmendjen e parĂ«, nuk ka nevojĂ« pĂ«r krijimin e tyre tĂ« qartĂ«.
  • Me ndihmĂ«n e XREADGROUP Ju mund tĂ« lexoni mesazhe nga disa rrjedha tĂ« ndryshme njĂ«herĂ«sh, megjithatĂ«, pĂ«r ta bĂ«rĂ« kĂ«tĂ«, duhet tĂ« krijoni paraprakisht grupe me tĂ« njĂ«jtin emĂ«r pĂ«r secilĂ«n rrjedhĂ« duke pĂ«rdorur XGROUP

Rindërtimi pas dështimit

Abonenti mund tĂ« rimarrĂ« pas njĂ« dĂ«shtimi dhe tĂ« rishikojĂ« listĂ«n e tij tĂ« mesazheve me statusin «pending». MegjithatĂ«, nĂ« botĂ«n reale abonentĂ«t mund tĂ« dĂ«shtojnĂ« pĂ«rfundimisht. ÇfarĂ« ndodh me mesazhet e ngrira tĂ« abonentit nĂ«se ai nuk arrin tĂ« rimarrĂ« pas njĂ« dĂ«shtimi?
Grupi i KonsumatorĂ«ve ofron njĂ« funksion qĂ« pĂ«rdoret pikĂ«risht pĂ«r kĂ«to raste — kur Ă«shtĂ« e nevojshme tĂ« ndryshohet pronari i mesazheve.

Hapi i parë është të thërrisni komandën XPENDING, e cila shfaq të gjitha mesazhet e grupit me statusin «pending». Në formën e saj më të thjeshtë, komanda thirret vetëm me dy argumente: emri i rrjedhës dhe emri i grupit:

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

Komanda nxori numrin e mesazheve të papërpunuara për të gjithë grupin dhe për secilin abonent. Ne kemi vetëm Bobin me dy mesazhe të papërpunuara, pasi mesazhi i vetëm i kërkuar nga Alisa u konfirmua përmes XACK.

Mund të kërkojmë informacion shtesë duke përdorur më shumë argumente:

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

{start-id} {end-id} — diapazoni i identifikatorĂ«ve (mund tĂ« pĂ«rdoren «-» dhe «+»)
{count} — numri i pĂ«rpjekjeve tĂ« dorĂ«zimit
{consumer-name} — emri i grupit

> 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

Tani kemi detajet për secilin mesazh: identifikatori, emri i abonentit, koha e papasur në milisekonda dhe, në fund, numri i përpjekjeve të dorëzimit. Ne kemi dy mesazhe nga Bobi, dhe ato janë në pritje për 74170458 milisekonda, rreth 20 orë.

Vini re se askush nuk na pengon të kontrollojmë çfarë kishte përmbajtja e mesazhit, thjesht duke përdorur XRANGE.

> XRANGE mystream 1526569498055-0 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "mesazh"
      2) "portokall"

Thjesht duhet të përsërisim të njëjtin identifikator dy herë në argumente. Tani që kemi një ide, Alisa mund të vendosë se pas 20 orësh papunësie Bobi, ndoshta, nuk do të rimarrë dhe është koha të kërkojmë këto mesazhe dhe të vazhdojmë përpunimin e tyre në vend të Bobit. Për këtë, ne përdorim komandën XCLAIM:

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

Me ndihmën e këtij komandimi, ne mund të marrim një mesazh "tjetër" që ende nuk është përpunuar, duke ndryshuar pronarin në {consumer}. Megjithatë, mund të ofrojmë gjithashtu një kohë minimale pushimi {min-idle-time}. Kjo ndihmon në shmangien e situatave kur dy klientë përpiqen të ndryshojnë pronarin e mesazheve të njëjta në të njëjtën kohë:

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

Klienti i parë do të përdorë kohën e pushimit dhe do ta rrisë numrin e dorëzimeve. Kështu që klienti i dytë nuk do të mund ta kërkojë atë.

> XCLAIM mystream mygroup Alice 3600000 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "mesazhi"
      2) "portokalli"

Mesazhi u mor me sukses nga Alice, e cila tani mund ta përpunojë mesazhin dhe ta konfirmojë atë.

Nga shembulli i mësipërm, shihet se ekzekutimi i suksesshëm i kërkesës kthen përmbajtjen e vetë mesazhit. Megjithatë, kjo nuk është e domosdoshme. Opzioni JUSTID mund të përdoret për të kthyer vetëm identifikuesit e mesazhit. Kjo është e dobishme nëse nuk jeni të interesuar për detajet e mesazhit dhe dëshironi të rritni performancën e sistemit.

Numri i dorëzimeve

Numri që po vëzhgoni në daljen XPENDING është numri i dorëzimeve për secilin mesazh. Ky numër rritet në dy mënyra: kur mesazhi kërkohet me sukses përmes XCLAIM apo kur përdoret thirrja XREADGROUP.

ËshtĂ« normale qĂ« disa mesazhe tĂ« dorĂ«zohen disa herĂ«. E rĂ«ndĂ«sishme Ă«shtĂ« qĂ« nĂ« fund tĂ« gjitha mesazhet tĂ« pĂ«rpunohen. NdonjĂ«herĂ« gjatĂ« pĂ«rpunimit tĂ« mesazhit ndodhin probleme pĂ«r shkak tĂ« dĂ«mtimit tĂ« vetĂ« mesazhit ose pĂ«rpunimi i mesazhit shkakton njĂ« gabim nĂ« kodin e pĂ«rpunuesit. NĂ« kĂ«tĂ« rast, mund tĂ« pĂ«rfundojĂ« qĂ« ky mesazh nuk do tĂ« jetĂ« nĂ« gjendje tĂ« pĂ«rpunojĂ« askush. Pasi kemi numrin e pĂ«rpjekjeve tĂ« dorĂ«zimit, mund tĂ« pĂ«rdorim kĂ«tĂ« numĂ«r pĂ«r tĂ« zbuluar situata tĂ« tilla. Prandaj, sa herĂ« qĂ« numri i dorĂ«zimeve arrin njĂ« numĂ«r tĂ« madh qĂ« keni caktuar, Ă«shtĂ« ndoshta mĂ« e arsyeshme tĂ« vendosni atĂ« mesazh nĂ« njĂ« rrjedhĂ« tjetĂ«r dhe tĂ« dĂ«rgoni njĂ« njoftim administratorit tĂ« sistemit.

Gjendja e rrjedhave

Ekipa XINFO përdoret për të kërkuar informacion të ndryshëm rreth rrjedhave dhe grupeve të saj. Për shembull, pamja themelore e komandës është kështu:

> 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) "mesazhi"
       2) "banane"

Komanda e mësipërme tregon informacionin e përgjithshëm për rrjedhën e caktuar. Tani një shembull pak më kompleks:

> XINFO GROUPS mystream
1) 1) emri
   2) "mygroup"
   3) konsumatorët
   4) (integer) 2
   5) në pritje
   6) (integer) 2
2) 1) emri
   2) "some-other-group"
   3) konsumatorët
   4) (integer) 1
   5) në pritje
   6) (integer) 0

Komanda e mësipërme tregon informacionin e përgjithshëm për të gjitha grupet e caktuara të rrjedhës.

> XINFO CONSUMERS mystream mygroup
1) 1) emri
   2) "Alice"
   3) në pritje
   4) (integer) 1
   5) i papunë
   6) (integer) 9104628
2) 1) emri
   2) "Bob"
   3) në pritje
   4) (integer) 1
   5) i papunë
   6) (integer) 83841983

Komanda e mësipërme tregon informacionin për të gjithë abonentët e rrjedhës dhe grupit të caktuar.
Nëse e harroni sintaksën e komandës, thjesht kërkoni ndihmë nga vetë komanda:

> XINFO HELP
1) XINFO {nënkomandë} arg arg ... arg. Nënkomandat janë:
2) KONSUMATORËT {çelĂ«si} {emri i grupit}  -- Trego grupet e konsumatorĂ«ve tĂ« grupit {emri i grupit}.
3) GRUPET {çelësi}                 -- Trego grupet e konsumatorëve të rrjedhës.
4) RRJEDHA {çelësi}                 -- Trego informacion rreth rrjedhës.
5) NDIHMA                         -- Printo këtë ndihmë.

Kufizimi i madhësisë së rrjedhës

Shumë aplikacione nuk duan të grumbullojnë të dhëna në rrjedhë përjetësisht. Shpesh është e dobishme të keni një numër maksimal të lejuar mesazhesh në rrjedhë. Në raste të tjera, është e dobishme të zhvendosni të gjitha mesazhet nga rrjedha në një ruajtje të qëndrueshme kur arrini madhësinë e caktuar të rrjedhës. Madhësia e rrjedhës mund të kufizohet me parametrin MAXLEN në komandë. XADD:

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

Kur përdoret MAXLEN, regjistrimet e vjetra fshihen automatikisht kur arrihet gjatësia e caktuar, kështu që rrjedha ka një madhësi konstante. Megjithatë, prishja në këtë rast ndodh në një mënyrë jo të efektshme në memorien e Redis. Situata mund të përmirësohet në këtë mënyrë:

XADD mystream MAXLEN ~ 1000 * ... fushat e hyrjes këtu ...

Argumenti ~ në shembullin e mësipërm do të thotë se nuk është e nevojshme të kufizohet gjatësia e rrjedhës në një vlerë të caktuar. Në shembullin tonë, kjo mund të jetë çdo numër më të madh ose të barabartë me 1000 (p.sh., 1000, 1010 ose 1030). Thjesht treguam qartë se duam që rrjedha jonë të ruajë të paktën 1000 regjistrime. Kjo e bën punën me memorien shumë më efikase brenda Redis.

Ka gjithashtu një komandë të veçantë XTRIM, që kryen të njëjtën funksion:

> XTRIM mystream MAXLEN 10

> XTRIM mystream MAXLEN ~ 10

Ruajtja e përhershme dhe replikimi

Redis Stream рДпlicohet asinkronisht nĂ« nodet slave dhe ruhet nĂ« skedarĂ« tĂ« tipit AOF (snapshot i tĂ« dhĂ«nave) dhe RDB (logu i tĂ« gjitha operacioneve tĂ« shkrimit). Po ashtu, mbĂ«shtetet replikimi i gjendjes sĂ« Grupeve tĂ« KonsumatorĂ«ve. KĂ«shtu, nĂ«se njĂ« mesazh Ă«shtĂ« nĂ« statusin "pending" nĂ« nodĂ«n master, atĂ«herĂ« nĂ« nodet slave, ky mesazh do tĂ« ketĂ« tĂ« njĂ«jtin status.

Heqja e elementeve të veçanta nga rrjedha

Për heqjen e mesazheve ekziston një komandë e veçantë XDEL. Komanda merr emrin e rrjedhës, ndjekur nga identifikuesit e mesazheve që duhet të hiqen:

> XRANGE mystream - + COUNT 2
1) 1) 1526654999635-0
   2) 1) "vlera"
      2) "2"
2) 1) 1526655000369-0
   2) 1) "vlera"
      2) "3"
> XDEL mystream 1526654999635-0
(integer) 1
> XRANGE mystream - + COUNT 2
1) 1) 1526655000369-0
   2) 1) "vlera"
      2) "3"

Kur përdoret kjo komandë, duhet të merret parasysh se në fakt, memoria do të lirohet jo menjëherë.

Rrymat me gjatësi zero

Dallimi midis rrymave dhe strukturave të tjera të të dhënave Redis është se kur strukturat e tjera të të dhënave nuk kanë më elemente brenda tyre, si efekt anësor, vetë struktura e të dhënave do të fshihet nga memoria. Kështu, për shembull, një grup i renditur do të fshihet plotësisht kur thirrja ZREM të heqë elementin e fundit. Në vend të kësaj, rrymave u lejohet të mbeten në memorie, edhe nëse nuk kanë asnjë element brenda.

Përfundim

Redis Stream është ideal për të krijuar brokerë mesazhesh, radhë mesazhesh, regjistra të unifikuar dhe sisteme bisedash që ruajnë histori.

Siç tha një herë Niklaus Wirth, programet përbëhen nga algoritme plus struktura të dhënash, dhe Redis tashmë ju jep të dyja.

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster