Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks» Tere, Habržiteli! See raamat sobib igasugustele arendajatele, kes soovivad mõista voogude töötlemist. Jagatud programmeerimise mõistmine aitab paremini uurida Kafka ja Kafka Streams'i. Oleks hea, kui teada ka Kafka raamistikku, kuid see pole kohustuslik: ma räägin teile kõik, mis vajalik. Kogenud Kafka arendajad, nagu ka algajad, saavad tänu sellele raamatule omandada oskusi huvitavate rakenduste loomisel voogude töötlemise kaudu Kafka Streams'i raamatukoguga. Keskmise ja kõrge tasemega Java-arendajad, kes on juba tuttavad selliste mõistetega nagu serialiseerimine, õpivad rakenduse loomise oskusi rakendama Kafka Streams'iga. Raamatu lähtekood on kirjutatud Java 8-s ja kasutab oluliselt Java 8 lambda-väljendite süntaksit, seega on kasulik osata töötada lambda-funktsioonidega (isegi kui neid rakendatakse teises programmeerimiskeeles).

Lõik. 5.3. Agregeerimine ja aknaoperatsioonid

Selles jaotises liigume edasi Kafka Streams'i kõige lootustandvamaid osi uurima. Seni oleme uurinud järgmisi aspekte Kafka Streams'ist:

  • töötlemise topoloogia loomine;
  • olekute kasutamine voograkendustes;
  • voogandmete ühenduste täitmine;
  • erinevused sündmuste voogude (KStream) ja värskenduste voogude (KTable) vahel.

Järgmistes näidetes toome need kõik elemendid kokku. Tutvute ka aknaoperatsioonidega — veel ühe suurepärase voograkenduste funktsionaalsusega. Meie esimene näide käsitleb lihtsat agregatsiooni.

5.3.1. Aktsiate müügimahtude agregatsioon tööstusharude kaupa

Agregatsioon ja rühmitamine on hädavajalikud tööriistad voogandmetega töötamisel. Üksikute kirje uurimine reaalajas osutub sageli ebapiisavaks. Andmetest täiendava teabe saamiseks on vajalik nende rühmitamine ja kombineerimine.

Selles näites peate proovima päevakaubanduse kaupleja kostüümi, kellel on vaja jälgida aktsiate müügimahte mitmest tööstusharust. Eelkõige huvitavad teid viis ettevõtet, millel on igas tööstusharus kõige suuremad aktsiate müügimahud.

Selle sarnase aggregeerimise jaoks on vajalikud järgmised sammud andmete soovitud vormi tõlkimiseks (kui rääkida laiemalt).

  1. Loo allikas, mis põhineb teemal, mis avaldab töötlemata teavet aktsiakauplemise kohta. Peame kujundama StockTransaction objekti ShareVolume objektiks. Asi on selles, et StockTransaction objekt sisaldab müügi metaandmeid, samas kui me vajame vaid andmeid müüdud aktsiate arvu kohta.
  2. Grupita ShareVolume andmed aktsia sümbolite kaupa. Pärast sümbolite järgi grupitamist saab need andmed kokku tõmmata aktsiate müügimahu vahepealseteks summadeks. Tuleb märkida, et meetod KStream.groupBy tagastab KGroupedStream tüüpi eksemplari. KTable eksemplari saamiseks saab seejärel kutsuda meetodi KGroupedStream.reduce.

Mis on KGroupedStream liides?

KStream.groupBy ja KStream.groupByKey meetodid tagastavad KGroupedStream'i eksemplari. KGroupedStream on vahepealne esitus sündmuste voos pärast gruppimist võtmete alusel. Seda ei ole mõeldud otseseks töötlemiseks. Selle asemel kasutatakse KGroupedStream'i kogumise operatsioonides, mille tulemus on alati KTable. Kuna kogumise operatsioonide tulemus on KTable ja neis kasutatakse olekuhoidlat, ei pruugi kõik uuendused edastada konveieri kaudu.

KTable.groupBy meetod tagastab sarnase KGroupedTable'i — vahepealse esitus voos, mis on uuendatud võtme alusel.

