Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

Tere, Habr!

Tuletame meelde, et koos raamatuga Kafka oleme vÀlja andnud mitte vÀhem huvitava teose raamatukogust Kafka Streams API.

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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 Kafka Streams 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: kuivaine intensiivsetes olukordades ĂŒhtsuse ja sĂŒnkroniseerimise tagamine muutub tĂ”eliseks vĂ€ljakutseks. Andmebaas vĂ”ib muutuda kitsaskohaks vĂ”i sattuda vĂ”istluseseisundisse ja kannatada ettearvamatuse all.

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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:

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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.

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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.

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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:

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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.

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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:

Mitte ainult töötlemine: Kuidas me muutsime Kafka Streams jaotatud andmebaasiks ja mis sellest vÀlja tuli

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

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster