Tere, Habrojekte! See raamat sobib igale arendajale, kes soovib voogude töötlemisse sĂŒveneda. Jaotatud programmeerimise mĂ”istmine aitab paremini mĂ”ista Kafka ja Kafka Streams'i. Oluline oleks tunda ka Kafka raamistikku, kuid see pole hĂ€davajalik: ma selgitan teile kĂ”ik, mida peate teadma. Kogenud Kafka arendajad, nagu ka algajad, saavad selle raamatu kaudu meisterdada huvitavaid voogude töötlemise rakendusi, kasutades Kafka Streams'i raamatukogu. Kesk- ja kĂ”rgema taseme Java arendajad, kes on juba tuttavad selliste mĂ”isted nagu serialiseerimine, Ă”pivad, kuidas rakendada oma oskusi Kafka Streams'i rakenduste loomisel. Raamatu lĂ€htekood on kirjutatud Java 8-s ja kasutab ulatuslikult Java 8 lambda-vĂ€ljendite sĂŒntaksit, seega tuleb kasuks oskused lambda-funktsioonide töötluseks (isegi teises programmeerimiskeeles).
LÔik. 5.3. Aggregeerimine ja aknoperatsioonid
Selles jaotises viime lÀbi uurimise kÔige lubavamatest osadest Kafka Streams'ist. Kuni praeguseni oleme kÀsitlenud jÀrgmisi Kafka Streams'i aspekte:
- töötlemise topoloogia loomine;
- oleku kasutamine voogude rakendustes;
- andmevoogude ĂŒhendamine;
- erinevused sĂŒndmusvoogude (KStream) ja uuenduste voogude (KTable) vahel.
JĂ€rgmistes nĂ€idetes koondame kĂ”ik need elemendid kokku. Samuti tutvute aknoperatsioonidega â veel ĂŒhe suurepĂ€rase vĂ”imalusega voogude rakendustes. Meie esimene nĂ€ide on lihtne aggregeerimine.
5.3.1. Aktsiate mĂŒĂŒgimahtude aggregeerimine tööstusharu kaupa
Aggregeerimine ja grupeerimine on eluliselt vajalikud tööriistad voogude andmete töötlemisel. Ăksikute rekordite uurimine reaalajas osutub tihti ebapiisavaks. Andmetest lisainfo vĂ€ljavĂ”tmiseks on vajalikud gruppeerimine ja kombineerimine.
Selles nĂ€ites peate end intraday kaupleja kostĂŒĂŒmile panema, kes peab jĂ€lgima aktsiate mĂŒĂŒgimahu mahtu mitmetes tööstusharudes. Eriti huvitavad teid viis ettevĂ”tet, mille aktsiate mĂŒĂŒgimaht on igas tööstusharus suurim.
Selle aggregeerimise jaoks on vaja mitmeid jĂ€rgmised samme andmete vajaliku vormi tĂ”lkimiseks (ĂŒldiselt öeldes).
- Loo source pĂ”hined teemarikkal, mis avaldab aktsiate kauplemise kohta toorinfot. Me peame kuvama StockTransaction tĂŒĂŒpi objekti ShareVolume tĂŒĂŒbiks. Asi on selles, et StockTransaction objekt sisaldab mĂŒĂŒgimetainfot, kuid me vajame ainult andmeid mĂŒĂŒdud aktsiate arvu kohta.
- Klasifitseerime ShareVolume andmeid aktsia sĂŒmbolite kaupa. PĂ€rast sĂŒmbolite jĂ€rgi grupeerimist saame need andmed kokku pöörata aktsiate mĂŒĂŒgi vahekokkuvĂ”teteks. Tuleb mĂ€rkida, et meetod KStream.groupBy tagastab KGroupedStream tĂŒĂŒpi eseme. KTable'i saame vĂ”tta, kutsudes edasi KGroupedStream.reduce meetodit.
Mis on KGroupedStream liides?
Meetodid KStream.groupBy ja KStream.groupByKey tagastavad KGroupedStream tĂŒĂŒpi eseme. KGroupedStream on vahepealne esitus sĂŒndmuste voost pĂ€rast grupeerimist vĂ”tmete jĂ€rgi. See ei ole mĂ”eldud otseseks töötamiseks. Selle asemel kasutatakse KGroupedStream'i kogumise toiminguteks, mille tulemuseks on alati KTable. Kuna kogumise tegemise tulemuseks on KTable ja see katab olekuhoidla, vĂ”ivad mitte kĂ”ik uuendused liikuda edasi torustiku kaudu.
Meetod KTable.groupBy tagastab sarnase KGroupedTable'i â vahepealse esitluse vĂ€rskenduste voost, mis on uuesti grupeeritud vĂ”tme jĂ€rgi.
Teeme vÀikese pausi ja vaatame joonist 5.9, mis nÀitab, mida me saavutasime. See topoloogia peaks teile juba tuttav olema.

Vaadakem nĂŒĂŒd koodi selle topoloogia jaoks (seda saab leida failist src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (loetelu 5.2).

Esitatud koodi iseloomustab lĂŒhidus ja suur hulk tehteid, mis tehakse mĂ”ne reaga. Meetodi builder.stream esimeses parameetris vĂ”ite mĂ€rgata midagi uut: loetletud tĂŒĂŒbi vÀÀrtus AutoOffsetReset.EARLIEST (on olemas ka LATEST), mis on mÀÀratud meetodi Consumed.withOffsetResetPolicy abil. Selle loetletud tĂŒĂŒbi abil saate mÀÀrata nihkestrateegia iga KStream vĂ”i KTable jaoks, see on eeliseks nihkemeetodi seadistusele konfiguratsioonis.
GroupByKey ja GroupBy
KStream liideses on kaks meetodit kirje gruppeeringuks: GroupByKey ja GroupBy. MĂ”lemad tagastavad KGroupedTable'i, seega vĂ”ite kĂŒsida: milline on nende vahe ja millal kumbagi kasutada?
GroupByKey meetodit rakendatakse, kui KStream'i vĂ”ti on juba olemas. Pealegi ei ole lipp 'nĂ”uab ĂŒmberjaotamist' kunagi seadistatud.
GroupBy meetod eeldab, et olete muutnud rĂŒhmitamise vĂ”tmeid, mistĂ”ttu on ĂŒmberjaotamise lipp seadistatud vÀÀrtusele true. GroupBy jĂ€rel teostatavad ĂŒhendused, agregatsioonid jne viivad automaatse ĂŒmberjaotamiseni.
KokkuvÔtteks: kui vÀhegi vÔimalik, tuleks kasutada GroupByKey meetodit, mitte GroupBy'd.
Kuna mapValues ja groupBy meetodite funktsioone on selge, vaatame sum() meetodit (see on leitav failist src/main/java/bbejeck/model/ShareVolume.java) (loetelu 5.3).

ShareVolume.sum meetod tagastab aktsiate mĂŒĂŒgimahtude vahepealse summa ning kogu arvutuste ahela tulemus on objekt KTable. NĂŒĂŒd mĂ”istate, millist rolli KTable mĂ€ngib. Kui ShareVolume objekte lisatakse, salvestatakse vastavas KTable objektis viimane asjakohane uuendus. On oluline mitte unustada, et kĂ”ik uuendused kajastuvad eelnevas shareVolumeKTable'is, kuid mitte kĂ”ik saadetakse edasi.
Edasi, kasutades seda KTable'i, teeme agregatsiooni (mĂŒĂŒdud aktsiate arvu alusel), et saada viis ettevĂ”tet, kellel on iga tööstusharu suurimad aktsiate mĂŒĂŒgivoldid. Meie tegevus on sarnane esimese agregatsiooniga.
- Tehke veel ĂŒks groupBy operatsioon, et rĂŒhmitada eraldi ShareVolume objekte tööstusharude jĂ€rgi.
- Alustage ShareVolume objektide summeerimisega. Seekord on agregatsiooni objectiks fikseeritud suurusega prioriteediköit. Sellises fikseeritud suurusega köites hoitakse ainult viit ettevĂ”tet, kellel on kĂ”ige rohkem mĂŒĂŒdud aktsiaid.
- Muuda eelneva punkti jĂ€rjekorrad stringivÀÀrtuseks ja tagasta viie kĂ”ige paremini mĂŒĂŒdud aktsia pĂ”hjal nende tööstusharude jĂ€rgi.
- Kanna tulemused stringivormis teemasse.
Joonisel 5.10 on kujutatud andmevoogude topoloogia graafik. Nagu nĂ€ete, on teise töötlemise tsĂŒkkel ĂŒsna lihtne.

NĂŒĂŒd, kui olete saanud selge arusaama selle teise töötlemise tsĂŒkli struktuurist, vĂ”ite tutvuda selle lĂ€htekoodiga (leiate selle failist src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (loetelu 5.4).
Selles initsialisaatoris on muutuja fixedQueue. See on kohandatud objekt â adapter java.util.TreeSet'ile, mida kasutatakse, et jĂ€lgida N suurimat mĂŒĂŒgihulka kahanevas jĂ€rjestuses.

Te olete juba tuttavad groupBy ja mapValues vÀljakutsetega, seega ei jÀÀme neile pikemalt peatuma (me kutsume KTable.toStream meetodit, kuna KTable.print meetod on vananenud). Kuid te pole veel nÀinud KTable versiooni meetodist aggregate(), seega kulutame veidi aega selle arutamiseks.
Nagu te meenute, erineb KTable sellest, et sama vĂ”tmega kirjeid peetakse uuendusteks. KTable asendab vana kirje uuega. Aggregeerimine toimub sarnaselt: aggregeeritakse viimased kirjed sama vĂ”tmega. Kirje saabudes lisatakse see FixedSizePriorityQueue klassi eksemplarisse koos summeerijaga (teine parameeter aggregate meetodi ĂŒleskutses), kuid kui sama vĂ”tmega on juba olemas teine kirje, siis vana kirje eemaldatakse vĂ€hendajaga (kolmas parameeter aggregate meetodi ĂŒleskutses).
See kĂ”ik tĂ€hendab, et meie aggregeerija, FixedSizePriorityQueue, ei aggregeeri kĂ”iki sama vĂ”tmega vÀÀrtusi, vaid hoiab N kĂ”ige enam mĂŒĂŒdud aktsiate liikide libisevat summat. Igas saabunud kirjes on ĂŒldine mĂŒĂŒdud aktsiate hulk seni. KTable annab teile teavet selle kohta, milliste ettevĂ”tete aktsiaid mĂŒĂŒakse praegu rohkem, libisev aggregeerimine iga uuenduse puhul ei ole vajalik.
Oleme Ôppinud tegema kahte olulist asja:
- gruppima KTable's vÀÀrtusi nende ĂŒhise vĂ”tme jĂ€rgi;
- teostama nendele rĂŒhmitatud vÀÀrtustele kasulikke toiminguid, nagu kokkuvĂ”te ja aggregeerimine.
Neid toiminguid osata on oluline, et mÔista andmete tÀhendust, mis liigub lÀbi Kafka Streams rakenduse, ja vÀlja selgitada, millist teavet nad kannavad.
Me oleme samuti kokku viinud mĂ”ned vĂ”tmekontseptsioonid, millest oleme varem selles raamatus rÀÀkinud. 4. peatĂŒkis rÀÀkisime, kui oluline on voogedastusrakenduse jaoks talletatud, lokaalne olek. Selles peatĂŒkis toodud esimene nĂ€ide demonstreeris, miks lokaalne olek on nii oluline â see vĂ”imaldab jĂ€lgida, millist teavet olete juba nĂ€inud. Lokaalne ligipÀÀs aitab vĂ€ltida vĂ”rgulisi viivitusi, muutes rakenduse tootlikumaks ja vigu taluvamaks.
Iga koondamis- vÔi agregatsioonitegevuse lÀbiviimisel tuleb tÀpsustada olekutehnika nime. Koondamis- ja agregatsioonitegevused tagastavad KTable'i koopia, samal ajal kui KTable kasutab olekutehnikat, et asendada vanad tulemused uutega. Nagu olete nÀinud, ei saadeta kaugeltki kÔiki vÀrskendusi torust edasi, ja see on oluline, kuna agregatsioonitegevused on mÔeldud lÔpptulemuse saamiseks. Ilma lokaalse olekuta saadab KTable edasi kÔik agregatsiooni ja koondamise tulemused.
JĂ€rgmises vaates vaatame lĂ€bi selliseid tegevusi nagu agregatsioon kindla ajavahemiku jooksul â nii nimetatud aktsioonitegevusi (windowing operations).
5.3.2. Aktsioonitegevused
Eelnevas peatĂŒkis tutvusime "libiseva" koondamise ja agregatsiooniga. Rakendus tegi pidevat koondamist aktsiate mĂŒĂŒgimahu kohta, jĂ€rgides viie kĂ”ige rohkem mĂŒĂŒdud aktsia koondamist.
MÔnikord on sarnased pidevad agregatsioonid ja koondamised vajalikud. Kuid mÔnikord tuleb teostada tegevusi ainult mÀÀratud ajavahemiku jooksul. NÀiteks vÀlja arvutada, kui palju börsil tehinguid on tehtud konkreetse ettevÔtte aktsiatega viimase 10 minuti jooksul. VÔi kui palju kasutajaid klÔpsas uuel reklaambÀnneril viimase 15 minuti jooksul. Rakendus vÔib neid tegevusi teha korduvalt, kuid tulemustega, mis on seotud ainult mÀÀratud ajavahemikega (ajaliste akendega).
Börsil tehingute arvestamine ostja jÀrgi
JĂ€rgmises nĂ€ites meie teema on börsil tehingute jĂ€lgimine mitme kaupleja puhul â kas suurte ettevĂ”tete vĂ”i nutikate individuaalsete finantsteadlaste puhul.
On olemas kaks vĂ”imalikku pĂ”hjust sarnase jĂ€lgimise jaoks. Ăks neist on vajadus teada, mida ostavad/mĂŒĂŒvad turuliidrid. Kui need suured mĂ€ngijad ja kogenud investorid nĂ€evad endale avanevaid vĂ”imalusi, on mĂ”istlik jĂ€rgida nende strateegiat. Teine pĂ”hjus on soov mĂ€rgata vĂ”imalikke ilmsi ebaseaduslikke tehingutega, mis kasutavad siseteavet. Selleks peate analĂŒĂŒsima suurte mĂŒĂŒgivoolude ja oluliste pressiteadete korrelatsiooni.
Selline jÀlgimine hÔlmab jÀrgmisi etappe:
- loome voolu lugemiseks teemast stock-transactions;
- grupime sissetulevad kirjed ostja ID ja aktsia sĂŒmboli jĂ€rgi. GroupBy meetodi kutsumine tagastab KGroupedStream klassi eksemplari;
- KGroupedStream.windowedBy meetodi abil tagastatakse ajavahemiku piiramiseks andmevoog, mis vĂ”imaldab teostada akna agregatsiooni. Olenemata akna tĂŒĂŒbist tagastatakse kas TimeWindowedKStream vĂ”i SessionWindowedKStream;
- tehingute arvestamine agregatsiooni toimingu jaoks. Aken mÀÀrab, kas konkreetne kande arvesse vÔtmine arvutamisel toimub;
- tulemuste salvestamine teemasse vÔi nende vÀljastamine konsooli arenduse kÀigus.
Rakenduse topoloogia on lihtne, kuid visuaalne esitus ei teeks kahju. Vaatame joonist 5.11.
Edasi vaatleme akna toimingute funktsionaalsust ja vastavat koodi.

Akna tĂŒĂŒbid
Kafka Streams'is on kolm tĂŒĂŒpi aknaid:
- seanssi;
- ârulluvadâ (tumbling);
- libisevad/âhĂŒppavadâ (sliding/hopping).
Mille valida â sĂ”ltub Ă€rinĂ”uetest. âRulluvadâ ja âhĂŒppavadâ aknad piirduvad ajaga, samas kui seansside piirangud on seotud kasutajate tegevustega â seansi (seansside) kestus mÀÀratakse ainuilt selle jĂ€rgi, kui aktiivselt kĂ€itub kasutaja. Peamine on mitte unustada, et kĂ”ik akna tĂŒĂŒbid pĂ”hinevad kannete kuupĂ€eva/aegade mĂ€rgistel, mitte sĂŒsteemi ajast.
Edasi rakendame meie topoloogiat iga akna tĂŒĂŒbi korral. TĂ€ielik kood toodakse ainult esimeses nĂ€ites, teiste akna tĂŒĂŒpide jaoks ei muutu midagi peale akna toimingu tĂŒĂŒbi.
Seanssi aknad
Seanside aknad erinevad oluliselt kĂ”igist teistest aknatĂŒĂŒpidest. Need on piiratud mitte niivĂ”rd ajaga, vaid kasutaja tegevusega (vĂ”i selle ĂŒksuse tegevusega, mida soovite jĂ€lgida). Seanside aknad eraldatakse tegevusetuse perioodide kaupa.
Joonisel 5.12 illustreeritakse seanside akna mÔistet. VÀiksem seanss sulandub temast vasakul asuva seansi sisse. Paremal asuv seanss on eraldi, kuna jÀrgneb pikale tegevusetuse perioodile. Seanside aknad pÔhinevad kasutajate tegevusel, kuid kasutavad kuupÀeva/kellaaja mÀrke, et mÀÀrata, mille seansi alla vastav kirje kuulub.

Seanside aknade kasutamine börsitehingute jÀlgimiseks
Kasutame seanside aknaid börsitehingute teabe kogumiseks. Seanside akna rakendamine on nÀidatud loendis 5.5 (mida vÔib leida failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKTableJoinExample.java).

Enamik selle topoloogia operatsioone on teile juba tuttavad, nii et pole vajalik neid siin uuesti kÀsitleda. Aga siin on ka mÔned uued elemendid, millest me praegu rÀÀgime.
Iga groupBy operatsiooni puhul tehakse tavaliselt mingisugune agregatsioonioperatsioon (agregatsioon, kokkusurumine vÔi arvu lugemine). Saame teostada kas kumulatiivset agregatsiooni jooksva kogusummaga vÔi akna agregatsiooni, kus arvesse vÔetakse kirjeid mÀÀratud ajavahemiku jooksul.
Loendi 5.5 kood loendab tehingute arvu seanside akende piires. Joonisel 5.13 analĂŒĂŒsitakse neid toiminguid samm-sammult.
Kasutades funktsiooni windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)), loome seanside akna 20-sekundilise tegevusetuse intervalliga ja 15-minutilise sÀilitamisintervalliga. 20-sekundiline tegevusetuse intervall tÀhendab, et rakendus hÔlmab igat kirjet, mis saabub 20 sekundi jooksul pÀrast praeguse seansi lÔppu vÔi algust, praegusesse (aktiivsesse) seanssi.

