Redis Stream — qëndrueshmëri dhe shkallëzueshmëri e sistemeve tuaja të mesazheve

Redis Stream — qëndrueshmëri dhe shkallëzueshmëri e sistemeve tuaja të mesazheve

Redis Stream — një tip i ri abstrakt të të dhënave, i prezantuar në Redis me daljen e versionit 5.0
Koncepte, Redis Stream është një Listë në të cilën mund të shtoni regjistrime. Çdo regjistrim ka një identifikues unik. Në mënyrë standarde, identifikuesi gjenerohet automatikisht dhe përfshin një markë kohe. Prandaj, mund të kërkoni intervale regjistrimesh sipas kohës ose të merrni të dhëna të reja ndërsa ato arrijnë në rrjedhë, si komandën Unix «tail -f» që lexon një skedar regjistrimi dhe ndalon për të pritur të dhëna të reja. Keni parasysh se disa klientë mund ta dëgjojnë të njëjtin rrjedhë njëkohësisht, si shumë procese «tail -f» që mund të lexojnë një skedar pa ndikuar njëri-tjetrin.

Për të kuptuar të gjitha përfitimet e këtij tipi të ri të të dhënave, le të përmendim shkurtimisht strukturat e vjetra të Redis, të cilat pjesërisht përsërisin funksionalitetin e Redis Stream.

Redis PUB/SUB

Redis Pub/Sub është një sistem i thjeshtë mesazhesh, tashmë i integruar në magazinën tuaj key-value. Megjithatë, për thjeshtësinë e tij paguani një çmim:

  • Nëse publisher-i për ndonjë arsye dështon, ai humbet të gjithë abonentët e tij
  • Publisher-i duhet të dijë adresën e saktë të të gjitha abonentëve të tij
  • Publisher-i mund të mbingarkojë punën e abonentëve të tij, nëse të dhënat publikohen më shpejt sesa përpunohen
  • Mesazhi hiqet nga buffer-i i publisher-it menjëherë pas publikimit, pavarësisht se sa abonentëve u është dorëzuar dhe sa shpejt ata kanë arritur të përpunojnë këtë mesazh.
  • Të gjithë abonentët do të marrin mesazhin njëkohësisht. Abonentët duhet të arrijnë njëfarë marrëveshjeje për rendin e përpunimit të mesazheve të njëjta.
  • Nuk ka mekanizëm të integruar për konfirmimin e përpunimit të suksesshëm të mesazhit nga abonenti. Nëse një abonent merr mesazhin dhe dështon gjatë përpunimit, publisher-i nuk do ta dijë këtë.

Redis List

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

  • Mesazhi dorëzohet një klienti. Klienti i bllokuar në lexim do të marrë të dhënat i pari.
  • Klienti duhet të initiojë vetë operacionin e leximit të çdo mesazhi. Lista nuk e di asgjë për klientët.
  • Mesazhet ruhet deri sa dikush t'i lexojë ose t'i fshijë në mënyrë të qartë. Nëse keni konfiguruar serverin Redis për ta ruajtur të dhënat në disk, besueshmëria e sistemit do të rritet ndjeshëm.

Hyrje në Stream

Shtimi i një regjistrimi në rrjedhë

Ekipa XADD shton një regjistrim të ri në rrjedhë. Regjistrimi — nuk është thjesht një varg, ai përbëhet nga një ose disa çiftë çelës-vlerë. Kështu, çdo regjistrim është tashmë i strukturuar dhe ngjan me strukturën e një skedari CSV.

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

Në shembullin e mësipërm, ne po shtojmë në rrjedhën me emrin (çelësi) «mystream» dy fusha: «sensor-id» dhe «temperature» me vlerat «1234» dhe «19.8» përkatësisht. Si argumentin e dytë, komanda merr identifikuesin që do t'i jepet regjistrimit — ky identifikues identifikon qartësisht çdo regjistrim në rrjedhë. Megjithatë, në këtë rast ne e kaluam *, sepse duam që Redis të gjenerojë për ne një identifikues të ri. Çdo identifikues i ri do të rritet. Prandaj, çdo regjistrim i ri do të ketë një identifikues më të madh në krahasim me regjistrimet e mëparshme.

Formati i identifikuesit

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

{millisecondsTime}-{sequenceNumber}

millisecondsTime — Koha Unix në milisekonda (koha e server Redis). Megjithatë, nëse koha aktuale rezulton të jetë e njëjtë ose më e vogël se koha e regjistrimit të mëparshëm, atëherë përdoret marka e kohës së regjistrimit të mëparshëm. Prandaj, nëse koha e serverit kthehet në të kaluarën, identifikuesi i ri ende do të ruajë pronën e rritjes.

sequenceNumber përdoret për regjistrimet e krijuara në të njëjtën milisekondë. sequenceNumber do të rritet me 1 në krahasim me regjistrimin e mëparshëm. Duke qenë se sequenceNumber ka madhësinë 64 bit, në praktikë nuk duhet të hasni asnjë kufizim në numrin e regjistrimeve që mund të krijohen në një milisekondë.

Formati i këtyre identifikuesve në shikim të parë mund të duket i çuditshëm. Një lexues që nuk beson mund të pyesë pse koha është pjesë e identifikuesit. Arsyetimi është se rrjedhat e Redis mbështesin kërkesat për intervale sipas identifikuesve. Duke qenë se identifikuesi është i lidhur me kohën e krijimit të regjistrimit, kjo ofron mundësinë për të kërkuar intervale kohe. Ne do të shqyrtojmë një shembull të caktuar kur të kalojmë në studimin e komandës XRANGE.

Nëse për ndonjë arsye përdoruesi duhet të specifikojë identifikuesin e tij të vet, i cili, për shembull, është i lidhur me një sistem të jashtëm, ne mund ta kalojmë atë ekipit. XADD ndërsa është përdorur shenjën * siç tregohet më poshtë:

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

Vini re se në këtë rast, ju duhet të kujdeseni vetë për rritjen e identifikuesit. Në shembullin tonë, identifikuesi minimal është "0-1", prandaj komanda nuk do të pranojë një identifikues tjetër që është i barabartë ose më i vogël se "0-1".

> XADD somestream 0-1 foo bar
(error) ERR ID i specifikuar në XADD është i barabartë ose më i vogël se elementi më i lartë i fluxit të synuar

Numri i shënimeve në flux

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

> XLEN somestream
(integer) 2

Kërkesat sipas gamës — XRANGE dhe XREVRANGE

Për të kërkuar të dhëna sipas gamës, na nevojitet të specifikojmë dy identifikues — fillimi dhe përfundimi i gamës. Gama e kthyer do të përfshijë të gjithë elementët, përfshirë kufijtë. Ka gjithashtu dy identifikues specialë "-" dhe "+", përkatësisht duke treguar identifikuesin më të vogël (shënimi i parë) dhe më të madhin (shënimi i fundit) në flux. Shembulli më poshtë do të nxjerrë të gjithë shënimet e fluxit.

> 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 shënim i kthyer përfaqëson një array me dy elementë: identifikuesi dhe lista e çiftëve çelës-vlerë. Ne tashmë kemi folur se identifikuesit e shënimit kanë lidhje me kohën. Prandaj, mund të kërkojmë një gamë të një periudhe të caktuar kohore. Megjithatë, ne mund të specifikojmë në kërkesë jo identifikuesin e plotë, por vetëm kohën Unix, duke lënë jashtë pjesën që i përket sequenceNumber. Pjesa e lënë jashtë e identifikuesit në mënyrë automatike do të barazohet me zero në fillimin e gamës dhe me vlerën maksimale në fundin e gamës. Më poshtë është një shembull se si mund të kërkojmë një gamë që është e barabartë me dy millisekonda.

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

Ne kemi vetëm një shënim në këtë gamë, megjithatë në grupe të dhënash reale, rezultati i kthyer mund të jetë i madh. Për këtë arsye XRANGE mbështet opsionin COUNT. Duke specifikuar numrin, ne mund të marrim thjesht shënimet e para N. Nëse na nevojitet të marrim shënimet e tjera N (paginim), ne mund të përdorim identifikuesin e fundit të marrë, ta rrisim atë sequenceNumber me një dhe ta 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 (le të supozojmë se fluxi mystream tashmë është plotësuar me 10 elemente). Për të filluar iterimin, duke marrë dy elemente për komandë, ne fillojmë me gamën 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 iterimin me dy elementët e ardhshëm, na nevojitet të zgjedhim identifikuesin e fundit të marrë, pra 1519073279157-0, dhe të shtojmë 1 në sequenceNumber.
Identifikuesin rezultat, në këtë rast 1519073279157-1, tani mund të përdoret si argumenti i ri i fillimit të gamës për thirrjen e ardhshme. 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"

