Tere, Habr!
Tulet meelde, et koos raamatuga oleme vÀlja andnud mitte vÀhem huvitava teose raamatukogust .

Kuni kogukond alles avastab selle vÔimsa tööriista vÔimalusi. Hiljuti ilmus artikkel, millega soovime teid tutvustada. Autor jagab oma kogemust, kuidas luua Kafka Streamsist hajutatud andmesalvestus. Head lugemist!
Apache'i raamatukogu kasutatakse kogu maailmas ettevĂ”tetes hajutatud voogude töötlemiseks Apache Kafka peal. Ăks selle raamistikku alahinnatud aspekte on see, et see vĂ”imaldab salvestada kohaliku oleku, mis on toodetud voogude töötlemise alusel.
Selles artiklis rÀÀgin, kuidas Ă”nnestus meil selle vĂ”imaluse kasulikult rakendada pilve rakenduste turvalisuse toote arendamisel. Kafka Streams'i abiga lĂ”ime mikroteenused, millel on jagatud olek, kus igaĂŒhes teenindatakse meid tĂ”rke- ja kĂ”rge kĂ€ttesaadavuse allikana usaldusvÀÀrse teabe saamiseks objekti oleku kohta sĂŒsteemis. Meie jaoks on see samm edasi nii usaldusvÀÀrsuse kui ka hoolduse mugavuse osas.
Kui olete huvitatud alternatiivsest lĂ€henemisest, mis vĂ”imaldab kasutada ĂŒhte kesksest andmebaasi objekti formaalse oleku toetamiseks â lugege edasi, see on huvitavâŠ
Miks me leidsime, et on aeg muuta oma lÀhenemisviise jagatud oleku kÀsitlemisel
Meil oli vaja toetada erinevate objektide olekut, tuginedes agentide aruannetele (nĂ€iteks: kas veebileht oli rĂŒnnaku alla sattunud)? Enne ĂŒleminekut Kafka Streamsile toetusid me sageli oleku haldamisel ĂŒhele keskasandmebaasile (+ teenuse API). Sellel lĂ€henemisel on oma puudused: talituse jĂ€rjepidevuse ja sĂŒnkroniseerimise tagamine muutub tĂ”eliseks vĂ€ljakutseks. Andmebaas vĂ”ib jĂ”uda kitsaskohaks vĂ”i sattuda ja kannatada ettearvamatuse all.

Illustratsioon 1: tĂŒĂŒpiline olukord, kus jagatud olek oli enne ĂŒleminekut
Kafka ja Kafka Streamsile: agendid edastavad oma arusaamad API kaudu, ajakohastatud olek arvutatakse keskse andmebaasi kaudu
Tutvustame Kafka Streamsi â nĂŒĂŒd on lihtne luua mikroteenuseid jagatud olekuga
Umbes aasta tagasi otsustasime pĂ”hjalikult ĂŒle vaadata oma jagatud oleku tööstsenaariumid, et mĂ”ista selliseid probleeme. Otsustasime kohe proovida Kafka Streams'i â on teada, kui skaleeritav, suure kĂ€ttesaadavuse ja tĂ”rke suhtes vastupidav see on, samuti kui rikkalik on selle voogedastuse funktsionaalsus (sealhulgas oleku sĂ€ilitamisega transformatsioonid). Just seda vajasime, rÀÀkimata, kui kĂŒps ja usaldusvÀÀrne on Kafka sĂ”numite vahetussĂŒsteem.
Iga loodud oleku sĂ€ilitamisega mikroteenus pĂ”hines pĂ€ris lihtsa topoloogiaga Kafka Streams'i instantsil. See koosnes 1) allikast 2) protsessorist koos pĂŒsiva vĂ”tme-vÀÀrtuse salvestusega 3) voost:

Illustreerimine 2: vaikimisi topoloogia meie voogedastuse instantside jaoks mikroteenustes, kus on oleku sÀilitamine. Pane tÀhele: siin on ka salvestus, kus hoitakse andmeid planeerimise kohta.
Uue lĂ€henemise kohaselt koostavad agentuurid sĂ”numid, mis saadetakse algsesse teemasse, ja tarbijad â nĂ€iteks meiliteenused â saavad arvutatud jagatud oleku voost (vĂ€ljaanneteema) tĂ€ielikult.