Edasi anname teada, millist agregatsiooni toimingut tuleb seanssi aknas teostada - antud juhul count. Kui sisenev kirje ĂŒletab tegevusetuse intervalli (mĂ”lemalt poolt ajatempli), siis rakendus loob uue seansi. Salvestamise intervall tĂ€hendab seansi hoidmist teatud aja jooksul ja vĂ”imaldab hilinenud andmeid, mis ĂŒletavad seansi tegevusetuse perioodi, kuid vĂ”ivad siiski liituda. Lisaks on uue seansi algus ja lĂ”pp, mis tulenevad ĂŒhendamisest, vastavalt varasemale ja hilisemale ajatempli.
Vaadakem mÔningaid kirjeid count meetodist, et nÀha, kuidas seansid toimivad (tabel 5.1).

Kirjete saabumise korral otsime juba olemasolevaid seansse sama vĂ”tme jaoks, mille lĂ”pp aeg on vĂ€hem kui praegune ajatempli - tegevusetuse intervall ja mille algusaeg on rohkem kui praegune ajatempli + tegevusetuse intervall. Selle jĂ€rel ĂŒhendatakse neli kirjet tabelist 5.1 ĂŒhte seanssi jĂ€rgmiselt.
1. Esimesena saabub kirje 1, nii et algusaeg on sama, mis lÔppaeg ja on 00:00:00.
2. JĂ€rgmisena saabub kirje 2, ja me otsime seansse, mis lĂ”pevad hiljemalt 23:59:55 ja algavad hiljemalt 00:00:35. Leiame kirje 1 ja ĂŒhendame seansid 1 ja 2. VĂ”tame seansi 1 algusaja (varasem) ja seansi 2 lĂ”ppaaja (hilisem), nii et meie uus seanss algab 00:00:00 ja lĂ”peb 00:00:15.
3. Saabub kirje 3, me otsime seansse ajavahemikus 00:00:30 kuni 00:01:10 ja ei leia ĂŒhtegi. Lisame teise seansi vĂ”tmega 123-345-654,FFBE, algusega ja lĂ”puga 00:00:50.
4. Saabub kirje 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.
KĂ€esolevas jaotises on mĂ”ningaid olulisi nĂŒansse, mida tuleks meeles pidada:
- seansid ei ole fikseeritud suurusega aknad. Seansi kestus mÀÀratakse aktiveerituse jÀrgi mÀÀratud ajavahemikus;
- ajatempli andmed mÀÀravad, kas sĂŒndmus kuulub olemasolevasse seanssi vĂ”i tegevusetuse intervalli.
Edasi arutame jĂ€rgmist akna tĂŒĂŒpi - âveerevadâ aknad.
âVeerevadâ aknad
âVeerlevadâ (tumbling) aknad haaravad sĂŒndmusi, mis toimuvad kindlas ajavahemikus. Kujutlege, et peate salvestama kĂ”iki mingisuguse ettevĂ”tte börsitehinguid iga 20 sekundi tagant, seega kogute kĂ”ik sĂŒndmused selle ajavahemiku jooksul. 20-sekundilise intervalli lĂ”ppedes âveerlebâ aken ja lĂ€heb uude 20-sekundilise vaatlusperioodi. Joonis 5.14 illustreerib seda olukorda.

Kuidas nĂ€ha, on kĂ”ik viimase 20 sekundi jooksul saabunud sĂŒndmused aknas. Selle ajavahemiku lĂ”ppedes luuakse uus aken.
Loendis 5.6 on toodud kood, mis demonstreerib âveerlevateâ akende kasutamist börsitehingute igakuiseks salvestamiseks iga 20 sekundi tagant (seda vĂ”ib leida failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

Selle vĂ€ikese muudatuse tĂ”ttu meetodi TimeWindows.of kutsumisel on vĂ”imalik kasutada âveerlevatâ akent. Antud nĂ€ites ei ole meetodi until() kutsumist, mistĂ”ttu kasutatakse vaikimisi salvestusintervalli, mis on 24 tundi.
LĂ”puks on aeg liikuda viimase aknatĂŒĂŒbi juurde â âhĂŒppavadâ (hopping) aknad.
Libisevad (âhĂŒppavadâ) aknad
Libisevad/âhĂŒppavadâ (sliding/hopping) aknad on sarnased âveerlevateleâ, 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 ooteaega, mis on lĂŒhem akna kestusest.
Erinevuste illustreerimiseks âveerlevateâ ja âhĂŒppavateâ akende vahel pöördume tagasi börsitehingute arvestuse nĂ€ite juurde. Meie eesmĂ€rk on endiselt arvestada tehingute arvu, kuid me ei soovi oodata kogu ajavahemikku enne arvu vĂ€rskendamist. Selle asemel vĂ€rskendame arvu lĂŒhemate ajavahemike tagant. NĂ€iteks loeme endiselt tehingute arvu iga 20 sekundi tagant, kuid vĂ€rskendame arvu iga 5 sekundi tagant, nagu on nĂ€idatud joonisel 5.15. Selle tulemusena on meil kolm tulemuste akent, millel on kattuvad andmed.

Loendis 5.7 on toodud kood libisevate akende seadistamiseks (seda vÔib leida failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

âVeerevâ akent saab muuta âhĂŒppavaksâ funktsiooni advanceBy() kutsumisega. Antud nĂ€ites on salvestusintervall 15 minutit.
Selles osas olete nÀinud, kuidas piirata tulemusi ajavahemikega. Eriti tasub meeles pidada kolm asja, mida sellelt, jaotiselt kaasa vÔtta:
- seanss akna suurus ei ole piiratud aja jooksul, vaid kasutajate aktiivsusega;
- âveerlevadâ aknad annavad ĂŒlevaate sĂŒndmustest antud ajavahemiku jooksul;
- âhĂŒppavateâ akende kestus on kindel, kuid need uuendatakse sageli ja vĂ”ivad kĂ”ikides akendes sisaldada kattuvaid kirjeid.
SeejĂ€rel Ă”ppime, kuidas KTable tagasi KStreami muuta ĂŒhendamiseks.
5.3.3. KStreami ja KTable'i ĂŒhendamine
Kapis 4 arutasime kahe KStreami objekti ĂŒhendamist. NĂŒĂŒd Ă”pime KTable'i ja KStreami ĂŒhendamist. Seda vĂ”ib vaja minna jĂ€rgmise lihtsa pĂ”hjuse tĂ”ttu. KStream on kirjevoog, KTable aga kirjeuuenduste voog, kuid mĂ”nikord vĂ”ib olla vajalik lisada voogudele konteksti KTable'i uuenduste kaudu.
VĂ”tame andmed börsitehingute arvu kohta ja ĂŒhendame need börsiuudistega vastavate tööstusharude jĂ€rgi. Siin on, mida tuleb teha, et saavutada seda olemasoleva koodiga.
- Muutke börsitehingute arvu KTable objekt KStreamiks, asendades vĂ”tme nii, et see tĂ€hendaks tööstusharu, mis vastab antud aktsiate sĂŒmbolile.
- Looge KTable objekt, mis loeb andmeid börsiuudiste teemast. See uus KTable klassifitseeritakse tööstusharude jÀrgi.
- Ăhendage uudiste uuendused börsitehingute arvu teabega vastavates tööstusharudes.
NĂŒĂŒd vaatame, kuidas seda tegevuskava ellu viia.
KTable'i muutmine KStreamiks
KTable'i KStreamiks muutmiseks tuleb teha jÀrgmist.
- Kutsuda vÀlja meetod KTable.toStream().
- Kasutades meetodi KStream.map kutset, asendada vÔti tööstusharu nimega ja seejÀrel vÀlja vÔtta TransactionSummary objekti.
Seome need toimingud ahelaks jÀrgmiselt (koodi leiate failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (loend 5.8).

Kuna me teeme operatsiooni KStream.map, toimub KStreami tagastatud eksemplari ĂŒmberjaotamine automaatselt, kui seda kasutatakse ĂŒhenduses.
Oleme lĂ”petanud muundamisprotsessi, nĂŒĂŒd peame looma KTable objekti börsiuudiste lugemiseks.
KTable loomine börsiuudiste jaoks
Ănneks piisab KTable objekti loomiseks ĂŒhest koodireast (seda koodi saab leida failist src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (loetelu 5.9).

Oluline on mÀrkida, et Serde objekte ei ole vaja mÀÀrata, kuna seadistustes kasutatakse stringide Serde-d. Samuti tÀidab tabel tÀnu EARLIEST loetlemise kasutamisele kirjeid alguses.
NĂŒĂŒd saame liikuda viimase sammu juurde â ĂŒhendamine.
Uudiste vĂ€rskenduste ja tehingute arvu andmete ĂŒhendamine
Ăhenduse loomine ei ole keeruline. Kasutame vasakut ĂŒhendust juhuks, kui vastavas tööstusharus pole börsiuudiseid (vajalik kood on saadaval failis src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (loetelu 5.10).

See leftJoin operaator on ĂŒsna lihtne. Erinevalt 4. peatĂŒkist, kus kasutati JoinWindow meetodit, ei kasutata KStream-KTable ĂŒhendamisel, kuna KTable'is on iga vĂ”tme kohta ainult ĂŒks kirje. Selline ĂŒhendamine ei ole ajaliselt piiratud: kirje kas on KTable'is vĂ”i pole. Peamine jĂ€reldus: KTable objektide abil saab KStreami rikastada harvemini uuendatud viidatud andmetega.
NĂŒĂŒd vaatame tĂ”husamat viisi KStreami sĂŒndmuste rikastamiseks.
5.3.4. GlobalKTable objektid
Kuidas te nĂŒĂŒd aru saite, on vaja rikastada sĂŒndmuste vooge vĂ”i lisada neile konteksti. 4. peatĂŒkis nĂ€gite kahe KStream objekti ĂŒhendamist, ja eelmises jaos KStreami ja KTable'i ĂŒhendamist. KĂ”igil nendel juhtudel on andmevoo ĂŒmberjaotamine vajalik, kui vĂ”tmed muudetakse uue tĂŒĂŒbi vĂ”i vÀÀrtuse jaoks. MĂ”nikord toimub ĂŒmberjaotamine selgelt, mĂ”nikord teeb Kafka Streams seda automaatselt. Ămberjaotamine on vajalik, kuna vĂ”tmed on muutunud ja kirjed peavad olema uuteks sektsioonideks, vastasel juhul ei saa ĂŒhendust teha (sellest arutati 4. peatĂŒkis jaotises "Andmete ĂŒmberjaotamine" alajaos 4.2.4).
Uuesti jaotamine on tasuline
Uuesti jaotamine nĂ”uab kulutusi â tĂ€iendavaid ressursikulu vahepealsete teemade loomisele, duplikaatandmete sĂ€ilitamisele veel ĂŒhes teemas; see toob kaasa ka viivituste suurenemise seoses andmete lugemise ja kirjutamisega sellesse teemasse. Lisaks, kui on vajalik teostada ĂŒhendust rohkem kui ĂŒhe aspekti vĂ”i mÔÔtme alusel, tuleb ĂŒhendused organiseerida ahelana, kuvada kirjed uute vĂ”tmete jĂ€rgi ja uuesti viia lĂ€bi uuesti jaotamise protsess.
Ăhendamine vĂ€iksemate andmekogumitega
MĂ”nes olukorras on viidatud andmete maht, millega plaanitakse ĂŒhendust, suhteliselt vĂ€ike, nii et nende tĂ€iendavad koopiad mahuvad igasse sĂ”lme kohapeal. Selle jaoks on Kafka Streams'is olemas klass GlobalKTable.
GlobalKTable'i eksemplarid on ainulaadsed, kuna rakendus replitseerib kĂ”ik andmed igasse sĂ”lmesse. Kuna igas sĂ”lmes on kĂ”ik andmed, ei ole vaja sĂŒndmuste voogu jaotada viidatud andmete vĂ”tme alusel, et need oleksid kergesti ligipÀÀsetavad kĂ”ikidele sektsioonidele. GlobalKTable objektidega on samuti vĂ”imalik teostada vĂ”tmeta ĂŒhendusi. Naaseme tagasi ĂŒhe eelneva nĂ€ite juurde, et demonstreerida seda vĂ”imalust.
KStream objektide ĂŒhendamine GlobalKTable objektidega
Alajaotises 5.3.2 viisime lĂ€bi akna agregatsiooni börsitehingute ĂŒle 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'}, 16Kuigi need tulemused vastasid seatud eesmĂ€rgile, oleks mugavam, kui seal oleks ka kliendi nimi ja ettevĂ”tte tĂ€ielik nimi. Kliendi nime ja ettevĂ”tte nime lisamiseks on vĂ”imalik teha tavalisi ĂŒhendusi, kuid see nĂ”uab kahe vĂ”tme kaardistamist ja uuesti jaotamist. GlobalKTable abil saab selliseid toimingute kulutusi vĂ€ltida.
Selleks kasutame objekti countStream nĂ€iteks 5.11 (vastav kood asub failis src/main/java/bbejeck/chapter_5/GlobalKTableExample.java), ĂŒhendades selle kahe GlobalKTable objektiga.

Oleme seda juba varem arutanud, seega ei hakka ma seda kordama. Kuid mainin, et funktsiooni toStream().map kood on loetavuse huvides abstraheeritud objekti-funktsiooni, mitte sisseehitatud lambdana.
JĂ€rgmine etapp on kahe GlobalKTable instantsi deklareerimine (toodud kood on saadaval failis src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (loetelu 5.12).

Pange tĂ€hele, et teema nimed on mÀÀratletud loendatavate tĂŒĂŒpide abil.
NĂŒĂŒd, kui oleme kĂ”ik komponendid ette valmistanud, jÀÀb kirjutada kood ĂŒhendamiseks (mida saab leida failist src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (loetelu 5.13).

Kuigi selles koodis on kaks ĂŒhendust, on need korraldatud ahelaks, kuna eraldi ei kasutata ĂŒhtegi nende tulemusi. Tulemused vĂ€ljastatakse kogu operatsiooni lĂ”pus.
KĂ€ivitades ĂŒlaltoodud ĂŒhenduse operatsiooni, saate jĂ€rgmiste sĂ”numite tulemusi:
{customer='Barney, Smith' company="Exxon", transactions= 17}Sisu pole muutunud, kuid need tulemused nÀevad vÀlja arusaadavamad.
Kui arvestada peatĂŒkki 4, olete juba nĂ€inud mitmeid tĂŒĂŒpe ĂŒhendusi tegevuses. Need on loetletud tabelis 5.2. See tabel kajastab ĂŒhenduste vĂ”imalusi, mis on seotud Kafka Streams versiooniga 1.0.0; tulevastes vĂ€ljalasemetes vĂ”ivad mĂ”ned asjad muutuda.

KokkuvĂ”tteks tahan meelde tuletada pĂ”hipunkti: saate ĂŒhendada sĂŒndmuste vooge (KStream) ja vĂ€rskenduste vooge (KTable) kohaliku oleku abil. Lisaks, kui kaasandmete maht pole liiga suur, saab kasutada objekti GlobalKTable. GlobalKTable kopeerib kĂ”ik osad iga Kafka Streams rakenduse sĂ”lme, tagades seelĂ€bi kĂ”igi andmete kĂ€ttesaadavuse sĂ”ltumata sellest, millisele osale vastab vĂ”ti.
Edasi vaatame Kafka Streams'i vÔimalust, mis vÔimaldab jÀlgida oleku muutusi ilma andmeid Kafka teemast tarbimata.
5.3.5. KĂŒsitav olek
Oleme juba teinud mitmeid olekuga seotud operatsioone ja alati oleme tulemused konsooli vÀljastanud (arenduse jaoks) vÔi salvestanud need teemasse (tööstuslikeks otstarveteks). Tulemuste teemasse kirjutamisel tuleb neid vaadata Kafka tarbija abil.
Nende teemade andmete lugemist vĂ”ib pidada materialiseeritud vaadete (materialized views) ĂŒheks vormiks. Meie jaoks vĂ”ime kasutada Wikipedia mÀÀratlust materialiseeritud vaate kohta: â⊠fĂŒĂŒsiline andmebaasi objekt, mis sisaldab pĂ€ringu tĂ€itmise tulemusi. NĂ€iteks vĂ”ib see olla kohaliku koopia kaugandmetest, vĂ”i ridade ja/vĂ”i veergude alamhulk tabelist vĂ”i ĂŒhendamise tulemustest, vĂ”i kogumisega saadud kokkuvĂ”tte tabel.â (https://en.wikipedia.org/wiki/Materialized_view).
Kafka Streams vÔimaldab samuti interaktiivsete pÀringute (interactive queries) tegemist olekuhoidlatest, mis annab vÔimaluse neid materialiseeritud vaateid otse lugeda. Oluline on mÀrkida, et pÀring olekuhoidlas on ainult lugemisoperatsioon. See tÀhendab, et saate olla kindel, et rakenduse andmete töötlemise kÀigus ei muuda te olekut ebatÀpsuseks.
Otseste pĂ€ringute vĂ”imalus olekuhoidlatesse on ÀÀrmiselt oluline. See tĂ€hendab, et saab luua rakendusi â informatiivseid paneele, ilma et peaksite esialgu andmeid Kafka tarbijalt hankima. Samuti suurendab see rakenduse efektiivsust, kuna andmeid ei pea uuesti salvestama:
- andmete lokaliseerimise tÔttu on neile kiire juurdepÀÀs;
- vÀlditakse andmete dubleerimist, kuna neid ei salvestata vÀlisesse andmehoidlasse.
Peamine, mida ma tahaksin, et te meeles peaksite: saate otse pĂ€rida olekust rakendusest. Selle vĂ”imaluse olulisust on raske ĂŒle hinnata. Selle asemel, et tarbida andmeid Kafka-st ja salvestada kirjeid andmebaasi rakenduse jaoks, saate teha pĂ€ringuid olekuhoidlatele sama tulemusega. Otsesed pĂ€ringud olekuhoidlatele tĂ€hendavad vĂ€hem koodi (tarbijat pole) ja vĂ€hem tarkvara (andmete tulemuste salvestamiseks pole vajadust andmebaasitabeli jĂ€rele).
Oleme selles peatĂŒkis kĂ€sitlenud suurt hulka teavet, seega teeme nĂŒĂŒd vĂ€ikese pausi interaktiivsete pĂ€ringute arutelust olekutehul. Ăra muretse: peatĂŒkis 9 loome lihtsa rakenduse - interaktiivsete pĂ€ringutega teabelehe. Sellel on mĂ”ned nĂ€ited kĂ€esolevast ja eelmise peatĂŒki interaktiivsetest pĂ€ringutest ja nende lisamisest Kafka Streams rakendustesse.
Elulookirjeldus
- KStream objektid kehastavad sĂŒndmuste vooge, mis sarnanevad andmebaasi lisamisele. KTable objektid kehastavad uuenduste vooge, nad sarnanevad rohkem andmebaasi uuendustega. KTable objekti suurus ei kasva, vanad kirjeldused asendatakse uutega.
- KTable objektid on vajalikud agregatsioonitegevuste jaoks.
- Aknategevuste abil saab agreggeeritud andmed jagada ajavahemike kaupa.
- GlobalKTable objektide kaudu on vÔimalik juurdepÀÀs viidatud andmetele rakenduse igas punktis sÔltumata osade jaotumisest.
- KStream, KTable ja GlobalKTable objektide vahel on lubatud ĂŒhendused.
Kuni praeguseni oleme keskendunud Kafka Streams rakenduste loomisele KStream kĂ”rgetasemelise DSL-i abil. Kuigi kĂ”rgetasemeline lĂ€henemine vĂ”imaldab luua puhtaid ja lĂŒhikesi programme, tĂ€hendab selle kasutamine teatud kompromissi. Töö DSL KStreamiga tĂ€hendab koodi lĂŒhendamise suurendamist kontrollitaseme alandamise arvelt. JĂ€rgmises peatĂŒkis vaatame madalama taseme API-d sĂ”lme töötlejatest ja proovime teisi kompromisse. Programmid saavad olema pikemad kui nad enne olid, aga meil on vĂ”imalus luua praktiliselt igasuguseid töötlejaid, mis meil kunagi vajalikud vĂ”ivad olla.
â Rohkem teavet raamatu kohta leiate aadressilt
â Habr Residentide jaoks allahindlus 25% kupongiga â Kafka Streams
â PĂ€rast trĂŒkitud raamatuga maksmist saadetakse elektrooniline raamat e-posti teel.
Allikas: habr.com
