Tere, Habr!
Tuletame 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. Nii ilmus hiljuti artikkel, millega tahame teid tutvustada. Autor jagab oma kogemusi, kuidas luua Kafka Streamsist jaotatud andmehoidla. Head lugemist!
Apache raamatukogu kasutatakse ĂŒle terve maailma ettevĂ”tetes jaotatud voogude töötlemiseks Apache Kafka peal. Ăks selle raamistikku alahinnatud aspekte on see, et see vĂ”imaldab hoida kohalikku olekut, mis on loodud voogude töötlemise pĂ”hjal.
Selles artiklis rÀÀgin, kuidas suutsime oma ettevĂ”ttes seda vĂ”imalust kasulikult rakendada, arendades pilve rakenduste turvalisuse tooteid. Kafka Streams'i abil lĂ”ime mikroteenuseid, millel on jagatud olek, millest igaĂŒhel on meie jaoks tĂ”rkedeta ja kĂ”rgelt kergesti saadaval teabeallikas sĂŒsteemis olevate objektide oleku kohta. Meie jaoks on see edasiminek nii usaldusvÀÀrsuse kui ka hoolduse mugavuse osas.
Kui teid huvitab alternatiivne lĂ€henemine, mis vĂ”imaldab kasutada ĂŒhtset keskvabasid andmebaasi teie objektide formaalse oleku toetamiseks â lugege, see on huvitav...
Miks me arvasime, et on aeg oma lÀhenemist jagatud olekuga muuta
Meil oli vaja toetada erinevate objektide olekute haldamist, toetudes agentide aruannetele (nĂ€iteks: kas saidi vastu tehti rĂŒnnak?). Enne Kafka Streamsile ĂŒleminekut tuginesime sageli oleku haldamiseks ĂŒhele kesksule andmebaasile (+ teenuste API). Selle lĂ€henemise miinused on: ĂŒhtsuse ja sĂŒnkroniseerimise tagamine muutub tĂ”eliseks vĂ€ljakutseks. Andmebaas vĂ”ib muutuda kitsaskohaks vĂ”i sattuda ja kannatada ettearvamatuse all.

Illustratsioon 1: tĂŒĂŒpiline stsenaarium jagatud olekuga, mis esines enne ĂŒleminekut
Kafka ja Kafka Streamsile: agentide aruanded edastatakse API kaudu, uuendatud olek arvutatakse keskses andmebaasis
Tutvuge Kafka Streams'iga â nĂŒĂŒd on lihtne luua mikroteenuseid jagatud olekuga
Umbes aasta tagasi otsustasime tĂ”siselt ĂŒle vaadata meie tööskeemid jagatud olekuga, et tegeleda selliste probleemidega. Otsustasime kohe proovida Kafka Streams'i â on teada, kui skaleeritav, kĂ”rge kĂ€ttesaadavuse ja tĂ”rketaluv sĂŒsteem see on ning kui rikkalik on selle voogude funktsionaalsus (kaasa arvatud, oleku salvestamise vĂ”imalused). Just seda me vajasime, rÀÀkimata sellest, kui kĂŒps ja usaldusvÀÀrne on sĂ”numite vahetamise sĂŒsteem Kafka-s.
Iga meie loodud oleku salvestamise mikroteenus ehitati soovitud Kafka Streams'i instantsi baasil, millel oli ĂŒsna lihtne topoloogia. See koosnes 1) allikast 2) töötlejast, millel on pidev vĂ”tmete ja vÀÀrtuste salvestus 3) vĂ€ljundist:

Illustratsioon 2: meie voogude instantside vaikimisi topoloogia oleku salvestamise mikroteenuste jaoks. Pange tÀhele: siin on ka salvestus, kus on sÀilitatud planeerimise metaandmed.
Uue lĂ€henemise korral koostavad agentuurid sĂ”numeid, mis edastatakse algsesse teema, ja tarbijad â ĂŒtleme, et postituseteenuse teenus â vastu vĂ”tavad arvutatud jagatud oleku vĂ€ljundi (vĂ€ljunditeema) kaudu.