Joonis 3: uus ĂŒlesannete voog jagatud mikroteenuste stsenaariumi jaoks: 1) agent genereerib sĂ”numi, mille kaudu jĂ”uab see Kafka algsesse teemasse; 2) jagatud olekuga mikroteenus (kasutades Kafka Streams'i) töötleb selle ja salvestab arvutatud oleku Kafka lĂ”ppteemasse; pĂ€rast seda 3) tarbijad saavad uue oleku.
Hei, see sisseehitatud vÔtmete ja vÀÀrtuste hoidja on tÔeliselt kasulik!
Kuidas mainitud, sisaldab meie jagatud olekute topoloogia vÔtmete ja vÀÀrtuste hoidjat. Oleme leidnud mitmeid kasutusvÔimalusi, millest kaks on allpool kirjeldatud.
Variant #1: vÔtmete ja vÀÀrtuste hoidja kasutamine arvutustes.
Meie esimene vĂ”tmete ja vÀÀrtuste hoidla sisaldas abiteavet, mis oli vajalik meie arvutuste jaoks. NĂ€iteks mĂ”ningatel juhtudel mÀÀrati jagatud olek âenamuse hÀÀlteâ pĂ”himĂ”tte alusel. Hoidlas sai hoida kĂ”iki viimaseid agente raportite kohta mingist objektist. Saades seejĂ€rel uue raporti ĂŒhelt vĂ”i teiselt agentilt, saime selle salvestada, vĂ€lja tĂ”mmata hoidlast teiste agentide raportid sama objekti kohta ning uued arvutused korrata.
Allpool on nÀidatud, kuidas avasime pÀÀsu vÔtmete ja vÀÀrtuste hoidla töötlevasse protsessori meetodisse, et seejÀrel saaks töödelda uut sÔnumit.

Illustratsioon 4: avame pÀÀsu vÔtmete ja vÀÀrtuste hoidla töötleva protsessori meetodi jaoks (pÀrast seda peab iga stsenaarium, mis töötab jagatud olekuga, implementima meetodi doProcess)
Variant #2: CRUD API loomine Kafka Streams'i kohal
PÀrast meie pÔhijÀrjestuse loomist oleme proovinud kirjutada RESTful CRUD API meie jagatud olekuga mikroteenuste jaoks. Soovisime, et oleks vÔimalik vÀlja tuua teatud vÔi kÔigi objektide olek ning samuti seada vÔi eemaldada objekti olekut (see on kasulik serveri toe pakkumisel).
Selleks, et toetada kĂ”iki API Get State, igal korral, kui oli vaja uuesti arvutada olekut töötlemise kĂ€igus, salvestasime selle pikaks ajaks sisseehitatud vĂ”tme-vÀÀrtuse salvestusse. Sellisel juhul on piisavalt lihtne rakendada sellist API ĂŒhe Kafka Streams eksemplariga, nagu on nĂ€idatud allpool olevas loendis.

Illustratsioon 5: sisseehitatud vÔtme-vÀÀrtuse salvestuse kasutamine objekti eelarvutatava oleku saamiseks.
Objekti oleku vÀrskendamine API kaudu on samuti lihtsalt rakendatav. PÔhimÔtteliselt on sellega vaid vajalik luua Kafka tootja ja selle abil teha kirje, milles on uus olek. Nii tagatakse, et kÔik API kaudu genereeritud sÔnumid töödeldakse tÀpselt samamoodi nagu muude tootjate (nt agentide) poolt saadetud sÔnumid.

Illustreerimine 6: objekti olek muutub kasutades Kafka produtsenti
VĂ€ike komplikatsioon: Kafka-l on palju partiisid
Soovisime koormust, mis oli seotud töötlemisega, jaotada ning kÀttesaadavust parandada, pakkudes igale stsenaariumile klastrit mikroteenuseid jagatud olekuga. Seade sujus meil lihtsasti: pÀrast seda, kui konfigureerisime kÔik instantsid töötama sama rakenduse ID-ga (ja samade alglaadimise serveritega), tehti praktiliselt kÔik muu automaatselt. Me mÀÀrasime ka, et iga lÀhteteema koosneb mitmest partiist, et igale instantsile saaks mÀÀrata partii alamkogumi.
Mainin ka, et siin on tavaline teha olekute hoiukoha varukoopia, et nĂ€iteks tĂ”rke korral saaks selle koopia teisele instantsile ĂŒle kanda. Iga oleku hoiukoha jaoks Kafka Streams'is luuakse replikatsioonitopik muudatuste logiga (kus jĂ€lgitakse kohalikke vĂ€rskendusi). Seega tagab Kafka pidevalt oleku hoiukoha turvalisuse. SeetĂ”ttu saab oleku hoiukohta kiiresti taastada teisel instantsil, kuhu vastavad partiid ĂŒleviivad. Meie testid nĂ€itasid, et see toimub mĂ”ne sekundi jooksul, isegi kui hoiukohta on salvestatud miljoneid kirjeid.
Ăleminek ĂŒhe jagatud olekuga mikroteenuselt klastrile mikroteenuseid muudab Get State API rakendamise vĂ€hem triviaalseks. Uues olukorras sisaldab iga mikroteenuse olekute salvestuses ainult osa tervikpildist (need objektid, mille vĂ”tmed vastavad konkreetsele partitsioonile). Pidi vĂ€lja selgitama, millises instantsis vajalik objekti olek asub, ja tegime seda voogude metaandmete alusel, nagu on nĂ€idatud allpool:

Illustratsioon 7: voogude metaandmete abil mÀÀrame, milliselt instantsilt kĂŒsida vajaliku objekti olekut; sarnast lĂ€henemist rakendati ka GET ALL API puhul.
Peamised jÀreldused
Kafka Striimis olevad olekute salvestused vÔivad de facto toimida jaotatud andmebaasina,
- mis on pidevalt replitseeritav Kafka-s.
- Selle sĂŒsteemi peal on lihtne ĂŒles ehitada CRUD API.
- Mitme partitsiooni töötlemine osutub veidi keerulisemaks.
- Samuti on vĂ”imalik lisada ĂŒks vĂ”i mitu olekute salvestust voogude topoloogiasse, et salvestada abiteavet. Sellist varianti vĂ”ib kasutada:
- Pikaajalise andmete salvestamiseks, mida vajatakse voogude töötlemise arvutustes.
- Pikaandmeid, mis vÔivad olla kasulikud jÀrgmise voogedastuse instantsi algseadmistes.
- palju muudâŠ
Nende ja teiste eeliste tĂ”ttu sobib Kafka Streams suurepĂ€raselt globaalse oleku toetamiseks sellises jaotatud sĂŒsteemis nagu meie. Kafka Streams on osutunud vĂ€ga usaldusvÀÀrseks tootmisjĂ€rgus (alates selle juurutamisest pole me praktiliselt kaotanud sĂ”numeid), ja me oleme kindlad, et sellevĂ”imalused ei piirdu sellega!
Allikas: habr.com
