{"id":35786,"date":"2019-10-31T22:06:19","date_gmt":"2019-10-31T19:06:19","guid":{"rendered":"https:\/\/prohoster.info\/blog\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\/"},"modified":"2019-10-31T22:06:19","modified_gmt":"2019-10-31T19:06:19","slug":"kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","status":"publish","type":"post","link":"https:\/\/prohoster.info\/et\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","title":{"rendered":"Raamat \"Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks\"","gt_translate_keys":[{"key":"rendered","format":"text"}]},"content":{"rendered":"<p><noindex><a rel=\"nofollow\" href=\"https:\/\/habr.com\/ru\/company\/piter\/blog\/457756\/\"><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/6f8bd2b31b87b0c760c1148515893c43.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/a><\/noindex> Tere, Habrojekte! See raamat sobib igale arendajale, kes soovib voogude t\u00f6\u00f6tlemisse s\u00fcveneda. Jaotatud programmeerimise m\u00f5istmine aitab paremini m\u00f5ista Kafka ja Kafka Streams'i. Oluline oleks tunda ka Kafka raamistikku, kuid see pole h\u00e4davajalik: ma selgitan teile k\u00f5ik, mida peate teadma. Kogenud Kafka arendajad, nagu ka algajad, saavad selle raamatu kaudu meisterdada huvitavaid voogude t\u00f6\u00f6tlemise rakendusi, kasutades Kafka Streams'i raamatukogu. Kesk- ja k\u00f5rgema taseme Java arendajad, kes on juba tuttavad selliste m\u00f5isted nagu serialiseerimine, \u00f5pivad, kuidas rakendada oma oskusi Kafka Streams'i rakenduste loomisel. Raamatu l\u00e4htekood on kirjutatud Java 8-s ja kasutab ulatuslikult Java 8 lambda-v\u00e4ljendite s\u00fcntaksit, seega tuleb kasuks oskused lambda-funktsioonide t\u00f6\u00f6tluseks (isegi teises programmeerimiskeeles).<br \/>\n<noindex><a rel=\"nofollow\" name=\"habracut\"><\/a><\/noindex><\/p>\n<h3>L\u00f5ik. 5.3. Aggregeerimine ja aknoperatsioonid<\/h3>\n<p>\nSelles jaotises viime l\u00e4bi uurimise k\u00f5ige lubavamatest osadest Kafka Streams'ist. Kuni praeguseni oleme k\u00e4sitlenud j\u00e4rgmisi Kafka Streams'i aspekte:<\/p>\n<ul>\n<li>t\u00f6\u00f6tlemise topoloogia loomine;<\/li>\n<li>oleku kasutamine voogude rakendustes;<\/li>\n<li>andmevoogude \u00fchendamine;<\/li>\n<li>erinevused s\u00fcndmusvoogude (KStream) ja uuenduste voogude (KTable) vahel.<\/li>\n<\/ul>\n<p>\nJ\u00e4rgmistes n\u00e4idetes koondame k\u00f5ik need elemendid kokku. Samuti tutvute aknoperatsioonidega \u2014 veel \u00fche suurep\u00e4rase v\u00f5imalusega voogude rakendustes. Meie esimene n\u00e4ide on lihtne aggregeerimine.<\/p>\n<h3>5.3.1. Aktsiate m\u00fc\u00fcgimahtude aggregeerimine t\u00f6\u00f6stusharu kaupa<\/h3>\n<p>\nAggregeerimine ja grupeerimine on eluliselt vajalikud t\u00f6\u00f6riistad voogude andmete t\u00f6\u00f6tlemisel. \u00dcksikute rekordite uurimine reaalajas osutub tihti ebapiisavaks. Andmetest lisainfo v\u00e4ljav\u00f5tmiseks on vajalikud gruppeerimine ja kombineerimine.<\/p>\n<p>Selles n\u00e4ites peate end intraday kaupleja kost\u00fc\u00fcmile panema, kes peab j\u00e4lgima aktsiate m\u00fc\u00fcgimahu mahtu mitmetes t\u00f6\u00f6stusharudes. Eriti huvitavad teid viis ettev\u00f5tet, mille aktsiate m\u00fc\u00fcgimaht on igas t\u00f6\u00f6stusharus suurim.<\/p>\n<p>Selle aggregeerimise jaoks on vaja mitmeid j\u00e4rgmised samme andmete vajaliku vormi t\u00f5lkimiseks (\u00fcldiselt \u00f6eldes).<\/p>\n<ol>\n<li>Loo source p\u00f5hined teemarikkal, mis avaldab aktsiate kauplemise kohta toorinfot. Me peame kuvama StockTransaction t\u00fc\u00fcpi objekti ShareVolume t\u00fc\u00fcbiks. Asi on selles, et StockTransaction objekt sisaldab m\u00fc\u00fcgimetainfot, kuid me vajame ainult andmeid m\u00fc\u00fcdud aktsiate arvu kohta.<\/li>\n<li>Klasifitseerime ShareVolume andmeid aktsia s\u00fcmbolite kaupa. P\u00e4rast s\u00fcmbolite j\u00e4rgi grupeerimist saame need andmed kokku p\u00f6\u00f6rata aktsiate m\u00fc\u00fcgi vahekokkuv\u00f5teteks. Tuleb m\u00e4rkida, et meetod KStream.groupBy tagastab KGroupedStream t\u00fc\u00fcpi eseme. KTable'i saame v\u00f5tta, kutsudes edasi KGroupedStream.reduce meetodit.<\/li>\n<\/ol>\n<p><\/p>\n<blockquote><p><b>Mis on KGroupedStream liides?<\/b><\/p>\n<p>Meetodid KStream.groupBy ja KStream.groupByKey tagastavad KGroupedStream t\u00fc\u00fcpi eseme. KGroupedStream on vahepealne esitus s\u00fcndmuste voost p\u00e4rast grupeerimist v\u00f5tmete j\u00e4rgi. See ei ole m\u00f5eldud otseseks t\u00f6\u00f6tamiseks. Selle asemel kasutatakse KGroupedStream'i kogumise toiminguteks, mille tulemuseks on alati KTable. Kuna kogumise tegemise tulemuseks on KTable ja see katab olekuhoidla, v\u00f5ivad mitte k\u00f5ik uuendused liikuda edasi torustiku kaudu.<\/p>\n<p>Meetod KTable.groupBy tagastab sarnase KGroupedTable'i \u2014 vahepealse esitluse v\u00e4rskenduste voost, mis on uuesti grupeeritud v\u00f5tme j\u00e4rgi.<\/p><\/blockquote>\n<p>\nTeeme v\u00e4ikese pausi ja vaatame joonist 5.9, mis n\u00e4itab, mida me saavutasime. See topoloogia peaks teile juba tuttav olema.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/9fd61317cde376362adcaeec72908919.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nVaadakem n\u00fc\u00fcd koodi selle topoloogia jaoks (seda saab leida failist src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (loetelu 5.2).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/f937287e448295fbd467c283ceca316a.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEsitatud koodi iseloomustab l\u00fchidus ja suur hulk tehteid, mis tehakse m\u00f5ne reaga. Meetodi builder.stream esimeses parameetris v\u00f5ite m\u00e4rgata midagi uut: loetletud t\u00fc\u00fcbi v\u00e4\u00e4rtus AutoOffsetReset.EARLIEST (on olemas ka LATEST), mis on m\u00e4\u00e4ratud meetodi Consumed.withOffsetResetPolicy abil. Selle loetletud t\u00fc\u00fcbi abil saate m\u00e4\u00e4rata nihkestrateegia iga KStream v\u00f5i KTable jaoks, see on eeliseks nihkemeetodi seadistusele konfiguratsioonis.<\/p>\n<blockquote><p><b>GroupByKey ja GroupBy<\/b><\/p>\n<p>KStream liideses on kaks meetodit kirje gruppeeringuks: GroupByKey ja GroupBy. M\u00f5lemad tagastavad KGroupedTable'i, seega v\u00f5ite k\u00fcsida: milline on nende vahe ja millal kumbagi kasutada?<\/p>\n<p>GroupByKey meetodit rakendatakse, kui KStream'i v\u00f5ti on juba olemas. Pealegi ei ole lipp 'n\u00f5uab \u00fcmberjaotamist' kunagi seadistatud.<\/p>\n<p>GroupBy meetod eeldab, et olete muutnud r\u00fchmitamise v\u00f5tmeid, mist\u00f5ttu on \u00fcmberjaotamise lipp seadistatud v\u00e4\u00e4rtusele true. GroupBy j\u00e4rel teostatavad \u00fchendused, agregatsioonid jne viivad automaatse \u00fcmberjaotamiseni.<br \/>\nKokkuv\u00f5tteks: kui v\u00e4hegi v\u00f5imalik, tuleks kasutada GroupByKey meetodit, mitte GroupBy'd.<\/p><\/blockquote>\n<p>\nKuna mapValues ja groupBy meetodite funktsioone on selge, vaatame sum() meetodit (see on leitav failist src\/main\/java\/bbejeck\/model\/ShareVolume.java) (loetelu 5.3).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/8e9a6f873594f9b5fef9a96c42353a61.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nShareVolume.sum meetod tagastab aktsiate m\u00fc\u00fcgimahtude vahepealse summa ning kogu arvutuste ahela tulemus on objekt KTable. N\u00fc\u00fcd m\u00f5istate, millist rolli KTable m\u00e4ngib. Kui ShareVolume objekte lisatakse, salvestatakse vastavas KTable objektis viimane asjakohane uuendus. On oluline mitte unustada, et k\u00f5ik uuendused kajastuvad eelnevas shareVolumeKTable'is, kuid mitte k\u00f5ik saadetakse edasi.<\/p>\n<p>Edasi, kasutades seda KTable'i, teeme agregatsiooni (m\u00fc\u00fcdud aktsiate arvu alusel), et saada viis ettev\u00f5tet, kellel on iga t\u00f6\u00f6stusharu suurimad aktsiate m\u00fc\u00fcgivoldid. Meie tegevus on sarnane esimese agregatsiooniga.<\/p>\n<ol>\n<li>Tehke veel \u00fcks groupBy operatsioon, et r\u00fchmitada eraldi ShareVolume objekte t\u00f6\u00f6stusharude j\u00e4rgi.<\/li>\n<li>Alustage ShareVolume objektide summeerimisega. Seekord on agregatsiooni objectiks fikseeritud suurusega prioriteedik\u00f6it. Sellises fikseeritud suurusega k\u00f6ites hoitakse ainult viit ettev\u00f5tet, kellel on k\u00f5ige rohkem m\u00fc\u00fcdud aktsiaid.<\/li>\n<li>Muuda eelneva punkti j\u00e4rjekorrad stringiv\u00e4\u00e4rtuseks ja tagasta viie k\u00f5ige paremini m\u00fc\u00fcdud aktsia p\u00f5hjal nende t\u00f6\u00f6stusharude j\u00e4rgi.<\/li>\n<li>Kanna tulemused stringivormis teemasse.<\/li>\n<\/ol>\n<p>\nJoonisel 5.10 on kujutatud andmevoogude topoloogia graafik. Nagu n\u00e4ete, on teise t\u00f6\u00f6tlemise ts\u00fckkel \u00fcsna lihtne.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/dabd1507eee267038edb7f8d76d8d8ae.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nN\u00fc\u00fcd, kui olete saanud selge arusaama selle teise t\u00f6\u00f6tlemise ts\u00fckli struktuurist, v\u00f5ite tutvuda selle l\u00e4htekoodiga (leiate selle failist src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (loetelu 5.4).<\/p>\n<p>Selles initsialisaatoris on muutuja fixedQueue. See on kohandatud objekt \u2014 adapter java.util.TreeSet'ile, mida kasutatakse, et j\u00e4lgida N suurimat m\u00fc\u00fcgihulka kahanevas j\u00e4rjestuses.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/488b072b0d91b81c925ca72291e69e48.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nTe olete juba tuttavad groupBy ja mapValues v\u00e4ljakutsetega, seega ei j\u00e4\u00e4me neile pikemalt peatuma (me kutsume KTable.toStream meetodit, kuna KTable.print meetod on vananenud). Kuid te pole veel n\u00e4inud KTable versiooni meetodist aggregate(), seega kulutame veidi aega selle arutamiseks.<\/p>\n<p>Nagu te meenute, erineb KTable sellest, et sama v\u00f5tmega kirjeid peetakse uuendusteks. KTable asendab vana kirje uuega. Aggregeerimine toimub sarnaselt: aggregeeritakse viimased kirjed sama v\u00f5tmega. Kirje saabudes lisatakse see FixedSizePriorityQueue klassi eksemplarisse koos summeerijaga (teine parameeter aggregate meetodi \u00fcleskutses), kuid kui sama v\u00f5tmega on juba olemas teine kirje, siis vana kirje eemaldatakse v\u00e4hendajaga (kolmas parameeter aggregate meetodi \u00fcleskutses).<\/p>\n<p>See k\u00f5ik t\u00e4hendab, et meie aggregeerija, FixedSizePriorityQueue, ei aggregeeri k\u00f5iki sama v\u00f5tmega v\u00e4\u00e4rtusi, vaid hoiab N k\u00f5ige enam m\u00fc\u00fcdud aktsiate liikide libisevat summat. Igas saabunud kirjes on \u00fcldine m\u00fc\u00fcdud aktsiate hulk seni. KTable annab teile teavet selle kohta, milliste ettev\u00f5tete aktsiaid m\u00fc\u00fcakse praegu rohkem, libisev aggregeerimine iga uuenduse puhul ei ole vajalik.<\/p>\n<p>Oleme \u00f5ppinud tegema kahte olulist asja:<\/p>\n<ul>\n<li>gruppima KTable's v\u00e4\u00e4rtusi nende \u00fchise v\u00f5tme j\u00e4rgi;<\/li>\n<li>teostama nendele r\u00fchmitatud v\u00e4\u00e4rtustele kasulikke toiminguid, nagu kokkuv\u00f5te ja aggregeerimine.<\/li>\n<\/ul>\n<p>\nNeid toiminguid osata on oluline, et m\u00f5ista andmete t\u00e4hendust, mis liigub l\u00e4bi Kafka Streams rakenduse, ja v\u00e4lja selgitada, millist teavet nad kannavad.<\/p>\n<p>Me oleme samuti kokku viinud m\u00f5ned v\u00f5tmekontseptsioonid, millest oleme varem selles raamatus r\u00e4\u00e4kinud. 4. peat\u00fckis r\u00e4\u00e4kisime, kui oluline on voogedastusrakenduse jaoks talletatud, lokaalne olek. Selles peat\u00fckis toodud esimene n\u00e4ide demonstreeris, miks lokaalne olek on nii oluline \u2013 see v\u00f5imaldab j\u00e4lgida, millist teavet olete juba n\u00e4inud. Lokaalne ligip\u00e4\u00e4s aitab v\u00e4ltida v\u00f5rgulisi viivitusi, muutes rakenduse tootlikumaks ja vigu taluvamaks.<\/p>\n<p>Iga koondamis- v\u00f5i agregatsioonitegevuse l\u00e4biviimisel tuleb t\u00e4psustada olekutehnika nime. Koondamis- ja agregatsioonitegevused tagastavad KTable'i koopia, samal ajal kui KTable kasutab olekutehnikat, et asendada vanad tulemused uutega. Nagu olete n\u00e4inud, ei saadeta kaugeltki k\u00f5iki v\u00e4rskendusi torust edasi, ja see on oluline, kuna agregatsioonitegevused on m\u00f5eldud l\u00f5pptulemuse saamiseks. Ilma lokaalse olekuta saadab KTable edasi k\u00f5ik agregatsiooni ja koondamise tulemused.<\/p>\n<p>J\u00e4rgmises vaates vaatame l\u00e4bi selliseid tegevusi nagu agregatsioon kindla ajavahemiku jooksul \u2013 nii nimetatud aktsioonitegevusi (windowing operations).<\/p>\n<h3>5.3.2. Aktsioonitegevused<\/h3>\n<p>\nEelnevas peat\u00fckis tutvusime \"libiseva\" koondamise ja agregatsiooniga. Rakendus tegi pidevat koondamist aktsiate m\u00fc\u00fcgimahu kohta, j\u00e4rgides viie k\u00f5ige rohkem m\u00fc\u00fcdud aktsia koondamist.<\/p>\n<p>M\u00f5nikord on sarnased pidevad agregatsioonid ja koondamised vajalikud. Kuid m\u00f5nikord tuleb teostada tegevusi ainult m\u00e4\u00e4ratud ajavahemiku jooksul. N\u00e4iteks v\u00e4lja arvutada, kui palju b\u00f6rsil tehinguid on tehtud konkreetse ettev\u00f5tte aktsiatega viimase 10 minuti jooksul. V\u00f5i kui palju kasutajaid kl\u00f5psas uuel reklaamb\u00e4nneril viimase 15 minuti jooksul. Rakendus v\u00f5ib neid tegevusi teha korduvalt, kuid tulemustega, mis on seotud ainult m\u00e4\u00e4ratud ajavahemikega (ajaliste akendega).<\/p>\n<h3>B\u00f6rsil tehingute arvestamine ostja j\u00e4rgi<\/h3>\n<p>\nJ\u00e4rgmises n\u00e4ites meie teema on b\u00f6rsil tehingute j\u00e4lgimine mitme kaupleja puhul \u2013 kas suurte ettev\u00f5tete v\u00f5i nutikate individuaalsete finantsteadlaste puhul.<\/p>\n<p>On olemas kaks v\u00f5imalikku p\u00f5hjust sarnase j\u00e4lgimise jaoks. \u00dcks neist on vajadus teada, mida ostavad\/m\u00fc\u00fcvad turuliidrid. Kui need suured m\u00e4ngijad ja kogenud investorid n\u00e4evad endale avanevaid v\u00f5imalusi, on m\u00f5istlik j\u00e4rgida nende strateegiat. Teine p\u00f5hjus on soov m\u00e4rgata v\u00f5imalikke ilmsi ebaseaduslikke tehingutega, mis kasutavad siseteavet. Selleks peate anal\u00fc\u00fcsima suurte m\u00fc\u00fcgivoolude ja oluliste pressiteadete korrelatsiooni.<\/p>\n<p>Selline j\u00e4lgimine h\u00f5lmab j\u00e4rgmisi etappe:<\/p>\n<ul>\n<li>loome voolu lugemiseks teemast stock-transactions;<\/li>\n<li>grupime sissetulevad kirjed ostja ID ja aktsia s\u00fcmboli j\u00e4rgi. GroupBy meetodi kutsumine tagastab KGroupedStream klassi eksemplari;<\/li>\n<li>KGroupedStream.windowedBy meetodi abil tagastatakse ajavahemiku piiramiseks andmevoog, mis v\u00f5imaldab teostada akna agregatsiooni. Olenemata akna t\u00fc\u00fcbist tagastatakse kas TimeWindowedKStream v\u00f5i SessionWindowedKStream;<\/li>\n<li>tehingute arvestamine agregatsiooni toimingu jaoks. Aken m\u00e4\u00e4rab, kas konkreetne kande arvesse v\u00f5tmine arvutamisel toimub;<\/li>\n<li>tulemuste salvestamine teemasse v\u00f5i nende v\u00e4ljastamine konsooli arenduse k\u00e4igus.<\/li>\n<\/ul>\n<p>\nRakenduse topoloogia on lihtne, kuid visuaalne esitus ei teeks kahju. Vaatame joonist 5.11.<\/p>\n<p>Edasi vaatleme akna toimingute funktsionaalsust ja vastavat koodi.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/67d9d8d8acb1370a7c7fe5cd9a8b6aa7.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Akna t\u00fc\u00fcbid<\/h3>\n<p>\nKafka Streams'is on kolm t\u00fc\u00fcpi aknaid:<\/p>\n<ul>\n<li>seanssi;<\/li>\n<li>\u201erulluvad\u201d (tumbling);<\/li>\n<li>libisevad\/\u201eh\u00fcppavad\u201d (sliding\/hopping).<\/li>\n<\/ul>\n<p>\nMille valida \u2014 s\u00f5ltub \u00e4rin\u00f5uetest. \u201eRulluvad\u201d ja \u201eh\u00fcppavad\u201d aknad piirduvad ajaga, samas kui seansside piirangud on seotud kasutajate tegevustega \u2014 seansi (seansside) kestus m\u00e4\u00e4ratakse ainuilt selle j\u00e4rgi, kui aktiivselt k\u00e4itub kasutaja. Peamine on mitte unustada, et k\u00f5ik akna t\u00fc\u00fcbid p\u00f5hinevad kannete kuup\u00e4eva\/aegade m\u00e4rgistel, mitte s\u00fcsteemi ajast.<\/p>\n<p>Edasi rakendame meie topoloogiat iga akna t\u00fc\u00fcbi korral. T\u00e4ielik kood toodakse ainult esimeses n\u00e4ites, teiste akna t\u00fc\u00fcpide jaoks ei muutu midagi peale akna toimingu t\u00fc\u00fcbi.<\/p>\n<h3>Seanssi aknad<\/h3>\n<p>\nSeanside aknad erinevad oluliselt k\u00f5igist teistest aknat\u00fc\u00fcpidest. Need on piiratud mitte niiv\u00f5rd ajaga, vaid kasutaja tegevusega (v\u00f5i selle \u00fcksuse tegevusega, mida soovite j\u00e4lgida). Seanside aknad eraldatakse tegevusetuse perioodide kaupa.<\/p>\n<p>Joonisel 5.12 illustreeritakse seanside akna m\u00f5istet. V\u00e4iksem seanss sulandub temast vasakul asuva seansi sisse. Paremal asuv seanss on eraldi, kuna j\u00e4rgneb pikale tegevusetuse perioodile. Seanside aknad p\u00f5hinevad kasutajate tegevusel, kuid kasutavad kuup\u00e4eva\/kellaaja m\u00e4rke, et m\u00e4\u00e4rata, mille seansi alla vastav kirje kuulub.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/d39a5db7af7d6aa2194802622b2b47fd.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Seanside aknade kasutamine b\u00f6rsitehingute j\u00e4lgimiseks<\/h3>\n<p>\nKasutame seanside aknaid b\u00f6rsitehingute teabe kogumiseks. Seanside akna rakendamine on n\u00e4idatud loendis 5.5 (mida v\u00f5ib leida failist src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKTableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/2dcbd9a36baec0e746aad165121451b3.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEnamik selle topoloogia operatsioone on teile juba tuttavad, nii et pole vajalik neid siin uuesti k\u00e4sitleda. Aga siin on ka m\u00f5ned uued elemendid, millest me praegu r\u00e4\u00e4gime.<\/p>\n<p>Iga groupBy operatsiooni puhul tehakse tavaliselt mingisugune agregatsioonioperatsioon (agregatsioon, kokkusurumine v\u00f5i arvu lugemine). Saame teostada kas kumulatiivset agregatsiooni jooksva kogusummaga v\u00f5i akna agregatsiooni, kus arvesse v\u00f5etakse kirjeid m\u00e4\u00e4ratud ajavahemiku jooksul.<\/p>\n<p>Loendi 5.5 kood loendab tehingute arvu seanside akende piires. Joonisel 5.13 anal\u00fc\u00fcsitakse neid toiminguid samm-sammult.<\/p>\n<p>Kasutades funktsiooni windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)), loome seanside akna 20-sekundilise tegevusetuse intervalliga ja 15-minutilise s\u00e4ilitamisintervalliga. 20-sekundiline tegevusetuse intervall t\u00e4hendab, et rakendus h\u00f5lmab igat kirjet, mis saabub 20 sekundi jooksul p\u00e4rast praeguse seansi l\u00f5ppu v\u00f5i algust, praegusesse (aktiivsesse) seanssi.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/9bd47b04698086872fd135b4c67eb938.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEdasi anname teada, millist agregatsiooni toimingut tuleb seanssi aknas teostada - antud juhul count. Kui sisenev kirje \u00fcletab tegevusetuse intervalli (m\u00f5lemalt poolt ajatempli), siis rakendus loob uue seansi. Salvestamise intervall t\u00e4hendab seansi hoidmist teatud aja jooksul ja v\u00f5imaldab hilinenud andmeid, mis \u00fcletavad seansi tegevusetuse perioodi, kuid v\u00f5ivad siiski liituda. Lisaks on uue seansi algus ja l\u00f5pp, mis tulenevad \u00fchendamisest, vastavalt varasemale ja hilisemale ajatempli.<\/p>\n<p>Vaadakem m\u00f5ningaid kirjeid count meetodist, et n\u00e4ha, kuidas seansid toimivad (tabel 5.1).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/ec04aae466d88c2d2349474069c8d541.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nKirjete saabumise korral otsime juba olemasolevaid seansse sama v\u00f5tme jaoks, mille l\u00f5pp aeg on v\u00e4hem kui praegune ajatempli - tegevusetuse intervall ja mille algusaeg on rohkem kui praegune ajatempli + tegevusetuse intervall. Selle j\u00e4rel \u00fchendatakse neli kirjet tabelist 5.1 \u00fchte seanssi j\u00e4rgmiselt.<\/p>\n<p>1. Esimesena saabub kirje 1, nii et algusaeg on sama, mis l\u00f5ppaeg ja on 00:00:00.<\/p>\n<p>2. J\u00e4rgmisena saabub kirje 2, ja me otsime seansse, mis l\u00f5pevad hiljemalt 23:59:55 ja algavad hiljemalt 00:00:35. Leiame kirje 1 ja \u00fchendame seansid 1 ja 2. V\u00f5tame seansi 1 algusaja (varasem) ja seansi 2 l\u00f5ppaaja (hilisem), nii et meie uus seanss algab 00:00:00 ja l\u00f5peb 00:00:15.<\/p>\n<p>3. Saabub kirje 3, me otsime seansse ajavahemikus 00:00:30 kuni 00:01:10 ja ei leia \u00fchtegi. Lisame teise seansi v\u00f5tmega 123-345-654,FFBE, algusega ja l\u00f5puga 00:00:50.<\/p>\n<p>4. Saabub kirje 4, ja me otsime seansse ajavahemikus 23:59:45 kuni 00:00:25. Seekord leitud m\u00f5lemad seansid - 1 ja 2. K\u00f5ik kolm seanssi \u00fchendatakse \u00fcheks, algusajaga 00:00:00 ja l\u00f5ppaeg 00:00:15.<\/p>\n<p>K\u00e4esolevas jaotises on m\u00f5ningaid olulisi n\u00fcansse, mida tuleks meeles pidada:<\/p>\n<ul>\n<li>seansid ei ole fikseeritud suurusega aknad. Seansi kestus m\u00e4\u00e4ratakse aktiveerituse j\u00e4rgi m\u00e4\u00e4ratud ajavahemikus;<\/li>\n<li>ajatempli andmed m\u00e4\u00e4ravad, kas s\u00fcndmus kuulub olemasolevasse seanssi v\u00f5i tegevusetuse intervalli.<\/li>\n<\/ul>\n<p>\nEdasi arutame j\u00e4rgmist akna t\u00fc\u00fcpi - \u201eveerevad\u201d aknad.<\/p>\n<h3>\u201eVeerevad\u201d aknad<\/h3>\n<p>\n\u201eVeerlevad\u201c (tumbling) aknad haaravad s\u00fcndmusi, mis toimuvad kindlas ajavahemikus. Kujutlege, et peate salvestama k\u00f5iki mingisuguse ettev\u00f5tte b\u00f6rsitehinguid iga 20 sekundi tagant, seega kogute k\u00f5ik s\u00fcndmused selle ajavahemiku jooksul. 20-sekundilise intervalli l\u00f5ppedes \u201eveerleb\u201c aken ja l\u00e4heb uude 20-sekundilise vaatlusperioodi. Joonis 5.14 illustreerib seda olukorda.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/e90e560d9ddda5e2e524b7c387ad9874.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nKuidas n\u00e4ha, on k\u00f5ik viimase 20 sekundi jooksul saabunud s\u00fcndmused aknas. Selle ajavahemiku l\u00f5ppedes luuakse uus aken.<\/p>\n<p>Loendis 5.6 on toodud kood, mis demonstreerib \u201eveerlevate\u201c akende kasutamist b\u00f6rsitehingute igakuiseks salvestamiseks iga 20 sekundi tagant (seda v\u00f5ib leida failist src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/e055bb1b288c7d500b64372fe3fbf064.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nSelle v\u00e4ikese muudatuse t\u00f5ttu meetodi TimeWindows.of kutsumisel on v\u00f5imalik kasutada \u201eveerlevat\u201c akent. Antud n\u00e4ites ei ole meetodi until() kutsumist, mist\u00f5ttu kasutatakse vaikimisi salvestusintervalli, mis on 24 tundi.<\/p>\n<p>L\u00f5puks on aeg liikuda viimase aknat\u00fc\u00fcbi juurde \u2013 \u201eh\u00fcppavad\u201c (hopping) aknad.<\/p>\n<h3>Libisevad (\u201eh\u00fcppavad\u201c) aknad<\/h3>\n<p>\nLibisevad\/\u201eh\u00fcppavad\u201c (sliding\/hopping) aknad on sarnased \u201eveerlevatele\u201c, kuid v\u00e4ikese erinevusega. Libisevad aknad ei oota ajavahemiku l\u00f5ppu enne uue akna loomist hiljutiste s\u00fcndmuste t\u00f6\u00f6tlemiseks. Nad k\u00e4ivitavad uued arvutused p\u00e4rast ooteaega, mis on l\u00fchem akna kestusest.<\/p>\n<p>Erinevuste illustreerimiseks \u201eveerlevate\u201c ja \u201eh\u00fcppavate\u201c akende vahel p\u00f6\u00f6rdume tagasi b\u00f6rsitehingute arvestuse n\u00e4ite juurde. Meie eesm\u00e4rk on endiselt arvestada tehingute arvu, kuid me ei soovi oodata kogu ajavahemikku enne arvu v\u00e4rskendamist. Selle asemel v\u00e4rskendame arvu l\u00fchemate ajavahemike tagant. N\u00e4iteks loeme endiselt tehingute arvu iga 20 sekundi tagant, kuid v\u00e4rskendame arvu iga 5 sekundi tagant, nagu on n\u00e4idatud joonisel 5.15. Selle tulemusena on meil kolm tulemuste akent, millel on kattuvad andmed.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/415c8cd9f2b60d453a1a01c3bc99331f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLoendis 5.7 on toodud kood libisevate akende seadistamiseks (seda v\u00f5ib leida failist src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/ab2d1a64380d256d3fb084e16597417c.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\n\u201eVeerev\u201d akent saab muuta \u201eh\u00fcppavaks\u201d funktsiooni advanceBy() kutsumisega. Antud n\u00e4ites on salvestusintervall 15 minutit.<\/p>\n<p>Selles osas olete n\u00e4inud, kuidas piirata tulemusi ajavahemikega. Eriti tasub meeles pidada kolm asja, mida sellelt, jaotiselt kaasa v\u00f5tta:<\/p>\n<ul>\n<li>seanss akna suurus ei ole piiratud aja jooksul, vaid kasutajate aktiivsusega;<\/li>\n<li>\u201eveerlevad\u201d aknad annavad \u00fclevaate s\u00fcndmustest antud ajavahemiku jooksul;<\/li>\n<li>\u201eh\u00fcppavate\u201d akende kestus on kindel, kuid need uuendatakse sageli ja v\u00f5ivad k\u00f5ikides akendes sisaldada kattuvaid kirjeid.<\/li>\n<\/ul>\n<p>\nSeej\u00e4rel \u00f5ppime, kuidas KTable tagasi KStreami muuta \u00fchendamiseks.<\/p>\n<h3>5.3.3. KStreami ja KTable'i \u00fchendamine<\/h3>\n<p>\nKapis 4 arutasime kahe KStreami objekti \u00fchendamist. N\u00fc\u00fcd \u00f5pime KTable'i ja KStreami \u00fchendamist. Seda v\u00f5ib vaja minna j\u00e4rgmise lihtsa p\u00f5hjuse t\u00f5ttu. KStream on kirjevoog, KTable aga kirjeuuenduste voog, kuid m\u00f5nikord v\u00f5ib olla vajalik lisada voogudele konteksti KTable'i uuenduste kaudu.<\/p>\n<p>V\u00f5tame andmed b\u00f6rsitehingute arvu kohta ja \u00fchendame need b\u00f6rsiuudistega vastavate t\u00f6\u00f6stusharude j\u00e4rgi. Siin on, mida tuleb teha, et saavutada seda olemasoleva koodiga.<\/p>\n<ol>\n<li>Muutke b\u00f6rsitehingute arvu KTable objekt KStreamiks, asendades v\u00f5tme nii, et see t\u00e4hendaks t\u00f6\u00f6stusharu, mis vastab antud aktsiate s\u00fcmbolile.<\/li>\n<li>Looge KTable objekt, mis loeb andmeid b\u00f6rsiuudiste teemast. See uus KTable klassifitseeritakse t\u00f6\u00f6stusharude j\u00e4rgi.<\/li>\n<li>\u00dchendage uudiste uuendused b\u00f6rsitehingute arvu teabega vastavates t\u00f6\u00f6stusharudes.<\/li>\n<\/ol>\n<p>\nN\u00fc\u00fcd vaatame, kuidas seda tegevuskava ellu viia.<\/p>\n<h3>KTable'i muutmine KStreamiks<\/h3>\n<p>\nKTable'i KStreamiks muutmiseks tuleb teha j\u00e4rgmist.<\/p>\n<ol>\n<li>Kutsuda v\u00e4lja meetod KTable.toStream().<\/li>\n<li>Kasutades meetodi KStream.map kutset, asendada v\u00f5ti t\u00f6\u00f6stusharu nimega ja seej\u00e4rel v\u00e4lja v\u00f5tta TransactionSummary objekti.<\/li>\n<\/ol>\n<p>\nSeome need toimingud ahelaks j\u00e4rgmiselt (koodi leiate failist src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (loend 5.8).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/0d43c2650f6e66e2816ed383da3a29c2.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nKuna me teeme operatsiooni KStream.map, toimub KStreami tagastatud eksemplari \u00fcmberjaotamine automaatselt, kui seda kasutatakse \u00fchenduses.<\/p>\n<p>Oleme l\u00f5petanud muundamisprotsessi, n\u00fc\u00fcd peame looma KTable objekti b\u00f6rsiuudiste lugemiseks.<\/p>\n<h3>KTable loomine b\u00f6rsiuudiste jaoks<\/h3>\n<p>\n\u00d5nneks piisab KTable objekti loomiseks \u00fchest koodireast (seda koodi saab leida failist src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (loetelu 5.9).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/6e83a393fdc9ab74fda4cbdddddb5213.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nOluline on m\u00e4rkida, et Serde objekte ei ole vaja m\u00e4\u00e4rata, kuna seadistustes kasutatakse stringide Serde-d. Samuti t\u00e4idab tabel t\u00e4nu EARLIEST loetlemise kasutamisele kirjeid alguses.<\/p>\n<p>N\u00fc\u00fcd saame liikuda viimase sammu juurde \u2014 \u00fchendamine.<\/p>\n<h3>Uudiste v\u00e4rskenduste ja tehingute arvu andmete \u00fchendamine<\/h3>\n<p>\n\u00dchenduse loomine ei ole keeruline. Kasutame vasakut \u00fchendust juhuks, kui vastavas t\u00f6\u00f6stusharus pole b\u00f6rsiuudiseid (vajalik kood on saadaval failis src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (loetelu 5.10).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/92ed70f98927d2f778ad14dbc2a5aa26.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nSee leftJoin operaator on \u00fcsna lihtne. Erinevalt 4. peat\u00fckist, kus kasutati JoinWindow meetodit, ei kasutata KStream-KTable \u00fchendamisel, kuna KTable'is on iga v\u00f5tme kohta ainult \u00fcks kirje. Selline \u00fchendamine ei ole ajaliselt piiratud: kirje kas on KTable'is v\u00f5i pole. Peamine j\u00e4reldus: KTable objektide abil saab KStreami rikastada harvemini uuendatud viidatud andmetega.<\/p>\n<p>N\u00fc\u00fcd vaatame t\u00f5husamat viisi KStreami s\u00fcndmuste rikastamiseks.<\/p>\n<h3>5.3.4. GlobalKTable objektid<\/h3>\n<p>\nKuidas te n\u00fc\u00fcd aru saite, on vaja rikastada s\u00fcndmuste vooge v\u00f5i lisada neile konteksti. 4. peat\u00fckis n\u00e4gite kahe KStream objekti \u00fchendamist, ja eelmises jaos KStreami ja KTable'i \u00fchendamist. K\u00f5igil nendel juhtudel on andmevoo \u00fcmberjaotamine vajalik, kui v\u00f5tmed muudetakse uue t\u00fc\u00fcbi v\u00f5i v\u00e4\u00e4rtuse jaoks. M\u00f5nikord toimub \u00fcmberjaotamine selgelt, m\u00f5nikord teeb Kafka Streams seda automaatselt. \u00dcmberjaotamine on vajalik, kuna v\u00f5tmed on muutunud ja kirjed peavad olema uuteks sektsioonideks, vastasel juhul ei saa \u00fchendust teha (sellest arutati 4. peat\u00fckis jaotises \"Andmete \u00fcmberjaotamine\" alajaos 4.2.4).<\/p>\n<h3>Uuesti jaotamine on tasuline<\/h3>\n<p>\nUuesti jaotamine n\u00f5uab kulutusi \u2014 t\u00e4iendavaid ressursikulu vahepealsete teemade loomisele, duplikaatandmete s\u00e4ilitamisele veel \u00fches teemas; see toob kaasa ka viivituste suurenemise seoses andmete lugemise ja kirjutamisega sellesse teemasse. Lisaks, kui on vajalik teostada \u00fchendust rohkem kui \u00fche aspekti v\u00f5i m\u00f5\u00f5tme alusel, tuleb \u00fchendused organiseerida ahelana, kuvada kirjed uute v\u00f5tmete j\u00e4rgi ja uuesti viia l\u00e4bi uuesti jaotamise protsess.<\/p>\n<h3>\u00dchendamine v\u00e4iksemate andmekogumitega<\/h3>\n<p>\nM\u00f5nes olukorras on viidatud andmete maht, millega plaanitakse \u00fchendust, suhteliselt v\u00e4ike, nii et nende t\u00e4iendavad koopiad mahuvad igasse s\u00f5lme kohapeal. Selle jaoks on Kafka Streams'is olemas klass GlobalKTable.<\/p>\n<p>GlobalKTable'i eksemplarid on ainulaadsed, kuna rakendus replitseerib k\u00f5ik andmed igasse s\u00f5lmesse. Kuna igas s\u00f5lmes on k\u00f5ik andmed, ei ole vaja s\u00fcndmuste voogu jaotada viidatud andmete v\u00f5tme alusel, et need oleksid kergesti ligip\u00e4\u00e4setavad k\u00f5ikidele sektsioonidele. GlobalKTable objektidega on samuti v\u00f5imalik teostada v\u00f5tmeta \u00fchendusi. Naaseme tagasi \u00fche eelneva n\u00e4ite juurde, et demonstreerida seda v\u00f5imalust.<\/p>\n<h3>KStream objektide \u00fchendamine GlobalKTable objektidega<\/h3>\n<p>\nAlajaotises 5.3.2 viisime l\u00e4bi akna agregatsiooni b\u00f6rsitehingute \u00fcle ostjate l\u00f5ikes. Selle agregatsiooni tulemused n\u00e4gid v\u00e4lja umbes j\u00e4rgmised:<\/p>\n<pre><code class=\"plaintext\">{customerId='074-09-3705', stockTicker='GUTM'}, 17\n{customerId='037-34-5184', stockTicker='CORK'}, 16<\/code><\/pre>\n<p>\nKuigi need tulemused vastasid seatud eesm\u00e4rgile, oleks mugavam, kui seal oleks ka kliendi nimi ja ettev\u00f5tte t\u00e4ielik nimi. Kliendi nime ja ettev\u00f5tte nime lisamiseks on v\u00f5imalik teha tavalisi \u00fchendusi, kuid see n\u00f5uab kahe v\u00f5tme kaardistamist ja uuesti jaotamist. GlobalKTable abil saab selliseid toimingute kulutusi v\u00e4ltida.<\/p>\n<p>Selleks kasutame objekti countStream n\u00e4iteks 5.11 (vastav kood asub failis src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java), \u00fchendades selle kahe GlobalKTable objektiga.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/fc4d91bbe062ceb94f5650224840b81e.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nOleme 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.<\/p>\n<p>J\u00e4rgmine etapp on kahe GlobalKTable instantsi deklareerimine (toodud kood on saadaval failis src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (loetelu 5.12).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/db3918895f174c8cfb5fc927b55f87a1.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p>Pange t\u00e4hele, et teema nimed on m\u00e4\u00e4ratletud loendatavate t\u00fc\u00fcpide abil.<\/p>\n<p>N\u00fc\u00fcd, kui oleme k\u00f5ik komponendid ette valmistanud, j\u00e4\u00e4b kirjutada kood \u00fchendamiseks (mida saab leida failist src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (loetelu 5.13).<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/799360cc99f1920c190a61fd4685d4ff.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nKuigi selles koodis on kaks \u00fchendust, on need korraldatud ahelaks, kuna eraldi ei kasutata \u00fchtegi nende tulemusi. Tulemused v\u00e4ljastatakse kogu operatsiooni l\u00f5pus.<\/p>\n<p>K\u00e4ivitades \u00fclaltoodud \u00fchenduse operatsiooni, saate j\u00e4rgmiste s\u00f5numite tulemusi:<\/p>\n<pre><code class=\"plaintext\">{customer='Barney, Smith' company=\"Exxon\", transactions= 17}<\/code><\/pre>\n<p>\nSisu pole muutunud, kuid need tulemused n\u00e4evad v\u00e4lja arusaadavamad.<\/p>\n<p>Kui arvestada peat\u00fckki 4, olete juba n\u00e4inud mitmeid t\u00fc\u00fcpe \u00fchendusi tegevuses. Need on loetletud tabelis 5.2. See tabel kajastab \u00fchenduste v\u00f5imalusi, mis on seotud Kafka Streams versiooniga 1.0.0; tulevastes v\u00e4ljalasemetes v\u00f5ivad m\u00f5ned asjad muutuda.<\/p>\n<p><img decoding=\"async\" alt=\"Raamat &quot;Kafka Streams tegevuses. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks&quot;\" src=\"\/wp-content\/uploads\/8e4cf35c64a8431bda43a5e748de275f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nKokkuv\u00f5tteks tahan meelde tuletada p\u00f5hipunkti: saate \u00fchendada s\u00fcndmuste vooge (KStream) ja v\u00e4rskenduste vooge (KTable) kohaliku oleku abil. Lisaks, kui kaasandmete maht pole liiga suur, saab kasutada objekti GlobalKTable. GlobalKTable kopeerib k\u00f5ik osad iga Kafka Streams rakenduse s\u00f5lme, tagades seel\u00e4bi k\u00f5igi andmete k\u00e4ttesaadavuse s\u00f5ltumata sellest, millisele osale vastab v\u00f5ti.<\/p>\n<p>Edasi vaatame Kafka Streams'i v\u00f5imalust, mis v\u00f5imaldab j\u00e4lgida oleku muutusi ilma andmeid Kafka teemast tarbimata.<\/p>\n<h3>5.3.5. K\u00fcsitav olek<\/h3>\n<p>\nOleme juba teinud mitmeid olekuga seotud operatsioone ja alati oleme tulemused konsooli v\u00e4ljastanud (arenduse jaoks) v\u00f5i salvestanud need teemasse (t\u00f6\u00f6stuslikeks otstarveteks). Tulemuste teemasse kirjutamisel tuleb neid vaadata Kafka tarbija abil.<\/p>\n<p>Nende teemade andmete lugemist v\u00f5ib pidada materialiseeritud vaadete (materialized views) \u00fcheks vormiks. Meie jaoks v\u00f5ime kasutada Wikipedia m\u00e4\u00e4ratlust materialiseeritud vaate kohta: \u201e\u2026 f\u00fc\u00fcsiline andmebaasi objekt, mis sisaldab p\u00e4ringu t\u00e4itmise tulemusi. N\u00e4iteks v\u00f5ib see olla kohaliku koopia kaugandmetest, v\u00f5i ridade ja\/v\u00f5i veergude alamhulk tabelist v\u00f5i \u00fchendamise tulemustest, v\u00f5i kogumisega saadud kokkuv\u00f5tte tabel.\u201d (https:\/\/en.wikipedia.org\/wiki\/Materialized_view).<\/p>\n<p>Kafka Streams v\u00f5imaldab samuti interaktiivsete p\u00e4ringute (interactive queries) tegemist olekuhoidlatest, mis annab v\u00f5imaluse neid materialiseeritud vaateid otse lugeda. Oluline on m\u00e4rkida, et p\u00e4ring olekuhoidlas on ainult lugemisoperatsioon. See t\u00e4hendab, et saate olla kindel, et rakenduse andmete t\u00f6\u00f6tlemise k\u00e4igus ei muuda te olekut ebat\u00e4psuseks.<\/p>\n<p>Otseste p\u00e4ringute v\u00f5imalus olekuhoidlatesse on \u00e4\u00e4rmiselt oluline. See t\u00e4hendab, et saab luua rakendusi \u2014 informatiivseid paneele, ilma et peaksite esialgu andmeid Kafka tarbijalt hankima. Samuti suurendab see rakenduse efektiivsust, kuna andmeid ei pea uuesti salvestama:<\/p>\n<ul>\n<li>andmete lokaliseerimise t\u00f5ttu on neile kiire juurdep\u00e4\u00e4s;<\/li>\n<li>v\u00e4lditakse andmete dubleerimist, kuna neid ei salvestata v\u00e4lisesse andmehoidlasse.<\/li>\n<\/ul>\n<p>\nPeamine, mida ma tahaksin, et te meeles peaksite: saate otse p\u00e4rida olekust rakendusest. Selle v\u00f5imaluse olulisust on raske \u00fcle hinnata. Selle asemel, et tarbida andmeid Kafka-st ja salvestada kirjeid andmebaasi rakenduse jaoks, saate teha p\u00e4ringuid olekuhoidlatele sama tulemusega. Otsesed p\u00e4ringud olekuhoidlatele t\u00e4hendavad v\u00e4hem koodi (tarbijat pole) ja v\u00e4hem tarkvara (andmete tulemuste salvestamiseks pole vajadust andmebaasitabeli j\u00e4rele).<\/p>\n<p>Oleme selles peat\u00fckis k\u00e4sitlenud suurt hulka teavet, seega teeme n\u00fc\u00fcd v\u00e4ikese pausi interaktiivsete p\u00e4ringute arutelust olekutehul. \u00c4ra muretse: peat\u00fckis 9 loome lihtsa rakenduse - interaktiivsete p\u00e4ringutega teabelehe. Sellel on m\u00f5ned n\u00e4ited k\u00e4esolevast ja eelmise peat\u00fcki interaktiivsetest p\u00e4ringutest ja nende lisamisest Kafka Streams rakendustesse.<\/p>\n<h3>Elulookirjeldus<\/h3>\n<p><\/p>\n<ul>\n<li>KStream objektid kehastavad s\u00fcndmuste vooge, mis sarnanevad andmebaasi lisamisele. KTable objektid kehastavad uuenduste vooge, nad sarnanevad rohkem andmebaasi uuendustega. KTable objekti suurus ei kasva, vanad kirjeldused asendatakse uutega.<\/li>\n<li>KTable objektid on vajalikud agregatsioonitegevuste jaoks.<\/li>\n<li>Aknategevuste abil saab agreggeeritud andmed jagada ajavahemike kaupa.<\/li>\n<li>GlobalKTable objektide kaudu on v\u00f5imalik juurdep\u00e4\u00e4s viidatud andmetele rakenduse igas punktis s\u00f5ltumata osade jaotumisest.<\/li>\n<li>KStream, KTable ja GlobalKTable objektide vahel on lubatud \u00fchendused.<\/li>\n<\/ul>\n<p>\nKuni praeguseni oleme keskendunud Kafka Streams rakenduste loomisele KStream k\u00f5rgetasemelise DSL-i abil. Kuigi k\u00f5rgetasemeline l\u00e4henemine v\u00f5imaldab luua puhtaid ja l\u00fchikesi programme, t\u00e4hendab selle kasutamine teatud kompromissi. T\u00f6\u00f6 DSL KStreamiga t\u00e4hendab koodi l\u00fchendamise suurendamist kontrollitaseme alandamise arvelt. J\u00e4rgmises peat\u00fckis vaatame madalama taseme API-d s\u00f5lme t\u00f6\u00f6tlejatest ja proovime teisi kompromisse. Programmid saavad olema pikemad kui nad enne olid, aga meil on v\u00f5imalus luua praktiliselt igasuguseid t\u00f6\u00f6tlejaid, mis meil kunagi vajalikud v\u00f5ivad olla.<\/p>\n<p>\u2192 Rohkem teavet raamatu kohta leiate aadressilt <noindex><a rel=\"nofollow\" href=\"https:\/\/www.piter.com\/collection\/best\/product\/kafka-streams-v-deystvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\">kirjastuse veebisaidil<\/a><\/noindex><\/p>\n<p>\u2192 Habr Residentide jaoks allahindlus 25% kupongiga \u2014 <b>Kafka Streams<\/b><\/p>\n<p>\u2192 P\u00e4rast tr\u00fckitud raamatuga maksmist saadetakse elektrooniline raamat e-posti teel.<br \/>\n<br \/>Allikas: <a content=\"nofollow\" rel=\"nofollow\" href=\"https:\/\/habr.com\/ru\/company\/piter\/blog\/457756\/\">habr.com<\/a><\/p>","protected":false,"gt_translate_keys":[{"key":"rendered","format":"html"}]},"excerpt":{"rendered":"<p>\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440\u043e\u0436\u0438\u0442\u0435\u043b\u0438! \u042d\u0442\u0430 \u043a\u043d\u0438\u0433\u0430 \u043f\u043e\u0434\u043e\u0439\u0434\u0435\u0442 \u0434\u043b\u044f \u043b\u044e\u0431\u043e\u0433\u043e \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0447\u0438\u043a\u0430, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0445\u043e\u0447\u0435\u0442 \u0440\u0430\u0437\u043e\u0431\u0440\u0430\u0442\u044c\u0441\u044f \u0432 \u043f\u043e\u0442\u043e\u043a\u043e\u0432\u043e\u0439 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0435. \u041f\u043e\u043d\u0438\u043c\u0430\u043d\u0438\u0435 \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u043d\u043e\u0433\u043e \u043f\u0440\u043e\u0433\u0440\u0430\u043c\u043c\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u043f\u043e\u043c\u043e\u0436\u0435\u0442 \u043b\u0443\u0447\u0448\u0435 \u0438\u0437\u0443\u0447\u0438\u0442\u044c Kafka \u0438 Kafka Streams. \u0411\u044b\u043b\u043e \u0431\u044b \u043d\u0435\u043f\u043b\u043e\u0445\u043e \u0437\u043d\u0430\u0442\u044c \u0438 \u0441\u0430\u043c \u0444\u0440\u0435\u0439\u043c\u0432\u043e\u0440\u043a Kafka, \u043d\u043e \u044d\u0442\u043e \u043d\u0435 \u043e\u0431\u044f\u0437\u0430\u0442\u0435\u043b\u044c\u043d\u043e: \u044f \u0440\u0430\u0441\u0441\u043a\u0430\u0436\u0443 \u0432\u0430\u043c \u0432\u0441\u0435, \u0447\u0442\u043e \u043d\u0443\u0436\u043d\u043e. \u041e\u043f\u044b\u0442\u043d\u044b\u0435 \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0447\u0438\u043a\u0438 Kafka, \u043a\u0430\u043a \u0438 \u043d\u043e\u0432\u0438\u0447\u043a\u0438, \u0431\u043b\u0430\u0433\u043e\u0434\u0430\u0440\u044f \u044d\u0442\u043e\u0439 \u043a\u043d\u0438\u0433\u0435 \u043e\u0441\u0432\u043e\u044f\u0442 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u0435 \u0438\u043d\u0442\u0435\u0440\u0435\u0441\u043d\u044b\u0445 \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u0439 [&hellip;]<\/p>\n","protected":false,"gt_translate_keys":[{"key":"rendered","format":"html"}]},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[688],"tags":[],"class_list":["post-35786","post","type-post","status-publish","format-standard","hentry","category-administrirovanie"],"aioseo_notices":[],"aioseo_head":"\n\t\t<!-- All in One SEO 5.0.1.1 - aioseo.com -->\n\t<meta name=\"robots\" content=\"max-image-preview:large\" \/>\n\t<meta name=\"author\" content=\"Yuri Gagarin\"\/>\n\t<link rel=\"canonical\" href=\"https:\/\/prohoster.info\/et\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\" \/>\n\t<meta name=\"generator\" content=\"All in One SEO (AIOSEO) 5.0.1.1\" \/>\n\t\t<meta property=\"og:locale\" content=\"et_EE\" \/>\n\t\t<meta property=\"og:site_name\" content=\"ProHoster | \u041a\u0443\u043f\u0438\u0442\u044c \u043d\u0430\u0434\u0435\u0436\u043d\u044b\u0439 \u0445\u043e\u0441\u0442\u0438\u043d\u0433 \u0434\u043b\u044f \u0441\u0430\u0439\u0442\u043e\u0432 \u0441 \u0437\u0430\u0449\u0438\u0442\u043e\u0439 \u043e\u0442 DDoS, VPS VDS \u0441\u0435\u0440\u0432\u0435\u0440\u044b\" \/>\n\t\t<meta property=\"og:type\" content=\"article\" \/>\n\t\t<meta property=\"og:title\" content=\"\ud83e\udd47\u041a\u043d\u0438\u0433\u0430 \u00abKafka Streams \u0432 \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0438. \u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f \u0438 \u043c\u0438\u043a\u0440\u043e\u0441\u0435\u0440\u0432\u0438\u0441\u044b \u0434\u043b\u044f \u0440\u0430\u0431\u043e\u0442\u044b \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u0438\u00bb | ProHoster\" \/>\n\t\t<meta property=\"og:url\" content=\"https:\/\/prohoster.info\/et\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\" \/>\n\t\t<meta property=\"og:image\" content=\"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg\" \/>\n\t\t<meta property=\"og:image:secure_url\" content=\"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg\" \/>\n\t\t<meta property=\"og:image:width\" content=\"350\" \/>\n\t\t<meta property=\"og:image:height\" content=\"350\" \/>\n\t\t<meta property=\"article:published_time\" content=\"2019-10-31T19:06:19+00:00\" \/>\n\t\t<meta property=\"article:modified_time\" content=\"2019-10-31T19:06:19+00:00\" \/>\n\t\t<meta property=\"article:publisher\" content=\"https:\/\/www.facebook.com\/prohoster\" \/>\n\t\t<meta property=\"article:author\" content=\"https:\/\/www.facebook.com\/prohoster\" \/>\n\t\t<!-- All in One SEO -->\n\n","aioseo_head_json":{"title":"\ud83e\udd47Raamat \u201eKafka Streams tegutsemisel. Rakendused ja mikroteenused reaalajas t\u00f6\u00f6tamiseks\u201c | ProHoster","description":"","canonical_url":"https:\/\/prohoster.info\/et\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","robots":"max-image-preview:large","keywords":"","webmasterTools":{"miscellaneous":""},"schema":null,"og:locale":"et_EE","og:site_name":"ProHoster | \u041a\u0443\u043f\u0438\u0442\u044c \u043d\u0430\u0434\u0435\u0436\u043d\u044b\u0439 \u0445\u043e\u0441\u0442\u0438\u043d\u0433 \u0434\u043b\u044f \u0441\u0430\u0439\u0442\u043e\u0432 \u0441 \u0437\u0430\u0449\u0438\u0442\u043e\u0439 \u043e\u0442 DDoS, VPS VDS \u0441\u0435\u0440\u0432\u0435\u0440\u044b","og:type":"article","og:title":"\ud83e\udd47\u041a\u043d\u0438\u0433\u0430 \u00abKafka Streams \u0432 \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0438. \u041f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f \u0438 \u043c\u0438\u043a\u0440\u043e\u0441\u0435\u0440\u0432\u0438\u0441\u044b \u0434\u043b\u044f \u0440\u0430\u0431\u043e\u0442\u044b \u0432 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c \u0432\u0440\u0435\u043c\u0435\u043d\u0438\u00bb | ProHoster","og:url":"https:\/\/prohoster.info\/et\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","og:image":"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg","og:image:secure_url":"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg","og:image:width":350,"og:image:height":350,"article:published_time":"2019-10-31T19:06:19+00:00","article:modified_time":"2019-10-31T19:06:19+00:00","article:publisher":"https:\/\/www.facebook.com\/prohoster","article:author":"https:\/\/www.facebook.com\/prohoster"},"aioseo_meta_data":{"post_id":"35786","title":null,"description":null,"keywords":null,"keyphrases":null,"primary_term":null,"canonical_url":null,"og_title":null,"og_description":null,"og_object_type":"default","og_image_type":"default","og_image_url":null,"og_image_width":null,"og_image_height":null,"og_image_custom_url":null,"og_image_custom_fields":null,"og_video":null,"og_custom_url":null,"og_article_section":null,"og_article_tags":null,"twitter_use_og":false,"twitter_card":"default","twitter_image_type":"default","twitter_image_url":null,"twitter_image_custom_url":null,"twitter_image_custom_fields":null,"twitter_title":null,"twitter_description":null,"schema":{"blockGraphs":[],"customGraphs":[],"default":{"data":{"Article":[],"Course":[],"Dataset":[],"FAQPage":[],"Movie":[],"Person":[],"Product":[],"ProductReview":[],"Car":[],"Recipe":[],"Service":[],"SoftwareApplication":[],"WebPage":[]},"graphName":"","isEnabled":true},"graphs":[]},"schema_type":null,"schema_type_options":null,"pillar_content":false,"robots_default":true,"robots_noindex":false,"robots_noarchive":false,"robots_nosnippet":false,"robots_nofollow":false,"robots_noimageindex":false,"robots_noodp":false,"robots_notranslate":false,"robots_max_snippet":null,"robots_max_videopreview":null,"robots_max_imagepreview":"large","priority":null,"frequency":null,"local_seo":null,"seo_analyzer_scan_date":"2026-01-22 00:45:19","breadcrumb_settings":null,"limit_modified_date":false,"reviewed_by":null,"ai":null,"created":"2021-03-01 01:56:32","updated":"2026-01-22 00:45:19","focus_keyword":null,"additional_keywords":null,"truseo_locale":null},"gt_translate_keys":[{"key":"link","format":"url"}],"_links":{"self":[{"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/posts\/35786","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/comments?post=35786"}],"version-history":[{"count":0,"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/posts\/35786\/revisions"}],"wp:attachment":[{"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/media?parent=35786"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/categories?post=35786"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/prohoster.info\/et\/wp-json\/wp\/v2\/tags?post=35786"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}