Etj. Duke marrë parasysh kompleksitetin XRANGE është O(log (N)) për kërkimin, e më pas O(M) për kthimin e M elementëve, çdo hap i iterimit është i shpejtë. Kështu, me XRANGE mund të iteroni efikasht fluxet.

Ekipa XREVRANGE është ekuivalente me XRANGE, por kthen elementet në rendin e kundërt:

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

Vini re se komanda XREVRANGE merr argumentet e gamës start dhe stop në rend të kundërt.

Leximi i shënimeve të reja me XREAD

Ndërsa shumë shpesh shfaqet detyra për t'u regjistruar në një flux dhe për të marrë vetëm mesazhet e reja. Kjo koncept mund të duket e ngjashme me Redis Pub/Sub ose listën bllokuese të Redis, por ka dallime thelbësore në mënyrën si përdoret Redis Stream:

  1. Çdo mesazh i ri, sipas parimit, dërgohet te çdo regjistrues. Ky qëndrim është ndryshe nga lista bllokuese e Redis, ku një mesazh i ri do të lexohet vetëm nga një regjistrues të vetëm.
  2. Ndërsa në Redis Pub/Sub të gjitha mesazhet harrohen dhe kurrë nuk ruhet, në Stream të gjitha mesazhet ruhet për një periudhë të pacaktuar (nëse klienti nuk kërkon eksplicitisht fshirjen).
  3. Redis Stream lejon ndarjen e aksesit në mesazhe brenda një fluksi. Një regjistrues i caktuar mund të shohë vetëm historinë e tij personale të mesazheve.

Mund të regjistroheni në flux dhe të merrni mesazhe të reja duke përdorur komandën XREAD. Kjo është pak më e komplikuar se XRANGE, prandaj së pari do të fillojmë me 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"

В примере выше указана неблокирующая форма XREAD. Обратите внимание, что опция COUNT не является обязательной. Фактически единственной обязательной опцией команды является опция STREAMS, которая задает список потоков вместе с соответствующим максимальным идентификатором. Мы написали «STREAMS mystream 0» — мы хотим получать все записи потока mystream с идентификатором больше чем «0-0». Как видно из примера, команда возвращает имя потока, потому что мы можем подписаться на несколько потоков одновременно. Мы могли бы написать, например, «STREAMS mystream otherstream 0 0». Обратите внимание, что после опции STREAMS нам нужно сперва предоставить имена всех нужных потоков и только затем список идентификаторов.

В этой простой форме команда не делает ничего особенного по сравнению с XRANGE. Однако интересно то, что мы можем легко превратить XREAD в блокирующую команду, указав аргумент BLOCK:

> XREAD BLOCK 0 STREAMS mystream $

В приведенном выше примере, указана новая опцию BLOCK с временем ожидания 0 миллисекунд (это означает бесконечное ожидание). Более того, вместо передачи обычного идентификатора для потока mystream, был передан специальный идентификатор $. Этот специальный идентификатор означает, что XREAD должен использовать в качестве идентификатора максимальный идентификатор в потоке mystream. Так что мы будем получать только новые сообщения, начиная с момента, когда мы начали прослушивание. В некотором смысле это похоже на Unix команду «tail -f».

Обратите внимание, что при использовании опции BLOCK нам не обязательно нужно использовать специальный идентификатор $. Мы можем использовать любой существующий в потоке идентификатор. Если команда сможет обслужить наш запрос немедленно, без блокировки, она сделает это, в противном случае она заблокируется.

Блокирующий XREAD также может прослушивать сразу несколько потоков, просто нужно указать их имена. В этом случае команда вернет запись первого потока, в который поступили данные. Первый подписчик, заблокированный для данного потока, будет получать данные первым.

Consumer Groups

В определенных задачах мы хотим разграничить доступ подписчиков к сообщениям внутри одного потока. Примером, когда это может быть полезно — очередь сообщений с воркерами, которые будут получать разные сообщения потока, позволяя масштабировать обработку сообщений.

Если мы представим, что у нас есть три подписчика C1, C2, C3 и поток, который содержит сообщения 1, 2, 3, 4, 5, 6, 7, то обслуживание сообщений будет происходить как на диаграмме ниже:

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

Чтобы получить этот эффект, Redis Stream использует концепцию, называемую Consumer Group. Эта концепция подобна псевдо-подписчику, который получает данные из потока, но фактически обслуживается несколькими подписчиками внутри группы, предоставляя определенные гарантии:

  1. Каждое сообщение доставляется разным подписчикам внутри группы.
  2. В пределах группы подписчики идентифицируются по имени, которое представляет собой строку с учетом регистра. Если какой-то подписчик временно выпадет из группы, то он может восстановиться в группу по собственному уникальному имени.
  3. Каждая Consumer Group следует концепции «первое непрочитанное сообщение». Когда подписчик запрашивает новые сообщения, он может получить только те сообщения, которые никогда ранее не доставлялись ни одному подписчику внутри группы.
  4. Существует команда явного подтверждения успешной обработки сообщения подписчиком. Пока не будет вызвана эта команда, запрошенное сообщение будет оставаться в статусе «pending» .
  5. Внутри Consumer Group каждый подписчик может запрашивать историю сообщений, которые были доставлены именно ему, но ещё не были обработаны (в статусе «pending»)

В некотором смысле, состояние группы может быть преставлено так:

+----------------------------------------+
| consumer_group_name: mygroup          
| consumer_group_stream: somekey        
| last_delivered_id: 1292309234234-92    
|                                                           
| consumers:                                          
|    "consumer-1" with pending messages  
|       1292309234234-4                          
|       1292309234232-8                          
|    "consumer-42" with pending messages 
|       ... (and so forth)                             
+----------------------------------------+

Теперь пришло время познакомиться с основными командами для Consumer Group, а именно:

  • XGROUP используется для создания, уничтожения и управления группами
  • XREADGROUP используется для чтения потока через группу
  • XACK — это команда позволяет подписчику пометить сообщение как успешно обработанное

Создание Consumer Group

Supozoni që rrjedha mystream tashmë ekziston. Atëherë komandat për krijimin e grupit do të duken kështu:

> XGROUP KRIJO mystream mygroup $
OK

Kur krijojmë një grup, ne duhet të kalojmë identifikuesin nga i cili grupi do të fillojë të marrë mesazhe. Nëse dëshirojmë të marrim thjesht të gjitha mesazhet e reja, mund të përdorim identifikuesin e veçantë $ (si në shembullin tonë më sipër). Nëse vendosim në vend të identifikuesit të veçantë 0, grupi do të ketë qasje në të gjitha mesazhet e rrjedhës.

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

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

> 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ë provojmë 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 mësipërme thotë fjalë për fjalë:

"Unë, abonenti Alice, anëtar i grupit mygroup, dua të lexoj nga rrjedha mystream një mesazh, i cili kurrë nuk është dërguar askujt më parë".

Çdo herë që një abonent kryen një operacion me grupin, ai duhet të specifikojë emrin e tij, duke identifikuar njëherë e mirë veten brenda grupit. Në komandën e mësipërme ka një detaj tjetër shumë të rëndësishëm — identifikuesin e veçantë „>”. Ky identifikues special filtroi mesazhet, duke lënë vetëm ata që nuk janë dërguar asnjëherë më parë.

Po ashtu, në raste të veçanta, ju mund të specifikoni një identifikues real, si 0 ose çdo identifikues tjetër të vlefshëm. Në këtë rast, komandat XREADGROUP do t'ju kthejë historinë e mesazheve me status „pending”, të cilat iu dërguan abonentit të specifikuar (Alice), por nuk janë konfirmuar ende me komandën XACK.

Ne mund ta verifikojmë këtë sjellje, duke specifikuar menjëherë identifikuesin 0, pa opsionin COUNT. Ne thjesht do të shohim mesazhin e vetëm në pritje, domethënë 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) (listë ose grup i zbrazët)

Tani është radha e Bob-it për të lexuar 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"

Bob, anëtar i grupit mygroup, kërkoi jo më shumë se dy mesazhe. Komanda raporton vetëm për mesazhet e pa dërguara për shkak të identifikuesit të veçantë „>”. Siç e shihni, mesazhi „apple” nuk shfaqet, pasi ai ka qenë tashmë i dërguar Alice-s, kështu që Bob merr „orange” dhe „strawberry”.

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

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

  • Pasi abonenti të lexojë mesazhin me komandën XREADGROUP, ky mesazh kalon në gjendjen „pending” dhe sigurohet për këtë abonent të veçantë. Abonentët e tjerë të grupit nuk do të mund ta lexojnë këtë mesazh.
  • Abonentët krijohen automatikisht me përmendjen e parë, nuk është e nevojshme të krijohen shprehimisht.
  • Me XREADGROUP mund të lexoni mesazhe nga disa rrjedha të ndryshme njëkohësisht, megjithatë, për të funksionuar, duhet të krijoni paraprakisht grupe me të njëjtin emër për çdo rrjedhë me anë të XGROUP

Rivendosja pas dështimit

Abonenti mund të rikuperohet pas dështimit dhe të rilexojë listën e tij të mesazheve me status „pending”. Megjithatë, në botën reale, abonentët mund të dështojnë përfundimisht. Çfarë ndodh me mesazhet e bllokuara të abonentit, nëse ai nuk mund të rikuperohet pas dështimit?
Grupi i konsumatorëve ofron një funksion që përdoret pikërisht për rastet e tilla — kur është e nevojshme të ndryshohet pronari i mesazheve.

Së pari, është e nevojshme të thërritet komanda XPENDING, e cila tregon të gjitha mesazhet e grupit me status „pending”. Në formën e saj më të thjeshtë, komanda thirret vetëm me dy argumente: emrin e rrjedhës dhe emrin e grupit:

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

Ekipi shfaqi numrin e mesazheve të papërpunuara për të gjithë grupin dhe për secilin anëtar të grupit. Kemi vetëm Bob me dy mesazhe të papërpunuara, pasi mesazhi i vetëm i kërkuar nga Alisa u konfirmua me 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} — intervali i identifikatorëve (mund të përdoren «-» dhe «+»)
{count} — numri i përpjekjeve për dorëzim
{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 çdo mesazh: identifikatori, emri i anëtarit, koha e papërdorur në milisekonda dhe, përfundimisht, numri i përpjekjeve për dorëzim. Kemi dy mesazhe nga Bob, dhe ato kanë qenë të papërdorura për 74170458 milisekonda, rreth 20 orë.

Kujdes, asgjë nuk na pengon të kontrollojmë se çfarë përmbante mesazhi, thjesht duke përdorur XRANGE.

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

Thjesht duhet të përsërisim të njëjtin identifikator dy herë në argumentet. Tani, kur kemi një ide, Alisa mund të vendosë se pas 20 orësh të papërdorur, Bob, ndoshta, nuk do të rikthehet, dhe është koha të kërkojmë këto mesazhe dhe të vazhdojmë përpunimin e tyre në vend të Bobit. Për këtë, përdorim komandën XCLAIM:

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

Me këtë komandë, mund të marrim një mesazh "tjetër" që ende nuk është përpunuar, duke e ndryshuar pronarin në {consumer}. Megjithatë, mund të ofrojmë gjithashtu kohën minimale të papërdorur {min-idle-time}. Kjo ndihmon në shmangien e situatës kur dy klientë përpiqen në të njëjtën kohë të ndryshojnë pronarin e të njëjtëve mesazhe:

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

Klienti i parë do të rikthejë kohën e papërdorur dhe do të rrisë numrin e dorëzimeve. Pra, klienti i dytë nuk do të mund ta kërkojë.

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

Mesazhi u kërkua me sukses nga Alisa, e cila tani mund të përpunojë mesazhin dhe ta konfirmojë atë.

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

Numri i dorëzimit

Numri që shihni në dalje XPENDING — është numri i dorëzimeve të çdo mesazhi. Ky numër rritet në dy mënyra: kur mesazhi kërkohet me sukses përmes XCLAIM ose kur bëhet thirrja XREADGROUP.

Është normale që disa mesazhe dorëzohen disa herë. E rëndësishme është që në fund të gjitha mesazhet të përpunohen. Ndonjëherë, kur përpunoni një mesazh, 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ë atë rast, mund të ndodhë që askush të mos jetë në gjendje ta përpunojë atë mesazh. Duke pasur një numër të përpjekjeve të dorëzimit, mund ta përdorim këtë numër për të zbuluar këto situata. Prandaj, sa herë që numri i dorëzimeve arrin një numër të madh që keni vendosur, ndoshta do të ishte më e mençur ta vendosni atë mesazh në një rrjedhë tjetër dhe të dërgoni një njoftim për administratorin e sistemit.

Gjendja e rrjedhave

Ekipa XINFO përdoret për të kërkuar informacion të ndryshëm rreth rrjedhës dhe grupeve të saj. Për shembull, pamja bazike e komandës është si më poshtë:

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

Komanda më sipër tregon informacionin e përgjithshëm për rrjedhën e caktuar. Tani një shembull pak më të ndërlikuar:

> 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

Komanda më sipër tregon informacionin e përgjithshëm për të gjitha grupet e rrjedhës së caktuar

> 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

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

> XINFO HELP
1) XINFO {subcommand} arg arg ... arg. Subkomandat janë:
2) CONSUMERS {key} {groupname}  -- Trego grupet e abonentëve të grupit {groupname}.
3) GROUPS {key}                 -- Trego grupet e abonentëve të rrjedhës.
4) STREAM {key}                 -- Trego informacion rreth rrjedhës.
5) HELP                         -- Printo këtë ndihmë.