Teeme väikese pausi ja vaatame joonist 5.9, mis näitab, mida oleme saavutanud. See topoloogia peaks olema teile juba tuttav.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Vaatame nüüd koodi selle topoloogia kohta (leiate selle failist src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (loetelu 5.2).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Esitatud kood eristub lühiduse ja mitme tegevuse samaaegse käivitamise suure mahuga. Esimeses parameetris meetodil builder.stream võite märgata midagi uut: loetletud tüübi AutoOffsetReset.EARLIEST väärtus (olemas on ka LATEST), mille määrab meetod Consumed.withOffsetResetPolicy. Selle loetletud tüübi abil saab määrata nihke lähtestamise strateegiat iga KStream või KTable jaoks, mis on prioriteediga kõrgem kui konfigureeritud lähtestamise parameeter.

GroupByKey ja GroupBy

KStream liideses on kaks meetodit kirjegruppideks: GroupByKey ja GroupBy. Mõlemad tagastavad KGroupedTable'i, mistõttu võib tekkida loogiline küsimus: milles seisneb nende erinevus ja millal kasutada ühte või teist?

Meetodit GroupByKey kasutatakse, kui KStreami võtmed on juba ei ole tühjad. Veelgi enam, lipp „nõuab ümberjaotamist” pole kunagi seadistatud.

Meetod GroupBy eeldab, et olete võtmed grupiseerimiseks muutnud, mistõttu on ümberjaotamise lipp seadistatud väärtusele true. Meetodi GroupBy järel toimunud liitmine, agregatsioon jne toob automaatselt kaasa ümberjaotamise.
Kokkuvõte: tuleks iga võimaluse korral kasutada GroupByKey, mitte GroupBy.

Mida teevad meetodid mapValues ja groupBy — see on selge, nii et vaatame meetodit sum() (seda võib leida failist src/main/java/bbejeck/model/ShareVolume.java) (loetelu 5.3).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Meetod ShareVolume.sum tagastab vahepealse aktsiate müügi mahu summa ning kogu arvutuste ahela tulemus on objekt KTable. Nüüd mõistate, milline roll on KTable'il. Kui objekti ShareVolume andmed saabuvad, salvestatakse vastavas KTable objekti viimased kehtivad uuendused. Oluline on meeles pidada, et kõik uuendused kajastuvad eelnevas shareVolumeKTable'is, kuid mitte kõik edastatakse edasi.

Seejärel kasutame selle KTable abil agregatsiooni (müüdud aktsiate arvu järgi), et saada viis ettevõtet, millel on suurim aktsiate müügimahu iga tööstusharu puhul. Meie tegevus on sarnane esimese agregatsiooni tegevusele.

  1. Tehke veel üks groupBy operatsioon, et grupeerida eraldi ShareVolume objekte tööstusharude kaupa.
  2. Alustage ShareVolume objektide kokkuvõtmiseks. Seekord on agregatsiooni objektiks prioriteetsete fikseeritud suurusega järjekord. Sellises fikseeritud suurusega järjekorras talletatakse ainult viis ettevõtet, kellel on kõige rohkem müüdud aktsiaid.
  3. Kuvage eelnevas punktis olevad järjekorrad sõnes ja tagastage viis kõige rohkem müüdud aktsiat tööstusharude kaupa.
  4. Salvestage tulemused sõna kujul teemasse.

Kujutisel 5.10 on kujutatud andmeedastuse topoloogia graafikut. Nagu näete, on teise töötlemise ring piisavalt lihtne.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Nüüd, olles selgelt mõistnud selle teise töötlemise ringi struktuuri, saame tutvuda selle lähtefailiga (leiate selle failist src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (loend 5.4).

Käesolevas initsialiseerijas on muutuja fixedQueue. See on kasutajaobjekt – adapter java.util.TreeSet jaoks, mida kasutatakse N suurima tulemuse jälgimiseks müüdud aktsiate arvu järgi.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Olete juba tutvunud groupBy ja mapValues väljakutsetega, seega ei peatume nende kallal (kutsume välja meetodi KTable.toStream, kuna meetodit KTable.print peetakse aegunuks). Kuid te ei ole veel näinud KTable versiooni meetodist aggregate(), seega kulutame natuke aega selle arutamiseks.

Kuidas te mäletate, eristab KTable seda, et sama võtmega salvestusi käsitletakse uuendustena. KTable asendab vana salvestuse uuega. Agregeerimine toimub sarnaselt: kogutakse viimaseid salvestusi sama võtmega. Kui salvestus saabub, lisatakse see klassi FixedSizePriorityQueue eksemplari läbi summuti (teine parameeter meetodi aggregate kutses), kuid kui sama võtmega salvestus juba eksisteerib, siis eemaldatakse vana salvestus läbi lahutaja (kolmas parameeter meetodi aggregate kutses).

See tähendab, et meie agregaat, FixedSizePriorityQueue, ei koguge kõiki väärtusi ühte võtmesse, vaid hoiab N kõige müüdava aktsia liiki läbiva summa. Igas sisenevas kirjes on kogu müüdud aktsiate arv. KTable annab teile teavet selle kohta, milliste ettevõtete aktsiaid praegu kõige rohkem müüakse, ja iga värskenduse libisev agregatsioon pole vajalik.

Oleme õppinud tegema kahte olulist asja:

  • gruppima KTable väärtusi nende ühiselt võtmega;
  • tegema nende grupeeritud väärtuste üle kasulikke operatsioone, nagu kokkuliitmine ja agregatsioon.

Nende operatsioonide teostamise oskus on oluline, et mõista andmete tähendust, mis liikuvad läbi Kafka Streams rakenduse ning välja selgitada, millist teavet need endaga kaasa toovad.

Oleme samuti kokku viinud mõned varem selle raamatu käigus arutlusele tulnud peamised mõisted. 4. peatükis rääkisime, kui oluline on voogedastusrakenduse jaoks usaldusväärne, lokaalne seisund. Esimene näide sellest peatükist näitas, miks lokaalne seisund on nii oluline — see võimaldab jälgida, millist teavet olete juba näinud. Lokaalne juurdepääs aitab vältida võrgu viivitusi, mistõttu rakendus muutub efektiivsemaks ja vigu taluvamaks.

Kogu sõlme või kogumise operatsiooni teostamisel on vaja näidata seisundipoe nime. Sõlme ja kogumise operatsioonid tagastavad KTable'i eksemplari, ning KTable kasutab seisundipood, et asendada vanad tulemused uutega. Nagu nägite, ei saadeta kõik värskendused edasi toru kaudu, ja see on oluline, kuna kogumise operatsioonid on mõeldud lõppinformatsiooni saamiseks. Kui lokaalset seisundit ei kasutata, saadab KTable edasi kõik kogumise ja sõlmimise tulemused.

Järgmisena vaatame, kuidas teostada selliseid tehinguid nagu aggregeerimine konkreetse ajaintervali piires — nn aknaoperatsioone (windowing operations).

5.3.2. Aknaoperatsioonid

Eelmises jaotises tutvusime „liuguvate“ kokkusurve ja aggregeerimisega. Rakendus teostas pidevat müügimahtude kokkusurvet ning seejärel aggregeeris viit kõige rohkem börsil müüdud aktsiat.

Mõnikord on sellised pidevad aggregeerimine ja tulemuste kokkusurve vajalikud. Teistel juhtudel tuleb operatsioone teostada ainult määratud ajavahemiku jooksul. Näiteks arvutada, kui palju börsitehinguid on toimunud konkreetse ettevõtte aktsiatega viimase 10 minuti jooksul. Või kui palju kasutajaid klikkis uuele reklaambännerile viimase 15 minuti jooksul. Rakendus võib selliseid operatsioone mitmeid kordi teostada, kuid ainult määratud ajavahemike (ajakenade) tulemuste jaoks.

Börsiterminite arvestamine ostjapõhiselt

Järgnevas näites keskendume börsitehingute jälgimisele mitme kaupleja — kas suurte organisatsioonide või nutikate üksikfinantseerijate — osas.

On kaks võimalikku põhjust selliseks jälgimiseks. Üks neist on vajadus teada, mida turu juhid ostavad/müüdavad. Kui need suured mängijad ja kogenud investorid näevad avanemas võimalusi, on mõistlik järgida nende strateegiaid. Teine põhjus on soov märgata mingeid võimalikud ebaseaduslikke tehingute märke, mis kasutavad siseinfot. Selleks peate analüüsima suurte müügipuhangute korrelatsiooni oluliste pressiteadete jaoks.

Selline jälgimine koosneb järgmistest etappidest:

  • looda voog, et lugeda teemast stock-transactions;
  • sisendandmete rühmitamine ostja ID ja aktsia sümboli järgi. meetodi groupBy kutsumine tagastab KGroupedStream klassi eksemplari;
  • KGroupedStream.windowedBy meetodi abil tagastatakse ajaliniku poolest piiratud andmevoog, mis võimaldab akna agregatsiooni teostamist. Aknatüübi järgi tagastatakse kas TimeWindowedKStream või SessionWindowedKStream;
  • tehingute arvestamine agregatsioonitegevuse jaoks. Aken andmevoos määrab, kas konkreetne kirje arvestatakse;
  • tulemuste salvestamine teemasse või nende väljastamine konsoolis arendamise ajal.

Selle rakenduse topoloogia on lihtne, kuid oleks kasulik näha selle visuaalit. Vaadake joonist 5.11.

Järgnevalt vaatleme aknaoperatsioonide funktsionaalsust ja vastavat koodi.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»

Aknatüübid

Kafka Streams'is on kolm aknatüüpi:

  • seansilised;
  • «kuvurdavad» (tumbling);
  • libisevad/«hüppavad» (sliding/hopping).

Valik sõltub äri nõudmistest. „Kohati hüplevaid“ ja „keerlevaid“ aknaid piiratakse ajaliselt, samas kui sessiooni piirmäärad on seotud kasutaja tegevustega — sessiooni (de) kestus määratakse ainult selle järgi, kui aktiivselt kasutaja käitub. Peaasi on meeles pidada, et kõik aknatüübid põhinevad kirje kuupäeva/kellaaja märkidel, mitte süsteemi ajast.

Järgmisena rakendame meie topoloogia iga aknatüübiga. Täielik kood tuuakse esimeses näites, muu aknatüüpide puhul ei muutu midagi, välja arvatud akna operatsiooni tüüp.

Sessiooni aknad

Sessiooni aknad erinevad oluliselt teistest aknatüüpidest. Need ei piirdu mitte niivõrd ajaga, kuivõrd kasutaja aktiivsusega (või selle tegevuse aktiivsusega, mida soovite jälgida). Sessiooni aknad eristuvad tegevusetuse perioodide järgi.

Joon 5.12 illustreerib seanssiakna mõistet. Väiksem seanss sulandub tema vasakul asuva seansiga. Paremal asuv seanss jääb aga eraldi, kuna sellele eelneb pikk vaikuseperiood. Seanssiaken põhineb kasutajate tegevusel, kuid kasutab kuupäeva/kellaaja tunnuseid kirjetest, et määrata, millise seansiga seondub üks või teine kirje.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»

Seanssiakende kasutamine börsitehingute jälgimiseks

Kasutame seanssiaken, et koguda teavet börsitehingute kohta. Seanssiakna rakendus on näidatud loendis 5.5 (mida saab leida failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKTableJoinExample.java).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Enamik selle topoloogia operatsioone on teile juba tuttavad, nii et pole mõtet neid siinkohal uuesti arutada. Siiski on siin mõned uued elemendid, mida me kohe arutame.

Iga groupBy operatsiooni puhul tehakse tavaliselt mingi agregatsiooni operatsioon (agregatsioon, kokkuliitmine või loendamine). Saate teha kas kumulatiivset agregatsiooni jooksva summa alusel või aknaalusel agregatsiooni, kus arvestatakse kirjeid määratud ajavahemiku jooksul.

Kood, mis on esitatud loendis 5.5, arvutab sessioonide aknas toimuva tehingute arvu. Joonisel 5.13 analüüsitakse neid toiminguid samm-sammult.

Kasutades käsku windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)), loome seansiaugu, mille inaktiivsuse intervall on 20 sekundit ja säilitamisinterval 15 minutit. 20-sekundiline inaktiivsuse intervall tähendab, et rakendus hõlmab kõiki kirjeid, mis saabuvad 20 sekundi jooksul pärast praeguse seansi lõppu või algust, aktiivsesse seanssi.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Seejärel täpsustame, millist agregatsioonitegevust tuleb seansiaugus teostada — antud juhul count. Kui sisenev kirje ületab inaktiivsuse intervalli (millisel pool ka ei jääks ajatempli), loob rakendus uue seansi. Säilitamisinterval tähendab seansi hoidmist kindla aja jooksul ja võimaldab hilinenud andmeid, mis ületavad seansi inaktiivsuse perioodi, kuid võivad siiski liituda. Lisaks sellele vastavad uue seansi algus ja lõpp, mis on saadud liitmise tulemusena, kõige varasemale ja kõige hilisemale ajatempli.

Vaatame mõningaid kirjeid meetodist count, et näha, kuidas sessioonid töötavad (vt tabel 5.1).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Kirjeid vastu võtades otsime juba olemasolevaid sessioone sama võtme, lõpuajamärgi all, mis on väiksem praegusest kuupäevast/aeg-märgist — inaktiivsuse intervall ja algusaeg märgiga suurem kui praegune kuupäev/aeg + inaktiivsuse intervall. Selle arvesse võttes sulanduvad tabelis 5.1 neli kirjet ühte sessiooni järgmiselt.

1. Esimesena saabub kiri 1, nii et algusaeg on sama kui lõpu aeg ja see on 00:00:00.

2. Järgmisena saabub kiri 2, ning otsime sessioone, mis lõppevad mitte enne 23:59:55 ja algavad mitte hiljem kui 00:00:35. Leiame kirje 1 ja ühendame sessioonid 1 ja 2. Võtame sessiooni 1 algusaja (varasema) ja sessiooni 2 lõpuaja (hilisema), nii et meie uus sessioon algab 00:00:00 ja lõpeb 00:00:15.

3. Saabub kiri 3, otsime sessioone ajavahemikus 00:00:30 kuni 00:01:10 ja ei leia ühtegi. Lisame teise sessiooni võtme 123-345-654,FFBE jaoks, mis algab ja lõppeb 00:00:50.

