Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale» Përshëndetje, Habrushitës! Ky libër është i përshtatshëm për çdo zhvillues që dëshiron të kuptojë përpunimin e rrjedhave. Kuptimi i programimit të shpërndarë do t'ju ndihmojë të mësoni më mirë Kafka dhe Kafka Streams. Do të ishte mirë të njihni ndonjëherë edhe vete strukturën Kafka, por nuk është e nevojshme: unë do t'ju tregoj gjithçka që ju nevojitet. Zhvilluesit e avancuar të Kafka, si dhe ata fillestarë, do të mësojnë të krijojnë aplikacione interesante për përpunimin e rrjedhave me ndihmën e bibliotekës Kafka Streams. Zhvilluesit JAVA të nivelit të mesëm dhe të lartë, që janë tashmë të njohur me koncepte si serializimi, do të mësojnë si të aplikojnë aftësitë e tyre për të krijuar aplikacione Kafka Streams. Kodi burimor i librit është shkruar në Java 8 dhe e përdor ndjeshëm sintaksën e shprehjeve lambda në Java 8, kështu që aftësia për të punuar me funksionet lambda (edhe në një gjuhë tjetër programuese) do t'ju shërbejë.

Shkurtimi. 5.3. Agregimi dhe operacionet mbi dritaret

Në këtë seksion do të kalojmë në studimin e pjesëve më premtuese të Kafka Streams. Deri tani kemi shqyrtuar këto aspekte të Kafka Streams:

  • krijimi i topologjisĂ« sĂ« pĂ«rpunimit;
  • pĂ«rdorimi i gjendjes nĂ« aplikacionet e rrjedhĂ«s;
  • kryerja e lidhjeve tĂ« rrjedhave tĂ« tĂ« dhĂ«nave;
  • dallimet midis rrymave tĂ« ngjarjeve (KStream) dhe rrymave tĂ« azhurnimeve (KTable).

NĂ« shembujt e ardhshĂ«m ne do tĂ« mbledhim tĂ« gjitha kĂ«to elemente sĂ« bashku. PĂ«r mĂ« tepĂ«r, do tĂ« njiheni me operacionet dritare — njĂ« mundĂ«si tjetĂ«r fantastike pĂ«r aplikacionet e rrjedhĂ«s. Shembulli ynĂ« i parĂ« do tĂ« jetĂ« agregimi i thjeshtĂ«.

5.3.1. Agregimi i vëllimit të shitjeve të aksioneve sipas industrive

Agregimi dhe grupimi janë mjete thelbësore kur punohet me të dhëna rrjedhëse. Kërkimi i regjistrimeve individuale përmes ardhjes shpesh nuk është i mjaftueshëm. Për të nxjerrë informacione shtesë nga të dhënat, nevojitet grupimi dhe kombinimi i tyre.

Në këtë shembull, ju do të vishni kostumin e tregtarit të brendshëm, i cili duhet të ndjekë vëllimet e shitjeve të aksioneve të kompanive në disa industri. Në veçanti, ju interesojnë pesë kompani me volumin më të lartë të shitjeve të aksioneve në secilën nga industrinë.

Për një agregim të tillë do të nevojiten disa hapa të ardhshëm për të përkthyer të dhënat në formën e duhur (kur flasim në përgjithësi).

  1. Krijoni njĂ« burim mbi njĂ« temĂ« qĂ« botoni informacion tĂ« papĂ«rpunuar nĂ« tregtimin e aksioneve. Ne do tĂ« duhet tĂ« shfaqim njĂ« objekt tĂ« tipit StockTransaction nĂ« njĂ« objekt tĂ« tipit ShareVolume. Çështja Ă«shtĂ« se objekti StockTransaction pĂ«rmban metadata tĂ« shitjeve, ndĂ«rsa ne na nevojiten vetĂ«m tĂ« dhĂ«nat rreth numrit tĂ« aksioneve tĂ« shitur.
  2. Gruponi të dhënat e ShareVolume sipas simboleve të aksioneve. Pas grupimit sipas simboleve, është e mundur të përmbledhni këto të dhëna në shumat ndërmjetëse të volumit të shitjeve të aksioneve. Vlen të theksohet se metoda KStream.groupBy kthen një instancë të tipit KGroupedStream. Ndërsa për të marrë një instancë të KTable, mund të thirrni më tej metodën KGroupedStream.reduce.

ÇfarĂ« Ă«shtĂ« ndĂ«rfaqja KGroupedStream

Metodat KStream.groupBy dhe KStream.groupByKey kthejnë një instancë të KGroupedStream. KGroupedStream është një paraqitje ndërmjetëse e rrjedhës së ngjarjeve pas grupimit sipas çelësave. Ajo nuk është aspak e destinuar për përdorim të drejtpërdrejtë. Në vend të kësaj, KGroupedStream përdoret për operacione agreguese, rezultat i të cilave është gjithmonë KTable. Dhe për sa kohë që rezultati i operacioneve agreguese është KTable dhe në to aplikohet një depo e shtetit, ndoshta jo të gjitha përditësimet rezultuese dërgohen më tutje në pajisje.

Metoda KTable.groupBy kthen njĂ« KGroupedTable tĂ« ngjashĂ«m — njĂ« pĂ«rfaqĂ«sim tĂ« mesĂ«m tĂ« njĂ« rrjedhe azhurnimesh tĂ« rregulluara sipas çelikut.