Kufizimi i madhësisë së rrjedhës

Shumë aplikacione nuk duan të mbledhin të dhëna në një rrjedhë për një kohë të gjatë. Shpesh është e dobishme të kesh numrin maksimal të lejuar të mesazheve në rrjedhë. Në raste të tjera, është e dobishme të transferosh të gjitha mesazhet nga rrjedha në një magazinë tjetër të përhershme kur arrihet një madhësi e caktuar e rrjedhës. Mund të kufizosh madhësinë e rrjedhës me parametrin MAXLEN në komandë 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"

Kur përdoret MAXLEN, regjistrimet e vjetra fshihen automatikisht kur arrihet gjatësia e caktuar, duke bërë që rrjedha të ketë një madhësi të përhershme. Megjithatë, mbushja në këtë rast nuk ndodh në mënyrën më të efektshme në memories Redis. Mund ta përmirësosh situatën si më poshtë:

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

Argumenti ~ në shembullin e mësipërm do të thotë se nuk është e nevojshme të kufizosh gjatësinë e rrjedhës në një vlerë specifike. 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 kemi treguar qartazi se duam që rrjedha jonë të ruajë të paktën 1000 regjistrime. Kjo e bën punën me memories shumë më efikase brenda Redis.

Po ashtu ekziston një komandë e veçantë XTRIM, që kryen të njëjtën gjë:

> XTRIM mystream MAXLEN 10

> XTRIM mystream MAXLEN ~ 10

Ruajtja dhe replikimi të përhershëm

Redis Stream replikon asinkronisht në nodet slave dhe ruhet në skedarët AOF (snapshot i të dhënave) dhe RDB (logu i të gjitha operacioneve të shkruara). Po ashtu, mbështetet replikimi i gjendjes së Grupeve të Konsumatorëve. Prandaj, nëse një mesazh është në statusin «pending» në noden master, atëherë në nodet slave ky mesazh do të ketë statusin e njëjtë.

Fshirja e elementeve të veçanta nga rrjedha

Për të fshirë mesazhet ekziston një komandë speciale XDEL. Komanda merr emrin e rrjedhës, duke ndjekur identifikuesit e mesazheve që duhen fshirë:

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

Kur përdoret kjo komandë, duhet të kihet parasysh se në të vërtetë memoria nuk do të lirohet menjëherë.

Rrjedhat me gjatësi zero

Differenca midis rrjedhave dhe strukturave të tjera të dhënash në Redis është se kur strukturat e tjera të dhënash nuk kanë më elemente brenda, si një efekt anësor, struktura e dhënash vetë do të fshihet nga memoria. Për shembull, një grup i renditur do të fshihet plotësisht kur thirrja ZREM të fshijë elementin e fundit. Në vend të kësaj, rrjedhave u lejohet të mbeten në memorie, edhe pa asnjë element brenda.

Përfundimi

Redis Stream është në përputhje të shkëlqyer për krijimin e brokerave mesazhesh, radhëve të mesazheve, regjistrave të unifikuar dhe sistemeve të bisedave që ruajnë historinë.

Siç ka thënë një herë Niklaus Wirth, programet janë algoritme plus struktura të dhënash, dhe Redis ju siguron të dyja.

Burimi: habr.com

Bleni hostim të besueshëm për faqe me mbrojtje nga DDoS, serverë VPS VDS 🔥 Bleni hostim të besueshëm për faqe me mbrojtje nga DDoS, serverë VPS VDS | ProHoster