Illustratsioon 3: uus ĂŒlesandev voog jagatud mikroteenuste stsenaariumi jaoks: 1) agent loob sĂ”numi, mis saadetakse Kafka algsesse teema; 2) jagatud oleku mikroteenus (kasutades Kafka Streams'i) töötleb seda ja salvestab arvutatud oleku Kafka lĂ”ppteemasse; seejĂ€rel 3) tarbijad vĂ”tavad uue oleku vastu.
Hei, see sisseehitatud vÔtmete ja vÀÀrtuste salvestus on tÔeliselt kasulik!
Nagu eespool mainitud, sisaldab meie jagatud olekuga topoloogia vÔtmete ja vÀÀrtuste salvestust. Oleme leidnud mitu vÔimalust selle kasutamiseks, ja kaks neist on allpool kirjeldatud.
Variant #1: vÔtmete ja vÀÀrtuste salvestuse kasutamine arvutustes
Meie esimene vĂ”tmete ja vÀÀrtuste hoidla sisaldas abiteavet, mida vajasime arvutamiseks. NĂ€iteks mÀÀrati mĂ”nel juhul jagatud olek âenamushÀÀletuseâ printsiibi jĂ€rgi. Hoidlas sai hoida kĂ”iki viimaseid agentide aruandeid teatud objekti oleku kohta. SeejĂ€rel, kui saime uue aruande mĂ”nelt agentidelt, saime salvestada selle, tĂ”mmata hoidlast vĂ€lja kĂ”ikide teiste agentide aruanded sama objekti oleku kohta ja arvutamist korrata.
Allpool, joonisel 4, on nÀidatud, kuidas avasime juurdepÀÀsu vÔtmete ja vÀÀrtuste hoidla pro tervendava meetodi jaoks, et hiljem oleks vÔimalik uut sÔnumit töödelda.

Joonis 4: avame juurdepÀÀsu vÔtmete ja vÀÀrtuste hoidla pro tervendava meetodi jaoks (pÀrast seda tuleb igas stsenaariumis, mis töötab jagatud olekuga, rakendada meetodit doProcess)
Variant #2: CRUD API loomine Kafka Streams'i kohal
PÀrast meie baatsooside voo seadistamist hakkasime proovima kirjutada RESTful CRUD API meie jagatud olekuga mikroteenustele. Soovisime, et oleks vÔimalik saada teatud vÔi kÔikide objektide olekut ning seadistada vÔi eemaldada objekti olek (see on kasulik serveripoolse toe sÀilitamiseks).
Kuna kĂ”ikide API Get State toetamiseks, kui meil oli vaja olekut uuesti arvutada töötlemise kĂ€igus, salvestasime selle pikaks ajaks sisseehitatud vĂ”tmete ja vÀÀrtuste hoidlasse. Sel juhul on piisav lihtsalt sellise API rakendamine ĂŒhe Kafka Streams'i eksemplari abil, nagu on nĂ€idatud allolevas loendis:

Joonis 5: sisseehitatud vÔtmete ja vÀÀrtuste hoidla kasutamine objekti eelnevalt arvutatud oleku saamiseks
Objekti oleku uuendamine API kaudu ei ole samuti keeruline. PÔhimÔtteliselt peab selle jaoks lihtsalt looma Kafka producer'i ja selle abil tegema kirje, milles on uus olek. Nii tagatakse, et kÔik API kaudu genereeritud sÔnumid töödeldakse tÀpselt samamoodi nagu teised producentide (nt agendid) kaudu saadetud sÔnumid.

Joonis 6: objekti olekut saab seadistada Kafka produceri abil
VĂ€ike keerukus: Kafka'l on palju partitsioone
SeejĂ€rel tahtsime jaotada töötluse koormust ja parandada kĂ€ttesaadavust, luues igale stsenaariumile klastrise mikroteenuste sĂŒsteemi, millel on jagatud olek. Seadistamine oli meie jaoks lihtne: pĂ€rast seda, kui konfigureerisime kĂ”ik instantsid töötama sama rakenduse ID-ga (ja samade alglaadimise serveritega), tehti praktiliselt kĂ”ik muu automaatselt. Samuti seadsime me, et iga sisenditeema koosneb mitmest partiisist, et igale instantsile saaks mÀÀrata partii alamhulga.
Tahan mainida ka, et siin on tavaks teha olekuhoidla varukoopia, et nÀiteks tÔrke taastamisel see koopia teisele instantsile viia. Iga olekuhoidla jaoks Kafka Streamsis luuakse kopeeritav teema muudatuste logiga (kus jÀlgitakse kohalikku vÀrskendust). Nii tagab Kafka pidevalt olekuhoidla kaitse. Seega, kui mÔne instantsi Kafka Streamsis esineb tÔrge, saab olekuhoidla kiiresti taastada teisel instantsil, kuhu liiguvad sobivad partiid. Meie testid nÀitasid, et see toimub mÔne sekundi jooksul isegi siis, kui olekus on miljoneid kirjeid.
Minna ĂŒhest jagatud olekuga mikroteenusest mikroteenuste klastrisse, ei ole Get State API rakendamine enam nii triviaalne. Uues olukorras sisaldab iga mikroteenuse olekuhoidla ainult osa kogu pildist (need objektid, mille vĂ”tmed seondusid konkreetsele partiile). Pidi mÀÀrama, millises instantsis oli vajalik objekti olek, ja tegime seda voogude metaandmete alusel, nagu on allpool nĂ€idatud:

Illustratsioon 7: voogude metaandmete abil mÀÀrame, milliselt instantsilt kĂŒsida vajaliku objekti olekut; sarnast lĂ€henemist kasutati GET ALL API-s.
Peamised jÀreldused
Olekuhoidlad Kafka Streamsis vÔivad de facto toimida jaotatud andmebaasina,
- mida kopeeritakse pidevalt Kafka
- Sellise sĂŒsteemi peale on lihtne luua CRUD API
- Mitme partii töötlemine osutub veidi keerulisemaks
- Samuti on vĂ”imalik lisada voogedastustopoloogiasse ĂŒks vĂ”i mitu olekuhoidlat abiteabe salvestamiseks. Seda varianti saab kasutada:
- Pikaajaline andmete salvestamine, mis on vajalik voogedastuse töötlemisel arvutuste jaoks
- Pikaajaline andmete salvestamine, mis vÔib olla kasulik jÀrgmise voogedastuse instantsi initsialiseerimisel
- paljude teiste jaoks...
Nende ja teiste eeliste tĂ”ttu sobib Kafka Streams suurepĂ€raselt globaalse oleku toetamiseks sellises jaotatud sĂŒsteemis nagu meie. Kafka Streams on tootmiskeskkonnas osutunud vĂ€ga usaldusvÀÀrseks (alates selle juurutamisest pole me praktiliselt ĂŒhtegi sĂ”numit kaotanud) ja oleme kindlad, et sellega ei piirdub tema vĂ”imekus!
Allikas: habr.com