Le të bëjmë një pushim të vogël dhe të shohim figurën 5.9, e cila tregon arritjet tona. Kjo topologji duhet të jetë tashmë e njohur për ju.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Tani le të shohim kodin për këtë topologji (mund ta gjeni në skedarin src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (listing 5.2).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Kodi i paraqitur veçohet për shkurtësinë dhe përbërjen e madhe të veprimeve të kryera në disa rreshta. Në parametrin e parë të metodës builder.stream mund të vëreni diçka të re për ju: vlerën e enumeruar të tipit AutoOffsetReset.EARLIEST (ka gjithashtu LATEST), e cila caktohet përmes metodës Consumed.withOffsetResetPolicy. Me këtë tip enumerimi mund të specifikoni strategjinë e rikthimit të pozicionit për çdo KStream ose KTable, që ka prioritet mbi parametrin e rikthimit të pozicionit nga konfigurimi.

GroupByKey dhe GroupBy

Në ndërfaqen KStream, ka dy metoda për grupimin e regjistrimeve: GroupByKey dhe GroupBy. Të dyja kthejnë KGroupedTable, kështu që mund të keni një pyetje të arsyeshme: çfarë ndan ata dhe kur të përdorni cilin prej tyre?

Metoda GroupByKey përdoret kur çelësat në KStream tashmë janë të mbushur. E rëndësishme është se flamuri "kërkon ripjesëzim" nuk është vendosur kurrë.

Metoda GroupBy parashikon që ju keni ndryshuar çelësat për grupim, kështu që flamuri i ripjesëzimit është vendosur në true. Ekzekutimi pas metodës GroupBy për lidhje, agregatë, etj. do të çojë në ripjesëzim automatik.
Përmbledhje: duhet të përdoret sa më shumë të jetë e mundur GroupByKey, dhe jo GroupBy.

ÇfarĂ« bĂ«jnĂ« metodat mapValues dhe groupBy — Ă«shtĂ« e qartĂ«, prandaj le t'i hedhim njĂ« vĂ«shtrim metodĂ«s sum() (e cila mund tĂ« gjendet nĂ« skedarin src/main/java/bbejeck/model/ShareVolume.java) (listimi 5.3).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Metoda ShareVolume.sum kthen shumĂ«n ndĂ«rmjetĂ«se tĂ« volumit tĂ« shitjeve tĂ« aksioneve, dhe rezultati i gjithĂ« zinxhirit tĂ« llogaritjeve paraqet njĂ« objekt KTable. Tani e kuptoni rolin qĂ« luan KTable. Kur objekte ShareVolume arrijnĂ«, pĂ«rditĂ«simi mĂ« i fundit i rĂ«ndĂ«sishĂ«m ruhet nĂ« objektin pĂ«rkatĂ«s KTable. ËshtĂ« e rĂ«ndĂ«sishme tĂ« mos e harroni se tĂ« gjitha pĂ«rditĂ«simet reflektohen nĂ« shareVolumeKTable paraprak, por jo tĂ« gjitha dĂ«rgohen mĂ« tej.

Më pas, me këtë KTable, ne do të kryejmë agregimin (sipër numrit të aksioneve të shitur), për të marrë pesë kompani me volumin më të madh të shitjeve të aksioneve në secilën nga industrive. Veprimet tona do të jenë të ngjashme me ato gjatë agregimit të parë.

  1. Të kryejmë një tjetër operacion groupBy për të grupuar objekte të veçanta ShareVolume sipas industrive.
  2. Të fillojmë me përllogaritjen e objekteve ShareVolume. Këtë herë, objekti i agregimit përfaqëson një radhë me përparësi të caktuar. Në një radhë me përparësi të caktuar ruhen vetëm pesë kompani me numrin më të madh të aksioneve të shitura.
  3. Të shfaqim radhët nga pika e mëparshme në një vlerë string dhe të kthejmë pesë aksionet më të shitura sipas numrit të aksioneve në industri.
  4. Të shkruajmë rezultatet në formë string në temë.

Në fig. 5.10, tregohet grafiku i topologjisë së lëvizjes së të dhënave. Siç e shihni, cikli i dytë i përpunimit është mjaft i thjeshtë.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Tani, pasi keni kuptuar qartë strukturën e këtij rrethi të dytë të përpunimit, mund të kaloni në kodin e tij origjinal (do ta gjeni në skedarin src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (listing 5.4).

Në këtë inicializues ka një variabël fixedQueue. Kjo është një objekt i personalizuar - një adapter për java.util.TreeSet, i cili përdoret për të ndjekur N rezultatet më të mëdha në rendin në rënie të numrit të aksioneve të shitura.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Ju keni hasur tashmë me thirrjet groupBy dhe mapValues, prandaj nuk do të ndalemi tek ato (ne thërrasim metodën KTable.toStream, pasi metoda KTable.print është e zastaruar). Por ende nuk keni parë versionin KTable të metodës aggregate(), kështu që do të kalojmë pak kohë për ta diskutuar atë.

Si e mbani mend, KTable dallon nga fakti se regjistrimet me çelësa të njëjtë merren si përditësime. KTable zëvendëson regjistrimin e vjetër me të riun. Agregimi ndodh në një mënyrë të ngjashme: regjistrimet më të fundit me një çelës mblidhen. Kur një regjistrim mbërrin, ai shtohet në instancën e klasës FixedSizePriorityQueue me ndihmën e somatorit (parametri i dytë në thirrjen e metodës aggregate), por nëse ekziston tashmë një regjistrim tjetër me të njëjtin çelës, regjistrimi i vjetër hiqet me ndihmën e zbritësit (parametri i tretë në thirrjen e metodës aggregate).

Kjo do të thotë se agregatori ynë, FixedSizePriorityQueue, nuk agregon të gjitha vlerat me një çelës të vetëm, por ruan shumën lëvizëse të sasisë së N llojeve më të shitura të aksioneve. Në çdo regjistrim që mbërrin përmban numrin total të aksioneve të shitura deri tani. KTable do t'ju japë informacion se cilat kompani kanë më shumë aksione që shiten në këtë moment, agregimi lëvizës i çdo përditësimi nuk është i nevojshëm.

Ne mësuam të bëjmë dy gjëra të rëndësishme:

  • grupojmĂ« vlerat nĂ« KTable sipas çelĂ«sit tĂ« pĂ«rbashkĂ«t;
  • kryejmĂ« operacione tĂ« dobishme mbi kĂ«to vlera tĂ« grupuara, siç janĂ« reduktimi dhe agregimi.

Aftësia për të kryer këto operacione është e rëndësishme për të kuptuar kuptimin e të dhënave që lëvizin përmes aplikacionit Kafka Streams dhe për të zbuluar se çfarë informacioni mbartin ato.

Ne gjithashtu e kemi lidhur sĂ« bashku disa nga konceptet kyçe qĂ« janĂ« diskutuar mĂ« parĂ« nĂ« kĂ«tĂ« libĂ«r. NĂ« kapitullin 4, folĂ«m pĂ«r rĂ«ndĂ«sinĂ« e qĂ«ndrueshmĂ«risĂ« dhe gjendjes lokale pĂ«r aplikacionet e rrjedhĂ«s. Shembulli i parĂ« nga ky kapitel demonstroi pse Ă«shtĂ« kaq e rĂ«ndĂ«sishme gjendja lokale — ajo ofron mundĂ«sinĂ« pĂ«r tĂ« ndjekur se çfarĂ« informacioni keni parĂ« tashmĂ«. Qasja lokale ndihmon pĂ«r tĂ« shmangur vonesat nĂ« rrjet, duke e bĂ«rĂ« aplikacionin mĂ« tĂ« shpejtĂ« dhe mĂ« tĂ« qĂ«ndrueshĂ«m ndaj gabimeve.

Kur realizoni çdo operacion të reduktimit ose përmbledhjes, duhet të tregoni emrin e depozitës së gjendjes. Operacionet e reduktimit dhe përmbledhjes kthejnë një ekzemplar të KTable, dhe KTable përdor depozitën e gjendjes për të zëvendësuar rezultatet e vjetra me të reja. Siç e keni parë, jo të gjitha përditësimet dërgohen më tej në tub, dhe kjo është e rëndësishme, sepse operacionet e përmbledhjes janë të destinuara për të marrë informacionin përfundimtar. Nëse nuk përdoret gjendja lokale, KTable do të dërgojë më tej të gjitha rezultatet e përmbledhjes dhe reduktimit.

MĂ« pas do tĂ« shohim si realizohen operacione tĂ« tilla si pĂ«rmbledhja brenda njĂ« intervali tĂ« caktuar kohor — ato qĂ« quhen operacione dritare (windowing operations).

5.3.2. Operacionet dritare

Në seksionin e mëparshëm u njohëm me reduktimin dhe përmbledhjen «lëvizëse». Aplikacioni realizoi një reduktim të vazhdueshëm të volumit të shitjeve të aksioneve duke e përmbledhur më pas pesë aksionet më të shitura në bursë.

Ndonjëherë, agregimi dhe kompresimi i vazhdueshëm i rezultateve është i nevojshëm. Ndërsa herë të tjera duhet të kryhen operacione vetëm për një interval të caktuar kohor. Për shembull, llogaritur se sa operacione tregtare janë kryer me aksionet e një kompanie të caktuar në 10 minutat e fundit. Ose sa përdorues klikuan në një banner të ri reklamor në 15 minutat e fundit. Aplikacioni mund të kryejë këto operacione përsëritësisht, por me rezultatet që lidhen vetëm me intervalet e caktuara kohore (dritaret e kohës).

Numërimi i transaksioneve bursore sipas blerësit

NĂ« shembullin e mĂ«poshtĂ«m, ne do tĂ« merremi me ndjekjen e transaksioneve bursore pĂ«r disa tregtarĂ« — ose organizata tĂ« mĂ«dha, ose financistĂ« tĂ« zgjuar dhe tĂ« vetmuar.

Ka ekzistojnë dy shkak të mundshëm për një ndjekje të tillë. Njëra është nevoja për të ditur se çfarë blejnë/shesin liderët e tregut. Nëse këta lojtarë të mëdhenj dhe investitorët e shkathët shohin mundësi që hapen për ta, ka kuptim të ndiqen strategjitë e tyre. Arsyeja tjetër është dëshira për të vënë re çdo shenjë të mundshme të marrëveshjeve të paligjshme duke përdorur informacionin e brendshëm. Për këtë, do t'ju nevojitet të analizoni korelacionin e shpërthimeve të mëdha të shitjeve me njoftimet e rëndësishme për shtyp.

Kjo ndjekje përbëhet nga faza të tilla si:

  • krijimi i njĂ« rrjedhe leximi nga tema stock-transactions;
  • grupimi i tĂ« dhĂ«nave hyrĂ«se sipas identifikuesit tĂ« blerĂ«sit dhe simbolit tĂ« tregtisĂ« sĂ« aksioneve. Thirrja e metodĂ«s groupBy kthen njĂ« instancĂ« tĂ« klasĂ«s KGroupedStream;
  • kthimi i njĂ« rrjedhe tĂ« dhĂ«nash nga KGroupedStream.windowedBy e kufizuar nĂ« njĂ« dritare kohore, e cila lejon realizimin e agregimit nĂ« dritare. NĂ« varĂ«si tĂ« llojit tĂ« dritares, kthehet ose TimeWindowedKStream ose SessionWindowedKStream;
  • numĂ«rimi i transaksioneve pĂ«r operacionin e agregimit. Rryma e dhĂ«nash nĂ« dritare pĂ«rcakton nĂ«se njĂ« regjistrim i caktuar pĂ«rfshihet nĂ« kĂ«tĂ« numĂ«rim;
  • shĂ«nimi i rezultateve nĂ« njĂ« temĂ« ose dalja e tyre nĂ« konsolĂ« gjatĂ« zhvillimit.

Topologjia e kësaj aplikacioni është e thjeshtë, por një imazh vizual nuk do të dëmtonte. Le të hedhim një sy në fig. 5.11.

Më pas do të shqyrtojmë funksionalitetin e operacioneve me dritare dhe kodin përkatës.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»

Llojet e dritareve

NĂ« Kafka Streams ka tre lloje dritaresh:

  • seancĂ«;
  • «lidhen» (tumbling);
  • lĂ«vizĂ«se/«skakuese» (sliding/hopping).

Cilin tĂ« zgjidhni varet nga kĂ«rkesat e biznesit. Dritaret «lidhen» dhe «skakuese» janĂ« kĂ«to me kufizime nĂ« kohĂ«, ndĂ«rsa kufizimet e dritareve tĂ« seancĂ«s varen nga veprimet e pĂ«rdoruesve — kohĂ«zgjatja e seancĂ«s (-ave) pĂ«rcaktohet vetĂ«m nga sa aktivisht vepron pĂ«rdoruesi. E rĂ«ndĂ«sishme Ă«shtĂ« tĂ« mos harrohet se tĂ« gjitha llojet e dritareve bazohen nĂ« etiketat e datĂ«s/koha tĂ« shĂ«nimeve, jo nĂ« kohĂ«n sistemike.

Më pas do ta realizojmë topologjinë tonë me secilin nga llojet e dritareve. Kodi i plotë do të jepet vetëm në shembullin e parë, për llojet e tjera të dritareve asgjë nuk do të ndryshojë, përveç llojit të operacionit të dritares.

Dritaret e seancës

Dritaret e sesioneve shumë ndahen nga të gjitha llojet e tjera të dritareve. Ato kufizohen jo aq sa në kohë, sa në aktivitetin e përdoruesit (ose aktivitetin e asaj entiteti që dëshironi të gjurmoni). Dritaret e sesioneve ndahen nga periudha të papunësisë.

Figura 5.12 ilustron konceptin e dritareve të sesioneve. Një sesion më i vogël do të bashkohet me sesionin në të majtë të tij. Ndërsa sesioni në të djathtë do të jetë i veçantë, pasi pason një periudhë të gjatë papunësie. Dritaret e sesioneve bazohen në veprimet e përdoruesve, por aplikojnë etiketat e datës/kohës nga regjistrimet për të përcaktuar se cilit sesion i përket regjistrimi.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»

Përdorimi i dritareve të sesioneve për të gjurmuar transaksionet në treg

Do të përdorim dritaret e sesioneve për të kapur informacion mbi transaksionet në treg. Zhvillimi i dritareve të sesioneve është demonstruar në listimin 5.5 (i cili mund të gjendet në skedarin src/main/java/bbejeck/chapter_5/CountingWindowingAndKTableJoinExample.java).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Shumica e operacioneve të kësaj topologjie tashmë i keni parë, kështu që nuk ka nevojë t'i shqyrtojmë ato këtu përsëri. Por ka disa elemente të reja këtu, të cilat do t'i diskutojmë tani.

Në çdo operacion groupBy zakonisht kryhet një operacion agregimi (agregim, përmbledhje ose numërim). Mund të bëhet ose agregim kumulativ me një total në rritje, ose agregim dritare ku merret parasysh regjistrohet brenda një dritare kohore të caktuar.

Kodi në listimin 5.5 kryen numërimin e transaksioneve brenda dritareve të seancës. Në fig. 5.13 këto veprime analizohen hap pas hapi.

Me thirrjen windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)) krijojmë një dritare seance me një interval inaktiviteti prej 20 sekondash dhe një periudhë ruajtjeje prej 15 minutash. Intervali i inaktivitetit prej 20 sekondash do të thotë se aplikacioni do të përfshijë çdo regjistrim që arrin brenda 20 sekondash nga përfundimi ose fillimi i seancës aktuale në këtë (aktive) seancë.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
MĂ« pas, ne tregojmĂ« se cila operacioni agregimi duhet tĂ« kryhet nĂ« dritaren e sesionit — nĂ« kĂ«tĂ« rast count. NĂ«se regjistrimi hyrĂ«s tejkalon intervalin e mosaktivitetit (nĂ« cilindo nga anĂ«t e timestamps), aplikacioni krijon njĂ« sesion tĂ« ri. Intervali i ruajtjes do tĂ« thotĂ« mbajtjen e sesionit pĂ«r njĂ« periudhĂ« tĂ« caktuar dhe lejon tĂ« dhĂ«na qĂ« vonohen dhe dalin jashtĂ« periudhĂ«s sĂ« mosaktivitetit tĂ« sesionit, por qĂ« ende mund tĂ« bashkohen. PĂ«r mĂ« tepĂ«r, fillimi dhe fundi i sesionit tĂ« ri, qĂ« rezulton nga bashkimi, pĂ«rputhen me timestamp-in mĂ« tĂ« hershĂ«m dhe mĂ« tĂ« vonĂ«.

Le të shqyrtojmë disa regjistrime nga metoda count, për të parë se si funksionojnë sesionet (Tabela 5.1).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Kur regjistrimet arrijnĂ«, ne kĂ«rkojmĂ« pĂ«r sesionet ekzistuese me tĂ« njĂ«jtin çelĂ«s, kohĂ«n e pĂ«rfundimit mĂ« tĂ« vogĂ«l se timestamp-i aktual — intervali i mosaktivitetit dhe kohĂ«n e fillimit mĂ« tĂ« madhe se timestamp-i aktual + intervali i mosaktivitetit. Duke marrĂ« kĂ«tĂ« parasysh, katĂ«r regjistrime nga Tabela 5.1 bashkohen nĂ« njĂ« sesion tĂ« vetĂ«m nĂ« kĂ«tĂ« mĂ«nyrĂ«.

1. E para që arrin është regjistrimi 1, kështu që koha e fillimit është e barabartë me kohën e përfundimit dhe është 00:00:00.

2. Më pas vjen regjistrimi 2, dhe kërkojmë seancat që përfundojnë jo më herët se 23:59:55 dhe fillojnë jo më vonë se 00:00:35. Gjejmë regjistrimin 1 dhe bashkojmë seancat 1 dhe 2. Marrim kohën e fillimit të seancës 1 (më e hershme) dhe kohën e përfundimit të seancës 2 (më të vonshme), pra seanca jonë e re fillon në 00:00:00 dhe përfundon në 00:00:15.

3. Vjen regjistrimi 3, kërkojmë seancat midis 00:00:30 dhe 00:01:10 dhe nuk gjejmë asnjë. Shtojmë seancën e dytë për çelësin 123-345-654,FFBE, që fillon dhe përfundon në 00:00:50.

4. Vjen regjistrimi 4, dhe kĂ«rkojmĂ« seancat midis 23:59:45 dhe 00:00:25. KĂ«tĂ« herĂ« gjejmĂ« tĂ« dyja seancat — 1 dhe 2. TĂ« gjitha tre seancat bashkohen nĂ« njĂ«, me kohĂ«n e fillimit 00:00:00 dhe kohĂ«n e pĂ«rfundimit 00:00:15.

Nga ajo çfarë është tregohet në këtë seksion, është e rëndësishme të mbani mend këto nuanca të rëndësishme:

  • seancat — nuk janĂ« dritare me pĂ«rmasa fikse. KohĂ«zgjatja e njĂ« seance pĂ«rcaktohet nga aktiviteti brenda njĂ« periudhe tĂ« caktuar kohore;
  • tĂ« dhĂ«nat e datĂ«s/kohĂ«s nĂ« tĂ« dhĂ«na pĂ«rcaktojnĂ« nĂ«se njĂ« ngjarje bĂ«n pjesĂ« nĂ« njĂ« seancĂ« ekzistuese ose nĂ« njĂ« periudhĂ« pa aktivitet.

MĂ« pas ne do tĂ« diskutojmĂ« pĂ«r llojin tjetĂ«r tĂ« dritareve — "tĂ« rrotulluara".

"Dritaret e rrotulluara"

Dritaret "tumbling" kapin ngjarjet që bien në një periudhë të caktuar kohe. Imagjinoni sikur duhet të kapni të gjitha transaksionet bursë të një kompanie çdo 20 sekonda, kështu që mbledhni të gjitha ngjarjet për këtë periudhë kohe. Pas përfundimit të intervalit 20-sekondësh, dritarja "kuvyrka" dhe kalon në një interval të ri 20-sekondësh të vëzhgimit. Figura 5.14 ilustron këtë situatë.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Siç mund ta shihni, të gjitha ngjarjet e ardhura gjatë 20 sekondave të fundit janë të përfshira në dritare. Pas përfundimit të kësaj periudhe kohore, krijohet një dritare e re.

Në listingun 5.6 është paraqitur kodi që demonstruar përdorimin e dritareve "kuvyrka" për të kapur transaksionet bursë çdo 20 sekonda (mund ta gjeni në skedarin src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Për shkak të këtij ndryshimi të vogël në thirrjen e metodës TimeWindows.of, mund të përdoret dritarja "kuvyrk". Në këtë shembull, nuk ka thirrje të metodës until(), duke përdorur intervalin e ruajtjes në mënyrë default, që është 24 orë.

SĂ« fundi, Ă«shtĂ« koha tĂ« kalojmĂ« te opsioni i fundit i dritareve — dritaret "hopping".

Dritaret lëvizëse ("hopping")

Dritaret e kalimit/"kërcitëse" janë të ngjashme me dritaret "rrotulluese", por me një ndryshim të vogël. Dritaret e kalimit nuk presin përfundimin e intervalit të kohës për të krijuar një dritare të re për të trajtuar ngjarjet e fundit. Ato fillojnë llogaritjet e reja pas një intervali pritjeje më të vogël se sa zgjatja e dritares.

Për të ilustruar dallimet midis dritareve "rrotulluese" dhe "kërcitëse", le të kthehemi te shembulli i numërimit të transaksioneve të bursës. Qëllimi ynë mbetet numërimi i numrit të transaksioneve, por nuk dëshirojmë të presim të gjithë intervalin e kohës përpara se të azhurnojmë numëruesin. Në vend të kësaj, do të azhurnojmë numëruesin çdo herë pas intervaleve më të shkurtër të kohës. Për shembull, do të numërojmë transaksionet çdo 20 sekonda, por do ta azhurnojmë numëruesin çdo 5 sekonda, siç tregohet në fig. 5.15. Kështu, do të kemi tre dritare rezultatesh me të dhëna të mbivendosura.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Në listimin 5.7 është paraqitur kodi për përcaktimin e dritareve të kalimit (mund të gjendet në skedarin src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Një dritare "fëlliqëse" mund të shndërrohet në një "kërcimtare" duke shtuar thirrjen e metodës advanceBy(). Në shembullin e dhënë, intervali i ruajtjes është 15 minuta.

Në këtë seksion, keni parë si të kufizoni rezultatet e agregimit me dritare të përkohshme. Sidomos, do të doja që të mbani mend nga ky seksion tri gjëra të mëposhtme:

  • pĂ«rmasat e dritareve tĂ« sesioneve kufizohen jo nga intervali i kohĂ«s, por nga aktiviteti i pĂ«rdoruesve;
  • dritaret "fĂ«lliqĂ«se" ofrojnĂ« njĂ« pasqyrĂ« tĂ« ngjarjeve brenda njĂ« periudhe tĂ« caktuar;
  • kohĂ«zgjatja e dritareve "kĂ«rcimtare" Ă«shtĂ« e fiksuar, por ato pĂ«rditĂ«sohen shpesh dhe mund tĂ« pĂ«rmbajnĂ« nĂ« tĂ« gjitha dritaret regjistrime qĂ« prekĂ«n.

Më pas do të mësojmë si të kthejmë KTable përsëri në KStream për lidhje.

5.3.3. Lidhja e objekteve KStream dhe KTable

Në kapitullin 4 ne diskutuam lidhjen e dy objekteve KStream. Tani na mbetet të mësojmë si të lidhim KTable me KStream. Kjo mund të nevojitet për një arsye të thjeshtë. KStream është një rrjedhë regjistrimesh, ndërsa KTable është një rrjedhë përditësimesh regjistrimesh, por ndonjëherë mund të jetë e nevojshme të shtohet një kontekst shtesë në rrjedhën e regjistrimeve me përditësime nga KTable.

Të marrim të dhënat mbi numrin e transaksioneve në bursë dhe t'i lidhim ato me lajmet e bursës për sektorët përkatës ekonomikë. Këto janë hapat që duhet të ndiqni për ta arritur këtë duke marrë parasysh kodin ekzistues.

  1. Të konvertojmë objektin KTable me të dhënat për numrin e transaksioneve në bursë në KStream, duke zëvendësuar çelësin me një çelës që tregon sektorin ekonomik të këtij simboli aksionesh.
  2. Të krijojmë një objekt KTable që lexon të dhënat nga topic-i me lajmet e bursës. Ky KTable i ri do të kategorizohet sipas sektorëve ekonomikë.
  3. Të lidhim azhurnimet e lajmeve me informacionin mbi numrin e transaksioneve në bursë sipas sektorëve ekonomikë.

Tani le të shohim se si ta zbatojmë këtë plan aksioni.

Konvertimi i KTable në KStream

Për të konvertuar KTable në KStream, duhet të bëni si më poshtë.

  1. Të thërrisni metodën KTable.toStream().
  2. Me thirrjen e metodës KStream.map, zëvendësoni çelësin me emrin e sektorit ekonomik dhe më pas nxirrni nga instanca Windowed objektin TransactionSummary.

Ne do të lidhim këto operacione në një zinxhir siç tregohet më poshtë (kodi mund të gjendet në skedarin src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (listing 5.8).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Duke ne realizojmë operacionin KStream.map, rikthimi i sekcionimit për instancën KStream të kthehet automatikisht kur përdoret në bashkim.

Ne përfunduam procesin e transformimit, tani nevojitet të krijojmë një objekt KTable për të lexuar lajmet e tregut.

Krijimi i KTable për lajmet e tregut

Fatmirësisht, për të krijuar objektin KTable mjafton një varg kodi (ky kod mund të gjendet në skedarin src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (listimi 5.9).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Vlen të përmendet se nuk kërkohet identifikimi i asnjë objekti Serde, pasi përdoren Serde me varg. Gjithashtu, falë përdorimit të enumeratës EARLIEST, tabela plotësohet me regjistrime që në fillim.

Tani mund tĂ« kalojmĂ« nĂ« hapin pĂ«rfundimtar — bashkimin.

Bashkimi i përditësimeve të lajmeve me të dhënat e numrit të transaksioneve

Krijimi i bashkimit nuk paraqet ndonjë vështirësi. Ne do të përdorim një bashkim të majtë për rastin se nuk ka lajme të tregut për përkatësinë sektoriale (kode e nevojshme mund të gjendet në skedarin src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (listimi 5.10).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Ky operatore leftJoin është mjaft i thjeshtë. Në dallim nga bashkimet në kapitullin 4, metoda JoinWindow nuk përdoret, sepse gjatë kryerjes së bashkimit KStream-KTable për çdo çelës në KTable ekziston vetëm një regjistrim. Ky bashkim nuk është i kufizuar në kohë: regjistrimi është ose në KTable, ose nuk ekziston. Përfundimi kryesor: përdorimi i objekteve KTable mundëson pasurimin e KStream me të dhëna referuese që për updates janë më pak të shpeshta.

Tani do të shqyrtojmë një mënyrë më efektive për pasurimin e ngjarjeve nga KStream.

5.3.4. Objektet GlobalKTable

Si e kuptoni, ekzistenca e nevojës për të pasuruar rrjedhat e ngjarjeve ose për të shtuar kontekstin për to. Në kapitullin 4, keni parë lidhjet e dy objekteve KStream, dhe në seksionin e mëparshëm - lidhjen e KStream dhe KTable. Në të gjitha këto raste, është e nevojshme të bëhet ri-sekcionimi i rrjedhës së të dhënave kur çelësat janë të shndërruar në një tip të ri ose vlerë. Ndonjëherë ri-sekcionimi bëhet në mënyrë eksplicite, dhe ndonjëherë Kafka Streams e bën këtë automatikisht. Rii-sekcionimi është i nevojshëm sepse çelësat janë ndryshuar dhe regjistrimet duhet të flasin në seksione të reja, përndryshe lidhja do të jetë e pamundur (kjo u diskutua në kapitullin 4, në pika «Rii-sekcionimi i të dhënave» nënseksioni 4.2.4).

Rii-sekcionimi ka një çmim të tij.

Riorganizimi kërkon shpenzime - shpenzime të mëtejshme për krijimin e temave ndërmjetëse, ruajtjen e të dhënave të dyfishta në një temë tjetër; kjo gjithashtu nënkupton një rritje të vonesës për shkak të shkrimit dhe leximit nga kjo temë. Për më tepër, nëse nevojitet të realizohet lidhja për më shumë se një aspekt ose dimension, duhet të organizohen lidhje në zinxhir, të shfaqen regjistrimet me çelësa të rinj dhe të përsëritet procesi i riorganizimit.

Lidhja me grupe të dhënash më të vogla

Në disa raste, sasia e të dhënave referuese me të cilat është planifikuar lidhja është relativisht e vogël, kështu që kopjet e plota të tyre mund të ndodhen lehtësisht lokalisht në secilin nga node-t. Për këto situata, Kafka Streams ofron klasën GlobalKTable.

Exemplarët e GlobalKTable janë unikë, pasi aplikacioni replikon të gjitha të dhënat në çdo nga nyjet. Dhe, pasi të gjitha të dhënat janë të pranishme në çdo nyje, nuk ka nevojë për të seksionuar rrjedhën e ngjarjeve sipas çelësit të të dhënave referuese për ta bërë atë të disponueshëm për të gjitha seksionet. Me objektet GlobalKTable, gjithashtu mund të kryhen bashkime pa çelës. Le të kthehemi në një nga shembujt e mëparshëm për të demonstruar këtë mundësi.

Bashkimi i objekteve KStream me objektet GlobalKTable

Në nënseksionin 5.3.2, ne kryem agregimin e dritareve të transaksioneve të bursës sipas blerësve. Rezultatet e këtij agregimi dukeshin përafërsisht kështu:

{customerId='074-09-3705', stockTicker='GUTM'}, 17
{customerId='037-34-5184', stockTicker='CORK'}, 16

Megjithëse këto rezultate përputhen me qëllimin e përcaktuar, do të ishte më e dobishme nëse gjithashtu do të shfaqej emri i klientit dhe emri i plotë i kompanisë. Për të shtuar emrin e blerësit dhe emrin e kompanisë, mund të kryhen bashkime të zakonshme, por do të nevojitet të bëhen dy përkthime çelesh dhe seksionimi i ri. Me GlobalKTable, mund të shmangen kostot për këto operacione.

Për këtë do ta përdorim objektin countStream nga lista 5.11 (kodin përkatës e gjeni në skedarin src/main/java/bbejeck/chapter_5/GlobalKTableExample.java), duke e lidhur atë me dy objekta GlobalKTable.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Kjo është diskutuar më parë, prandaj nuk do të përsëris. Por dua të theksoj se kodi në funksionin toStream().map është abstrahuar në një objekt-funksion për lehtësinë e lexueshmërisë, në vend të një shprehjeje lambda të vendosur.

Hapi tjetër është shpallja e dy instancave GlobalKTable (kodin e dhënë e gjeni në skedarin src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (lista 5.12).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»

Vini re se emrat e temave përshkruhen me ndihmën e llojeve enumeruese.

Tani që kemi përgatitur të gjithë komponentët, mbetet të shkruajmë kodin për lidhjen (të cilin e gjeni në skedarin src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (lista 5.13).

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Megjithatë, në këtë kod ka dy lidhje, të cilat janë organizuar në formë zinxhiri, pasi asnjë nga rezultatet e tyre nuk përdoret veçmas. Rezultatet shfaqen në fund të gjithë operacionit.

Kur të ekzekutoni operacionin e lidhjes të mësipërm, do të merrni rezultate të këtij lloji:

{customer='Barney, Smith' company="Exxon", transactions= 17}

Përmbajtja nuk ka ndryshuar, por këto rezultate duken më të qarta.

Nëse përfshijmë kapitullin 4, tashmë keni parë disa lloje lidhjesh në veprim. Ato janë listuar në tabela 5.2. Kjo tabelë reflekton mundësitë e lidhjes që janë të rëndësishme për versionin 1.0.0 të Kafka Streams; në botimet e ardhshme, ndoshta do të ketë disa ndryshime.

Libri «Kafka Streams në Veprim. Aplikacione dhe Mikroservis për Punën në Kohë Reale»
Në përfundim, dua të theksoj thelbësoren: mund të lidhni rrjedhat e ngjarjeve (KStream) dhe rrjedhat e përditësimeve (KTable) duke përdorur gjendjen lokale. Për më tepër, nëse madhësia e të dhënave ndihmëse nuk është tepër e madhe, mund të përdorni objektin GlobalKTable. GlobalKTable replikon të gjitha seksionet në secilën nga nodet e aplikacionit Kafka Streams, duke siguruar në këtë mënyrë disponueshmërinë e të gjitha të dhënave pavarësisht nga cilit seksion i përket çelësi.

Më pas do të shohim një mundësi të Kafka Streams, e cila lejon që të vëzhgoni ndryshimet në gjendje pa konsumuar të dhëna nga topic-i Kafka.

5.3.5. Gjendja e disponueshme për kërkesa

Ne kemi kryer disa operacione që përfshijnë gjendjen dhe gjithmonë kemi shfaqur rezultatet në konsolë (për qëllime zhvillimi) ose i kemi regjistruar ato në një temë (për qëllime industriale). Kur regjistrojmë rezultatet në një temë, duhet të përdorim konsumatorin Kafka për t'i parë ato.

Leximi i tĂ« dhĂ«nave nga kĂ«to tema mund tĂ« konsiderohet si njĂ« lloj vizionesh tĂ« materializuara (materialized views). PĂ«r ne, mund tĂ« pĂ«rdorim pĂ«rkufizimin e vizionit tĂ« materializuar nga «Wikipedia»: « njĂ« objekt fizik i bazĂ«s sĂ« tĂ« dhĂ«nave qĂ« pĂ«rmban rezultatet e ekzekutimit tĂ« njĂ« kĂ«rkese. PĂ«r shembull, ajo mund tĂ« jetĂ« njĂ« kopje lokale e tĂ« dhĂ«nave tĂ« largĂ«ta, ose njĂ« nĂ«ngrup rreshtash dhe/ose kolonash tĂ« njĂ« tabele ose rezultateve tĂ« bashkimit, apo njĂ« tabelĂ« e pĂ«rmbledhur e marrĂ« pĂ«rmes agregimit» (https://en.wikipedia.org/wiki/Materialized_view).

Kafka Streams gjithashtu lejon ekzekutimin e pyetjeve interaktive ndaj depozitave tĂ« gjendjes, duke ofruar mundĂ«sinĂ« e leximit tĂ« drejtpĂ«rdrejtĂ« tĂ« kĂ«tyre pamjeve tĂ« materializuara. ËshtĂ« e rĂ«ndĂ«sishme tĂ« theksohet se pyetja ndaj depozitave tĂ« gjendjes ka karakterin e njĂ« operacioni 'vetĂ«m pĂ«r lexim'. KĂ«shtu, mund tĂ« mos shqetĂ«soheni pĂ«r aksidentalisht tĂ« bĂ«ni gjendjen tĂ« papajtueshme gjatĂ« pĂ«rpunimit tĂ« tĂ« dhĂ«nave nga aplikacioni.

MundĂ«sia e pyetjeve tĂ« drejtpĂ«rdrejta ndaj depozitave tĂ« gjendjes ka rĂ«ndĂ«si tĂ« madhe. Kjo do tĂ« thotĂ« se mund tĂ« krijoni aplikacione — tabela informuese pa e nevojitur mĂ« parĂ« marrjen e tĂ« dhĂ«nave nga konsumatori i Kafka. Ajo rrit edhe efikasitetin e aplikacionit, pĂ«r shkak se nuk Ă«shtĂ« e nevojshme tĂ« shkruhen pĂ«rsĂ«ri tĂ« dhĂ«nat:

  • pĂ«r shkak tĂ« lokalitetit tĂ« tĂ« dhĂ«nave, ato mund tĂ« aksesohen shpejt;
  • pĂ«rjashtohet dublimi i tĂ« dhĂ«nave, pasi ato nuk shkruhen nĂ« njĂ« depo tĂ« jashtme.

Gjëja më e rëndësishme që do doja që ju ta mbani mend është: mund të bëni kërkesa direkt në magazinat e gjendjes nga aplikacioni. Asnjëherë nuk mund të nënvlerësoni mundësitë që kjo ju ofron. Në vend që të konsumoni të dhëna nga Kafka dhe të ruani regjistrat në një bazë të dhënash për aplikacionin, mund të bëni kërkesa në magazinat e gjendjes me të njëjtin rezultat. Kërkesat e drejtpërdrejta në magazinat e gjendjes nënkuptojnë më pak kod (pa nevojën për një konsumator) dhe më pak softuer (pa nevojën për një tabelë të bazës së të dhënave për ruajtjen e rezultateve).

Ne kemi mbuluar njĂ« sasi tĂ« konsiderueshme informacioni nĂ« kĂ«tĂ« kapitull, prandaj pĂ«r njĂ« kohĂ« do ta ndĂ«rpresim diskutimin tonĂ« mbi kĂ«rkesat interaktive nĂ« magazinat e gjendjes. Por mos u shqetĂ«soni: nĂ« kapitullin 9 do tĂ« krijojmĂ« njĂ« aplikacion tĂ« thjeshtĂ« — njĂ« panel informacioni me kĂ«rkesa interaktive. PĂ«r tĂ« demonstruar kĂ«rkesat interaktive dhe mundĂ«sitĂ« e tyre pĂ«r t'u shtuar nĂ« aplikacionet Kafka Streams, do tĂ« pĂ«rdoren disa nga shembujt e kĂ«tij kapitulli dhe tĂ« kapitujve tĂ« mĂ«parshĂ«m.

Curriculum Vitae

  • Objektet KStream pĂ«rfaqĂ«sojnĂ« rrjedha eventi, tĂ« ngjashme me futjet nĂ« njĂ« bazĂ« tĂ« dhĂ«nash. Objektet KTable pĂ«rfaqĂ«sojnĂ« rrjedha azhurnimesh, janĂ« mĂ« shumĂ« tĂ« ngjashme me azhurnimet nĂ« njĂ« bazĂ« tĂ« dhĂ«nash. MadhĂ«sia e objektit KTable nuk rritet, tĂ« dhĂ«nat e vjetra zĂ«vendĂ«sohen nga tĂ« rejat.
  • Objektet KTable janĂ« tĂ« nevojshme pĂ«r operacionet e agregimit.
  • Me ndihmĂ«n e operacioneve tĂ« dritareve, mund tĂ« ndahen tĂ« dhĂ«nat e agreguara nĂ« kosha kohorĂ«.
  • Me ndihmĂ«n e objekteve GlobalKTable, mund tĂ« aksesohen tĂ« dhĂ«nat referuese nĂ« çdo pikĂ« tĂ« aplikacionit, pavarĂ«sisht ndarjes sipas sekcioneve.
  • JanĂ« tĂ« mundshme lidhjet mes objekteve KStream, KTable dhe GlobalKTable.

Derët e fundit, ne kemi përqendruar vëmendjen tonë në krijimin e aplikacioneve Kafka Streams duke përdorur DSL të nivelit të lartë KStream. Megjithëse qasja e nivelit të lartë lejon krijimin e programeve të pastra dhe të përmbledhura, përdorimi i saj përfaqëson një kompromis të caktuar. Puna me DSL KStream do të thotë se po e rrisim përmbledhjen e kodit duke ulur shkallën e kontrollit. Në kapitullin e ardhshëm do të shqyrtojmë API-në e nivelit të ulët të nyjave përpunuese dhe do të provojmë kompromise të tjera. Programet do të bëhen më të gjata se sa ishin deri tani, por kështu do të kemi mundësinë të krijojmë pothuajse çdo nyje përpunuese që na nevojitet.

→ MĂ« shumĂ« informacione pĂ«r librin mund tĂ« gjenden nĂ« faqen e botuesit

→ PĂ«r pĂ«rdoruesit e Habra, zbritje 25% me kodin — Kafka Streams

→ Pas pagesĂ«s sĂ« versionit tĂ« printuar tĂ« librit, dĂ«rgohet njĂ« libĂ«r elektronik nĂ« e-mail.

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