4. Saabub sissekanne 4, ja me otsime seansse ajavahemikus 23:59:45 kuni 00:00:25. Seekord leitud mõlemad seansid — 1 ja 2. Kõik kolm seanssi ühendatakse üheks, algusajaga 00:00:00 ja lõppaeg 00:00:15.

Selles osas räägitust tasub meelde jätta järgmised olulised nüansid:

  • seansid — ei ole fikseeritud suurusega aknad. Seansi kestus määratakse antud ajavahemiku jooksul toimuva aktiivsuse järgi;
  • kuupäeva/aegade märgid andmetes määravad, kas sündmus langeb olemasolevasse seanssi või passiivsuse perioodi.

Järgmine, millest me räägime, on järgmine akende tüüp — "rulluvad" aknad.

"Rulluvad" aknad

"Rulluvad" (tumbling) aknad haaravad sündmusi, mis langevad kindlasse ajavahemikku. Kujutage ette, et peate haarama kõik börsitransaktsioonid mingi ettevõtte kohta iga 20 sekundi tagant, nii et kogute kõik sündmused selle ajavahemiku jooksul. Pärast 20-sekundilist vahemikku akna "rullub" ja liigub uuele 20-sekundilisele jälgimisintervallile. Kujutis 5.14 illustreerib seda olukorda.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Nagu näete, on kõik viimase 20 sekundi jooksul toimunud sündmused aknas. Selle ajavahemiku lõppedes luuakse uusaken.

Loendis 5.6 on esitatud kood, mis demonstreerib „hüppavate” akende kasutamist börsietehingute igakümne sekundi jooksul jäädvustamiseks (seda leiab failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Selle väikese muudatuse tõttu meetodi TimeWindows.of kutsumises saab kasutada „hüppavat” akent. Antud näites ei ole meetodi until() kutsumist, mistõttu kasutatakse vaikimisi salvestusintervalle, mis on 24 tundi.

Lõpuks on aeg liikuda viimase aknatüübi juurde — „hüppavad” aknad.

Libisevad („hüppavad”) aknad

Libisevad/„hüppavad” aknad on sarnased „hüppavatele”, kuid väikese erinevusega. Libisevad aknad ei oota ajavahemiku lõppu enne uue akna loomist hiljutiste sündmuste töötlemiseks. Nad käivitavad uued arvutused pärast ooteperioodi, mis on lühem kui akna kestus.

Erinevuste illustreerimiseks «veerlevate» ja «hüppavate» akende vahel naaseme näite juurde börsitehingute arvestamisest. Meie eesmärk on endiselt tehingute arvu arvestamine, kuid me ei soovi oodata kogu ajavahemikku enne arvesti uuendamist. Selle asemel uuendame arvestit lühemate ajavahemike järel. Näiteks arvestame tehingute arvu endiselt iga 20 sekundi järel, kuid uuendame arvestit iga 5 sekundi järel, nagu on kujutatud joonisel 5.15. Sellega saame kolme tulemuse akna, millel on kattuvad andmed.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Loendis 5.7 on toodud kood liugakende määramiseks (selle leiab failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
«Veerlev» aken saab muuta «hüppavaks», lisades meetodi advanceBy() kutsumise. Antud näites on säilitamise periood 15 minutit.

Selles jaos nägite, kuidas piirata kokkupanemise tulemusi ajavahemikega. Eriti sooviksin, et mälestaksite sellest jaost kolme asja:

  • sessioonide suuruse piirang ei sõltu mitte ajavahemikust, vaid kasutajate aktiivsusest;
  • "veerevad" aknad annavad ülevaate sündmustest määratud ajavahemikus;
  • "hüppavate" akende töö kestus on fikseeritud, kuid need uuendatakse sageli ja võivad sisaldada kõigis akendes kattuvaid kirjeid.

Järgmises osas saame teada, kuidas muuta KTable taas KStreamiks ühendamiseks.

5.3.3. KStreami ja KTable'i ühendamine

4. peatükis arutasime kahe KStreami objekti ühendamist. Nüüd peame õppima, kuidas KTable'i ja KStreami ühendamine toimub. Seda võib olla vajalik teha järgmistel lihtsatel põhjustel. KStream on kirjete voog, samas kui KTable on kirjete uuenduste voog, kuid mõnikord võib osutuda vajalikuks lisada kirjete voole täiendav kontekst KTable'i uuenduste kaudu.

Võtame andmed börsitehingute arvu kohta ja ühendame need börsiuudistega vastavate tööstusharude kohta. Siin on, mida tuleb teha, et seda olemasoleva koodi arvesse võttes saavutada.

  1. Muutame KTable'i, kus on andmed börsitehingute arvu kohta, KStream'iks, asendades seejärel võtme põhitegevuse valdkonna võtmega, mis vastab antud aktsiasümbolile.
  2. Loome KTable'i objekti, mis loeb andmeid börsiuudiste teemast. See uus KTable kategoriseeritakse tööstusharude järgi.
  3. Seome uudiste uuendused teabe kogustega börsitehingute alusel tööstusharude kohta.

Nüüd vaatame, kuidas ellu viia seda tegevuskava.

KTable'i konverteerimine KStream'iks

KTable'i konverteerimiseks KStream'iks tuleb teha järgmist.

  1. Käivitage meetod KTable.toStream().
  2. Kasutades meetodi KStream.map kutsumist, asendage võtta tööstusharu nimega ja seejärel võtke Windowed objekti TransactionSummary eksemplar.

Seome need operatsioonid ahelana järgmiselt (kood on saadaval failis src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (loend 5.8).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Kuna täidame operatsiooni KStream.map, toimub KStream'i tagastatud eksemplari põhiülesanne automaatselt, kui seda kasutatakse ühenduses.

Oleme lõpetanud konverteerimisprotsessi, nüüd peame looma KTable objekti börsitarvete lugemiseks.

KTable'i loomine börsitarvete jaoks

Õnneks on KTable objekti loomiseks piisav vaid üks koodirida (see kood on saadaval failis src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (loend 5.9).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
On oluline märkida, et Serde objekte ei pea määrama, kuna seadistustes kasutatakse stringiserde. Samuti tänu ENUM'i EARLIEST kasutamisele täidetakse tabel algsete kantidega.

Nüüd saame liikuda viimase sammu juurde — ühendamine.

Uudiste uuenduste ühendamine tehingute arvu andmetega

Ühenduse loomine ei sea mingeid keerukusi. Kasutame vasakut ühendust, et juhul kui vastavas tööstusharus ei ole börsitarbeid (vajalik kood on saadaval failis src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (loend 5.10).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
See operaator leftJoin on piisavalt lihtne. Erinevalt peatükist 4, ei kasutata Meetodi JoinWindow puhul, kuna KStream-KTable ühendamise käigus on KTable iga võtme jaoks ainult üks kirje. See ühendamine ei ole ajaliselt piiratud: kirje on KTable'is olemas või puudub. Peamine järeldus: KTable objektide abil saab KStream'i rikastada harvemini uuendatavate viidatud andmetega.

Ja nüüd vaatame tõhusamat viisi KStream'i sündmuste rikastamiseks.

5.3.4. GlobalKTable objektid

Nagu mõistsite, on vajalik sündmustevoolude rikastamine või nendele konteksti lisamine. 4. peatükis nägite kahte KStream'i objekti ühendamist ning eelnevas osas KStream'i ja KTable'i ühendamist. Kõigis nendes juhtumites on vajalik andmevoogude uuesti sektsioonideks jagamine, et kaardistada võtmed uue tüübi või väärtuse juurde. Mõnikord toimub uuesti sektsioonideks jagamine selgelt, teinekord teeb Kafka Streams seda automaatselt. Uuesti sektsioonideks jagamine on vajalik, kuna võtmed on muutunud ja kirjed peavad olema uutes sektsioonides, vastasel juhul ei ole ühendamine võimalik (sellest rääkisime 4. peatükis, lõigus „Andmete uuesti sektsioonideks jagamine“ alal 4.2.4).

Uuesti sektsioonideks jagamine on kulukas

Uuesti partitsioneerimine nõuab kulutusi — lisakulusid ressursside suurendamiseks vahepealsete teemade loomisel, dubleeritud andmete säilitamisel veel ühes teemasse; see toob kaasa ka viivituse suurenemise, kuna andmeid salvestatakse ja loetakse sellest teemast. Lisaks, kui on vajalik teostada ühendusi rohkem kui ühe aspekti või mõõtme osas, tuleb korraldada ühendused ahelana, kuvada kirjed uute võtmetega ja uuesti käivitada partitsioneerimise protsess.

Ühendus väiksemate andmehulkadega

Mõnel juhul on viidatud andmehulkade maht suhteliselt väike, nii et nende täiskopeerimist on võimalik iga sõlme kohapeal hoida. Selliste olukordade jaoks on Kafka Streams'is ette nähtud klass GlobalKTable.

GlobalKTable'i eksemplarid on unikaalsed, kuna rakendus replitseerib kõik andmed igas sõlmes. Kuna igas sõlmes on kõik andmed, ei ole vaja sündmuste voogu jagada andmeviidendi võtme järgi, et see oleks kõigile osakondadele kergesti kättesaadav. GlobalKTable'i objektide abil on võimalik teostada ka võtmeta ühendusi. Vaatame tagasi ühele varasemale näitele, et demonstreerida seda võimalust.

KStream'i objektide ühendamine GlobalKTable'i objektidega

Alamjaotises 5.3.2 tegime akna agregatsiooni börsitehingute põhjal ostjate lõikes. Selle agregatsiooni tulemused nägid välja umbes järgmised:

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

Kuigi need tulemused vastasid seatud eesmärgile, oleks mugavam, kui kliendi nimi ja ettevõtte täisnimi oleksid ka väljundis. Ostja nime ja ettevõtte nime lisamiseks saab teostada tavalisi ühendusi, kuid lahtise võtme ja uuesti jaotamise tegemine on vajalik. GlobalKTable'i abil on võimalik selliste toimingute kulusid vältida.

Kasutame selleks objekti countStream, mis on loendis 5.11 (vastav kood on saadaval failis src/main/java/bbejeck/chapter_5/GlobalKTableExample.java), ühendades selle kahe GlobalKTable objektiga.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Oleme seda juba varem arutanud, seega ei hakka kordama. Täiendavalt mainin, et toStream().map funktsiooni kood on lugemise mugavuse huvides abstrakteeritud objekti funktsiooniks, mitte sisseehitatud lambda-väljendiks.

Järgmine etapp on kahe GlobalKTable eksemplari kuulutamine (antud kood on saadaval failis src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (loendis 5.12).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»

Märkusena: teemasid kirjeldatakse loetletud tüüpide abil.

Nüüd, kui oleme kõik komponendid ette valmistanud, tuleb kirjutada kood ühendamiseks (mida leiate failist src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (loendis 5.13).

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Kuigi selles koodis on kaks ühendust, on need korraldatud ahelana, kuna eraldi nende tulemusi ei kasutata. Tulemused kuvatakse kogu operatsiooni lõpus.

Kui käitate eeltoodud ühendusoperatsiooni, saate tulemused järgmise kujul:

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

Sisu ei ole muutunud, kuid need tulemused näevad nüüd arusaadavamad välja.

Kui arvestada peatükki 4, olete juba näinud mitmeid ühendustüüpe tegutsemas. Need on loetletud tabelis 5.2. See tabel peegeldab ühendamise võimalusi, mis on kehtivad versioonile 1.0.0 Kafka Streams; tulevastes väljalasetel võib midagi muutuda.

Raamat «Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas töötamiseks»
Kokkuvõtteks meenutan peamist: saate siduda sündmuste vooge (KStream) ja uuenduste vooge (KTable) kohaliku olekuga. Lisaks, kui viidata isenduse suurus ei ole liiga suur, saate kasutada objekti GlobalKTable. GlobalKTable replitseerib kõik sektsioonid iga Kafka Streams rakenduse sõlme, tagades seeläbi kõikide andmete kättesaadavuse sõltumata sellest, millisele sektsioonile key vastab.

Edasi näeme Kafka Streams'i võimalust, mis võimaldab jälgida olekumuutusi ilma andmete tarbimiseta Kafka teemas.

5.3.5. Küsimiseks saadaval olek

Oleme juba teinud mitu toimingut, mis hõlmavad olekuteavet, ja oleme alati tulemused konsoolile edastanud (arenduseesmärkidel) või salvestanud need teemale (tööstuslikuks kasutamiseks). Tulemuste teemasse salvestamisel tuleb kasutada Kafka tarbijat nende vaatamiseks.

Andmete lugemist nendest teemadest võib pidada materialiseeritud vaadete variandiks. Meie vajaduste jaoks saame kasutada materialiseeritud vaate määratlust „Wikipediast”: „…andmebaasi füüsiline objekt, mis sisaldab päringu täitmise tulemusi. Näiteks võib see olla kaugandmete kohalik koopia, või ridade ja / või veergude alamhulk tabelis, või ühendamise tulemuste alusel koostatud kokkuvõtte tabel.” (https://en.wikipedia.org/wiki/Materialized_view)

Kafka Streams võimaldab ka interaktiivsete päringute (interactive queries) tegemist olekute salvestustele, mis annab võimaluse nende materialiseeritud vaadete otse lugemiseks. Oluline on märkida, et päring olekute salvestusele on «ainult lugemise» operatsioon. See tähendab, et te ei pea kartma, et juhuslikult muudate oleku kooskõlas andmete töötlemisel rakenduses.

Olekute salvestustele otse tehtavad päringud on väga olulised. Need võimaldavad luua rakendusi — informatiivseid paneele, ilma et oleks esmalt vaja andmeid Kafka tarbijalt hankida. Samuti suurendab see rakenduse tõhusust, kuna andmeid pole vaja uuesti salvestada:

  • andmete lokaliteeduse tõttu on neile kiire ligipääs;
  • välistatakse andmete dubleerimine, kuna neid ei salvestata välistesse salvestustesse.

Peamine, mida soovin, et sa meeles peaksid: saate otse rakendusest teha päringuid olekuteavet. Selle võimaluste üle ei saa üle hinnata. Selle asemel, et tarbida andmeid Kafka'st ja salvestada kirjeid rakenduse andmebaasi, saate teha päringuid olekuhoidlatele sama tulemuse saavutamiseks. Otsesed päringud olekuhoidlatele tähendavad vähem koodi (tarbijat pole) ja vähem tarkvara (puudub vajadus andmetabeli järele tulemuste salvestamiseks).

Oleme käesolevas peatükis käsitlenud suurt hulka teavet, seega lõpetame nüüd interaktiivsete päringute arutelu olekuhoidlatele. Kuid ärge muretsege: peatükis 9 loome lihtsa rakenduse — infopaneeli interaktiivsete päringutega. Interaktiivsete päringute ja nende rakendustesse lisamise võimaluste demonstreerimiseks kasutatakse seal mõningaid selle ja eelnevate peatükkide näiteid.

Kokkuvõte

  • KStream objektid esindavad sündmuste vooge, mis on sarnased andmebaasi sisestustele. KTable objektid esindavad uuenduste vooge ja need sarnanevad rohkem andmebaasi uuendustega. KTable objekti suurus ei kasva, vana teave asendatakse uuega.
  • KTable objektid on vajalikud aggregatsioonioperatsioonide jaoks.
  • Aknate operatsioonide abil on võimalik jagada kogutud andmed ajavahemike kaupa.
  • GlobalKTable objektide kaudu on võimalik juurde pääseda käsiraamatutele igas rakenduse kohas, olenemata osade jaotamisest.
  • KStream, KTable ja GlobalKTable objektide omavaheline ühendamine on võimalik.

Siiani oleme keskendunud Kafka Streams'i rakenduste loomisele, kasutades kõrgetasemelist DSL KStream. Kuigi kõrgetasemeline lähenemine võimaldab luua selgeid ja lühikesi programme, toob selle kasutamine kaasa teatud kompromissi. DSL KStream'i kasutamine suurendab koodi lühidust, kuid vähendab kontrolli taset. Järgmisel peatükil uurime madalama taseme API töötlemismooduleid ja proovime muid kompromisse. Programmid saavad olema pikemad kui seni, kuid see annab meile võimaluse luua praktiliselt igasuguseid töötlemismooduleid, mis võivad osutuda vajalikuks.

→ Raamatuga saab lähemalt tutvuda väljaandja veebisaidil

→ Habr'i kasutajatele 25% allahindlus kupongiga — Kafka Streams

→ Pärast füüsilise raamatu tasumist saadetakse elektrooniline raamat e-posti aadressile.

Allikas: habr.com

Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid 🔥 Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster