{"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\/pl\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","title":{"rendered":"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d","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=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/6f8bd2b31b87b0c760c1148515893c43.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/a><\/noindex> Cze\u015b\u0107, Habrzytelnicy! Ta ksi\u0105\u017cka jest odpowiednia dla ka\u017cdego dewelopera, kt\u00f3ry chce zrozumie\u0107 przetwarzanie strumieniowe. Zrozumienie programowania rozproszonego pomo\u017ce lepiej pozna\u0107 Kafka i Kafka Streams. By\u0142oby mi\u0142o zna\u0107 same ramy Kafka, ale to nie jest konieczne: opowiem Wam wszystko, co potrzebne. Do\u015bwiadczeni programi\u015bci Kafka, jak i nowicjusze, dzi\u0119ki tej ksi\u0105\u017cce opanuj\u0105 tworzenie interesuj\u0105cych aplikacji do przetwarzania strumieniowego za pomoc\u0105 biblioteki Kafka Streams. Programi\u015bci Java \u015bredniego i wy\u017cszego poziomu, ju\u017c przyzwyczajeni do takich poj\u0119\u0107 jak serializacja, naucz\u0105 si\u0119 stosowa\u0107 swoje umiej\u0119tno\u015bci do tworzenia aplikacji Kafka Streams. Kod \u017ar\u00f3d\u0142owy ksi\u0105\u017cki napisany jest w Java 8 i intensywnie wykorzystuje sk\u0142adni\u0119 wyra\u017ce\u0144 lambda Java 8, wi\u0119c umiej\u0119tno\u015b\u0107 pracy z funkcjami lambda (nawet w innym j\u0119zyku programowania) przyda si\u0119.<br \/>\n<noindex><a rel=\"nofollow\" name=\"habracut\"><\/a><\/noindex><\/p>\n<h3>Fragment. 5.3. Agregacja i operacje okienne<\/h3>\n<p>\nW tej sekcji przejdziemy do badania najbardziej obiecuj\u0105cych cz\u0119\u015bci Kafka Streams. Jak dot\u0105d przyjrzeli\u015bmy si\u0119 nast\u0119puj\u0105cym aspektom Kafka Streams:<\/p>\n<ul>\n<li>tworzenie topologii przetwarzania;<\/li>\n<li>u\u017cywanie stanu w aplikacjach strumieniowych;<\/li>\n<li>wykonywanie po\u0142\u0105cze\u0144 strumieni danych;<\/li>\n<li>r\u00f3\u017cnice mi\u0119dzy strumieniami wydarze\u0144 (KStream) a strumieniami aktualizacji (KTable).<\/li>\n<\/ul>\n<p>\nW nast\u0119pnych przyk\u0142adach po\u0142\u0105czymy wszystkie te elementy w ca\u0142o\u015b\u0107. Ponadto zapoznacie si\u0119 z operacjami okiennymi \u2014 jeszcze jedn\u0105 wspania\u0142\u0105 mo\u017cliwo\u015bci\u0105 aplikacji strumieniowych. Naszym pierwszym przyk\u0142adem b\u0119dzie prosta agregacja.<\/p>\n<h3>5.3.1. Agregacja wolumenu sprzeda\u017cy akcji wed\u0142ug bran\u017c przemys\u0142u<\/h3>\n<p>\nAgregacja i grupowanie to niezb\u0119dne narz\u0119dzia w pracy z danymi strumieniowymi. Badanie pojedynczych wpis\u00f3w w miar\u0119 ich pojawiania si\u0119 cz\u0119sto okazuje si\u0119 niewystarczaj\u0105ce. Aby wydoby\u0107 dodatkowe informacje z danych, konieczne jest ich grupowanie i \u0142\u0105czenie.<\/p>\n<p>W tym przyk\u0142adzie przyjmiesz rol\u0119 intraday tradera, kt\u00f3ry musi \u015bledzi\u0107 wolumeny sprzeda\u017cy akcji firm w kilku bran\u017cach przemys\u0142u. W szczeg\u00f3lno\u015bci interesuje ci\u0119 pi\u0119\u0107 firm z najwi\u0119kszym wolumenem sprzeda\u017cy akcji w ka\u017cdej z bran\u017c przemys\u0142u.<\/p>\n<p>Aby przeprowadzi\u0107 tak\u0105 agregacj\u0119, wymagane b\u0119d\u0105 nast\u0119puj\u0105ce kroki przetwarzania danych w odpowiedni spos\u00f3b (m\u00f3wi\u0105c og\u00f3lnie).<\/p>\n<ol>\n<li>Utw\u00f3rz \u017ar\u00f3d\u0142o na podstawie tematu, publikuj\u0105ce surowe informacje o handlu akcjami. Musimy przekszta\u0142ci\u0107 obiekt typu StockTransaction w obiekt typu ShareVolume. Problem polega na tym, \u017ce obiekt StockTransaction zawiera metadane sprzeda\u017cy, a potrzebne s\u0105 nam tylko dane o liczbie sprzedawanych akcji.<\/li>\n<li>Grupuj dane ShareVolume wed\u0142ug symboli akcji. Po grupowaniu wed\u0142ug symboli mo\u017cna zredukowa\u0107 te dane do po\u015brednich sum wolumen\u00f3w sprzeda\u017cy akcji. Warto zauwa\u017cy\u0107, \u017ce metoda KStream.groupBy zwraca instancj\u0119 typu KGroupedStream. Mo\u017cna uzyska\u0107 instancj\u0119 KTable, wywo\u0142uj\u0105c nast\u0119pnie metod\u0119 KGroupedStream.reduce.<\/li>\n<\/ol>\n<p><\/p>\n<blockquote><p><b>Czym jest interfejs KGroupedStream<\/b><\/p>\n<p>Metody KStream.groupBy i KStream.groupByKey zwracaj\u0105 instancj\u0119 KGroupedStream. KGroupedStream jest po\u015bredni\u0105 reprezentacj\u0105 strumienia zdarze\u0144 po grupowaniu wed\u0142ug kluczy. Nie jest przeznaczony do bezpo\u015bredniej pracy z nim. Zamiast tego KGroupedStream s\u0142u\u017cy do operacji agregacji, kt\u00f3rych wynikiem zawsze jest KTable. Poniewa\u017c wynikiem operacji agregacji jest KTable i wykorzystuje si\u0119 w nich pami\u0119\u0107 stanu, mo\u017cliwe, \u017ce nie wszystkie aktualizacje w wyniku s\u0105 przesy\u0142ane dalej w potoku.<\/p>\n<p>Metoda KTable.groupBy zwraca analogiczny KGroupedTable \u2014 po\u015bredni\u0105 reprezentacj\u0119 strumienia aktualizacji, przegrupowanych wed\u0142ug klucza.<\/p><\/blockquote>\n<p>\nZr\u00f3bmy ma\u0142\u0105 przerw\u0119 i sp\u00f3jrzmy na rys. 5.9, na kt\u00f3rym pokazano, co osi\u0105gn\u0119li\u015bmy. Ta topologia powinna by\u0107 Wam ju\u017c dobrze znana.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/9fd61317cde376362adcaeec72908919.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nPrzyjrzyjmy si\u0119 teraz kodowi dla tej topologii (mo\u017cna go znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (listing 5.2).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/f937287e448295fbd467c283ceca316a.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nPodany kod r\u00f3\u017cni si\u0119 zwi\u0119z\u0142o\u015bci\u0105 i du\u017cym zakresem dzia\u0142a\u0144 wykonywanych w kilku linijkach. W pierwszym parametrze metody builder.stream mo\u017cecie zauwa\u017cy\u0107 co\u015b nowego: warto\u015b\u0107 enumeracji AutoOffsetReset.EARLIEST (istnieje r\u00f3wnie\u017c LATEST), okre\u015blona za pomoc\u0105 metody Consumed.withOffsetResetPolicy. Dzi\u0119ki tej enumeracji mo\u017cna wskaza\u0107 strategi\u0119 resetowania offset\u00f3w dla ka\u017cdego KStream lub KTable, ma ona pierwsze\u0144stwo przed parametrem resetu offset\u00f3w z konfiguracji.<\/p>\n<blockquote><p><b>GroupByKey i GroupBy<\/b><\/p>\n<p>W interfejsie KStream znajduj\u0105 si\u0119 dwie metody do grupowania rekord\u00f3w: GroupByKey i GroupBy. Obie zwracaj\u0105 KGroupedTable, wi\u0119c mo\u017ce pojawi\u0107 si\u0119 naturalne pytanie: jaka jest r\u00f3\u017cnica mi\u0119dzy nimi i kiedy u\u017cywa\u0107 kt\u00f3rej z nich?<\/p>\n<p>Metoda GroupByKey jest stosowana, gdy klucze w KStream s\u0105 ju\u017c niepuste. Co wi\u0119cej, flaga \u201ewymaga ponownego sekcjonowania\u201d nigdy nie by\u0142a ustawiana.<\/p>\n<p>Metoda GroupBy zak\u0142ada, \u017ce zmieni\u0142e\u015b klucze do grupowania, wi\u0119c flaga ponownego sekcjonowania jest ustawiona na true. Wykonanie po metodzie GroupBy operacji \u0142\u0105czenia, agregacji itp. spowoduje automatyczne ponowne sekcjonowanie.<br \/>\nPodsumowanie: zawsze nale\u017cy, gdy to mo\u017cliwe, u\u017cywa\u0107 GroupByKey, a nie GroupBy.<\/p><\/blockquote>\n<p>\nTo, co robi\u0105 metody mapValues i groupBy, jest jasne, wi\u0119c przyjrzyjmy si\u0119 metodzie sum() (mo\u017cna j\u0105 znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/model\/ShareVolume.java) (listing 5.3).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/8e9a6f873594f9b5fef9a96c42353a61.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nMetoda ShareVolume.sum zwraca po\u015bredni\u0105 sum\u0119 wolumenu sprzeda\u017cy akcji, a wynik ca\u0142ego ci\u0105gu oblicze\u0144 stanowi obiekt KTable. Teraz rozumiesz, jak\u0105 rol\u0119 pe\u0142ni KTable. Po przyj\u0119ciu obiekt\u00f3w ShareVolume w odpowiednim obiekcie KTable przechowywana jest ostatnia aktualizacja. Wa\u017cne jest, aby pami\u0119ta\u0107, \u017ce wszystkie aktualizacje s\u0105 odzwierciedlane w poprzednim shareVolumeKTable, ale nie wszystkie s\u0105 wysy\u0142ane dalej.<\/p>\n<p>Nast\u0119pnie za pomoc\u0105 tego KTable wykonujemy agregacj\u0119 (wed\u0142ug liczby sprzedanych akcji), aby uzyska\u0107 pi\u0119\u0107 firm z najwy\u017cszymi wolumenami sprzeda\u017cy akcji w ka\u017cdej bran\u017cy. Nasze dzia\u0142ania b\u0119d\u0105 podobne do dzia\u0142a\u0144 przy pierwszej agregacji.<\/p>\n<ol>\n<li>Wykonaj kolejn\u0105 operacj\u0119 groupBy w celu grupowania poszczeg\u00f3lnych obiekt\u00f3w ShareVolume wed\u0142ug bran\u017c.<\/li>\n<li>Przyst\u0105p do sumowania obiekt\u00f3w ShareVolume. Tym razem obiekt agreguj\u0105cy stanowi kolejk\u0119 o ustalonej wielko\u015bci. W takiej kolejce o ustalonej wielko\u015bci przechowywane s\u0105 tylko pi\u0119\u0107 firm z najwy\u017cszymi liczbami sprzedanych akcji.<\/li>\n<li>Przekszta\u0142\u0107 kolejki z poprzedniego punktu na warto\u015b\u0107 \u0142a\u0144cuchow\u0105 i zwr\u00f3\u0107 pi\u0119\u0107 najlepiej sprzedaj\u0105cych si\u0119 akcji wed\u0142ug bran\u017c.<\/li>\n<li>Zapisz wyniki w formie tekstowej w temacie.<\/li>\n<\/ol>\n<p>\nNa rys. 5.10 przedstawiona jest mapa topologii ruchu danych. Jak wida\u0107, drugi cykl przetwarzania jest do\u015b\u0107 prosty.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/dabd1507eee267038edb7f8d76d8d8ae.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nTeraz, maj\u0105c jasno okre\u015blon\u0105 struktur\u0119 tego drugiego cyklu przetwarzania, mo\u017cna przej\u015b\u0107 do jego kodu \u017ar\u00f3d\u0142owego (znajdziesz go w pliku src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (listing 5.4).<\/p>\n<p>W tym inicjatorze znajduje si\u0119 zmienna fixedQueue. To obiekt u\u017cytkownika \u2014 adapter dla java.util.TreeSet, kt\u00f3ry jest u\u017cywany do \u015bledzenia N najwi\u0119kszych wynik\u00f3w w kolejno\u015bci malej\u0105cej liczby sprzedanych akcji.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/488b072b0d91b81c925ca72291e69e48.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nJu\u017c spotka\u0142e\u015b si\u0119 z wywo\u0142aniami groupBy i mapValues, wi\u0119c nie b\u0119dziemy si\u0119 na nich zatrzymywa\u0107 (wywo\u0142ujemy metod\u0119 KTable.toStream, poniewa\u017c metoda KTable.print jest uznawana za przestarza\u0142\u0105). Ale jeszcze nie widzia\u0142e\u015b KTable-wersji metody aggregate(), wi\u0119c po\u015bwi\u0119cimy chwil\u0119 na jej om\u00f3wienie.<\/p>\n<p>Jak pami\u0119tasz, KTable r\u00f3\u017cni si\u0119 tym, \u017ce rekordy z tymi samymi kluczami s\u0105 traktowane jako aktualizacje. KTable zast\u0119puje stary rekord nowym. Agregacja odbywa si\u0119 w podobny spos\u00f3b: agregowane s\u0105 najnowsze rekordy z tym samym kluczem. Po otrzymaniu rekordu jest on dodawany do instancji klasy FixedSizePriorityQueue za pomoc\u0105 sumatora (drugi parametr w wywo\u0142aniu metody aggregate), ale je\u015bli ju\u017c istnieje inny rekord z tym samym kluczem, stary rekord jest usuwany za pomoc\u0105 reduktora (trzeci parametr w wywo\u0142aniu metody aggregate).<\/p>\n<p>To wszystko oznacza, \u017ce nasz agregator, FixedSizePriorityQueue, wcale nie agreguje wszystkich warto\u015bci z tym samym kluczem, tylko przechowuje przesuwaj\u0105c\u0105 si\u0119 sum\u0119 ilo\u015bci N najlepiej sprzedaj\u0105cych si\u0119 rodzaj\u00f3w akcji. Ka\u017cdy przychodz\u0105cy rekord zawiera \u0142\u0105czn\u0105 liczb\u0119 sprzedanych dotychczas akcji. KTable dostarczy ci informacji o tym, jakie akcje firm sprzedaj\u0105 si\u0119 obecnie najlepiej, a przesuwaj\u0105ca agregacja ka\u017cdego z aktualizacji nie jest konieczna.<\/p>\n<p>Nauczyli\u015bmy si\u0119 robi\u0107 dwie wa\u017cne rzeczy:<\/p>\n<ul>\n<li>grupowa\u0107 warto\u015bci w KTable wed\u0142ug wsp\u00f3lnego klucza;<\/li>\n<li>wykonywa\u0107 na tych zgromadzonych warto\u015bciach przydatne operacje, takie jak redukcja i agregacja.<\/li>\n<\/ul>\n<p>\nUmiej\u0119tno\u015b\u0107 wykonywania tych operacji jest wa\u017cna dla zrozumienia sensu danych przep\u0142ywaj\u0105cych przez aplikacj\u0119 Kafka Streams i ustalenia, jakie informacje nios\u0105 ze sob\u0105.<\/p>\n<p>Po\u0142\u0105czyli\u015bmy w jedn\u0105 ca\u0142o\u015b\u0107 niekt\u00f3re z kluczowych poj\u0119\u0107, omawianych wcze\u015bniej w tej ksi\u0105\u017cce. W rozdziale 4 opowiedzieli\u015bmy, jak wa\u017cne jest lokalne, odporne na awarie stanie dla aplikacji strumieniowych. Pierwszy przyk\u0142ad z tego rozdzia\u0142u pokaza\u0142, dlaczego lokalne stanie jest tak istotne \u2014 umo\u017cliwia \u015bledzenie, jakie informacje ju\u017c zosta\u0142y zobaczone. Lokalne podej\u015bcie pozwala unikn\u0105\u0107 op\u00f3\u017anie\u0144 sieciowych, co czyni aplikacj\u0119 bardziej wydajn\u0105 i odporn\u0105 na b\u0142\u0119dy.<\/p>\n<p>Podczas wykonywania jakiejkolwiek operacji agregacji lub redukcji nale\u017cy poda\u0107 nazw\u0119 magazynu stanu. Operacje redukcji i agregacji zwracaj\u0105 instancj\u0119 KTable, a KTable korzysta z magazynu stanu do zast\u0119powania starych wynik\u00f3w nowymi. Jak zauwa\u017cy\u0142e\u015b, nie wszystkie aktualizacje s\u0105 wysy\u0142ane dalej w \u0142a\u0144cuchu, co jest wa\u017cne, poniewa\u017c operacje agregacji s\u0105 przeznaczone do uzyskania informacji ko\u0144cowych. Je\u015bli nie b\u0119dziemy stosowa\u0107 lokalnego stanu, KTable b\u0119dzie przesy\u0142a\u0107 wszystkie wyniki agregacji i redukcji.<\/p>\n<p>Teraz przyjrzymy si\u0119 wykonaniu operacji, takich jak agregacja, w okre\u015blonym przedziale czasowym \u2014 tak zwanych operacjach okna (windowing operations).<\/p>\n<h3>5.3.2. Operacje okna<\/h3>\n<p>\nW poprzedniej sekcji zapoznali\u015bmy si\u0119 z \u201eprzesuwaj\u0105c\u0105 si\u0119\u201d redukcj\u0105 i agregacj\u0105. Aplikacja przeprowadza\u0142a ci\u0105g\u0142\u0105 redukcj\u0119 wolumenu sprzeda\u017cy akcji, a nast\u0119pnie agregacj\u0119 pi\u0119ciu najbardziej sprzedawanych akcji na gie\u0142dzie.<\/p>\n<p>Czasami takie ci\u0105g\u0142e agregacje i redukcje wynik\u00f3w s\u0105 niezb\u0119dne. A czasami trzeba wykona\u0107 operacje tylko w obr\u0119bie okre\u015blonego przedzia\u0142u czasowego. Na przyk\u0142ad, obliczy\u0107, ile transakcji gie\u0142dowych z akcjami konkretnej firmy mia\u0142o miejsce w ci\u0105gu ostatnich 10 minut. Lub ile u\u017cytkownik\u00f3w klikn\u0119\u0142o w nowy baner reklamowy w ci\u0105gu ostatnich 15 minut. Aplikacja mo\u017ce wykonywa\u0107 takie operacje wielokrotnie, ale z wynikami odnosz\u0105cymi si\u0119 tylko do okre\u015blonych przedzia\u0142\u00f3w czasowych (okien czasowych).<\/p>\n<h3>Liczenie transakcji gie\u0142dowych wed\u0142ug klienta<\/h3>\n<p>\nW nast\u0119pnym przyk\u0142adzie zajmiemy si\u0119 \u015bledzeniem transakcji gie\u0142dowych dla kilku trader\u00f3w \u2014 zar\u00f3wno du\u017cych organizacji, jak i bystrych inwestor\u00f3w indywidualnych.<\/p>\n<p>Istniej\u0105 dwie mo\u017cliwe przyczyny takiego \u015bledzenia. Jedn\u0105 z nich jest potrzeba wiedzy o tym, co kupuj\u0105\/sprzedaj\u0105 liderzy rynku. Je\u015bli ci duzi gracze i do\u015bwiadczeni inwestorzy dostrzegaj\u0105 otwieraj\u0105ce si\u0119 przed nimi mo\u017cliwo\u015bci, sensowne jest pod\u0105\u017canie za ich strategi\u0105. Drug\u0105 przyczyn\u0105 jest ch\u0119\u0107 dostrzegania jakichkolwiek mo\u017cliwych oznak nielegalnych transakcji z wykorzystaniem informacji poufnych. W tym celu b\u0119dziesz musia\u0142 przeanalizowa\u0107 korelacj\u0119 du\u017cych wzrost\u00f3w sprzeda\u017cy z wa\u017cnymi komunikatami prasowymi.<\/p>\n<p>Takie \u015bledzenie sk\u0142ada si\u0119 z takich etap\u00f3w jak:<\/p>\n<ul>\n<li>tworzenie strumienia do odczytu z tematu stock-transactions;<\/li>\n<li>grupowanie przychodz\u0105cych wpis\u00f3w wed\u0142ug identyfikatora nabywcy i symbolu gie\u0142dowego akcji. Wywo\u0142anie metody groupBy zwraca instancj\u0119 klasy KGroupedStream;<\/li>\n<li>zwracanie metod\u0105 KGroupedStream.windowedBy strumienia danych ograniczonego oknem czasowym, co pozwala na wykonywanie agregacji okien. W zale\u017cno\u015bci od typu okna zwracane jest albo TimeWindowedKStream, albo SessionWindowedKStream;<\/li>\n<li>zliczanie transakcji dla operacji agregacji. Okienny strumie\u0144 danych okre\u015bla, czy konkretny wpis jest brany pod uwag\u0119 w tym zliczeniu;<\/li>\n<li>zapisywanie wynik\u00f3w w temacie lub wyprowadzanie ich na konsol\u0119 podczas rozwoju.<\/li>\n<\/ul>\n<p>\nTopologia tej aplikacji jest prosta, ale wizualizacja jej obrazu nie zaszkodzi. Przyjrzyjmy si\u0119 rys. 5.11.<\/p>\n<p>Nast\u0119pnie om\u00f3wimy funkcjonalno\u015b\u0107 operacji okiennych oraz odpowiedni kod.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/67d9d8d8acb1370a7c7fe5cd9a8b6aa7.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Typy okien<\/h3>\n<p>\nW Kafka Streams istniej\u0105 trzy typy okien:<\/p>\n<ul>\n<li>sesyjne;<\/li>\n<li>w \u201eturlaj\u0105cych si\u0119\u201d (tumbling);<\/li>\n<li>w \u201eprzesuwaj\u0105cych si\u0119\/\u201eskacz\u0105cych\u201d (sliding\/hopping).<\/li>\n<\/ul>\n<p>\nWyb\u00f3r zale\u017cy od wymaga\u0144 biznesowych. \u201eTurlaj\u0105ce si\u0119\u201d i \u201eskacz\u0105ce\u201d okna s\u0105 ograniczone czasowo, podczas gdy ograniczenia okien sesyjnych zwi\u0105zane s\u0105 z dzia\u0142aniami u\u017cytkownik\u00f3w \u2014 d\u0142ugo\u015b\u0107 sesji(-y) okre\u015bla si\u0119 wy\u0142\u0105cznie na podstawie aktywno\u015bci u\u017cytkownika. Najwa\u017cniejsze jest, aby pami\u0119ta\u0107, \u017ce wszystkie typy okien opieraj\u0105 si\u0119 na znacznikach daty\/czasu wpis\u00f3w, a nie na czasie systemowym.<\/p>\n<p>Nast\u0119pnie zaimplementujemy nasz\u0105 topologi\u0119 z ka\u017cdym z typ\u00f3w okien. Pe\u0142ny kod zostanie podany tylko w pierwszym przyk\u0142adzie, dla innych typ\u00f3w okien nic si\u0119 nie zmieni, poza typem operacji okiennej.<\/p>\n<h3>Okna sesyjne<\/h3>\n<p>\nOkna sesyjne znacz\u0105co r\u00f3\u017cni\u0105 si\u0119 od wszystkich innych typ\u00f3w okien. Ograniczaj\u0105 si\u0119 nie tyle czasowo, co aktywno\u015bci\u0105 u\u017cytkownika (lub aktywno\u015bci\u0105 jednostki, kt\u00f3r\u0105 chcia\u0142by\u015b \u015bledzi\u0107). Okna sesyjne s\u0105 wyznaczane okresem bezczynno\u015bci.<\/p>\n<p>Rysunek 5.12 ilustruje poj\u0119cie okien sesyjnych. Mniejsza sesja \u0142\u0105czy si\u0119 z sesj\u0105 po lewej stronie. Sesja po prawej stronie b\u0119dzie osobna, poniewa\u017c nast\u0119puje po d\u0142ugim okresie bezczynno\u015bci. Okna sesyjne opieraj\u0105 si\u0119 na dzia\u0142aniach u\u017cytkownik\u00f3w, ale stosuj\u0105 znaczniki daty\/czasu z zapis\u00f3w, aby okre\u015bli\u0107, do kt\u00f3rej sesji odnosi si\u0119 dany zapis.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/d39a5db7af7d6aa2194802622b2b47fd.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Wykorzystanie okien sesyjnych do \u015bledzenia transakcji gie\u0142dowych.<\/h3>\n<p>\nWykorzystamy okna sesyjne do uchwycenia informacji o transakcjach gie\u0142dowych. Implementacja okien sesyjnych jest przedstawiona w listing 5.5 (kt\u00f3ry mo\u017cna znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKTableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/2dcbd9a36baec0e746aad165121451b3.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nWi\u0119kszo\u015b\u0107 operacji w tej topologii ju\u017c znasz, wi\u0119c nie ma potrzeby ponownego ich omawiania. Jednak s\u0105 te\u017c tutaj pewne nowe elementy, kt\u00f3re teraz om\u00f3wimy.<\/p>\n<p>Przy ka\u017cdej operacji groupBy zazwyczaj wykonywana jest jaka\u015b operacja agregacji (agregacja, redukcja lub zliczanie). Mo\u017cna wykona\u0107 zar\u00f3wno agregacj\u0119 skumulowan\u0105 z narastaj\u0105cym wynikiem, jak i agregacj\u0119 okienn\u0105, w kt\u00f3rej uwzgl\u0119dnia si\u0119 zapisy w okre\u015blonym czasie.<\/p>\n<p>Kod z listingu 5.5 wykonuje zliczanie transakcji w obr\u0119bie okien sesyjnych. Na rys. 5.13 te dzia\u0142ania s\u0105 analizowane krok po kroku.<\/p>\n<p>Za pomoc\u0105 wywo\u0142ania windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)) tworzymy okno sesyjne z okresem bezczynno\u015bci 20 sekund oraz okresem przechowywania 15 minut. Okres bezczynno\u015bci 20 sekund oznacza, \u017ce aplikacja uwzgl\u0119dni ka\u017cdy zapis, kt\u00f3ry wp\u0142ynie w ci\u0105gu 20 sekund od zako\u0144czenia lub rozpocz\u0119cia bie\u017c\u0105cej sesji w obecnej (aktywnej) sesji.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/9bd47b04698086872fd135b4c67eb938.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nNast\u0119pnie okre\u015blamy, jak\u0105 operacj\u0119 agregacji nale\u017cy wykona\u0107 w oknie sesji \u2014 w tym przypadku count. Je\u015bli przychodz\u0105cy zapis wykracza poza granice interwa\u0142u bezczynno\u015bci (z dowolnej strony od znacznika daty\/czasu), aplikacja tworzy now\u0105 sesj\u0119. Interwa\u0142 przechowywania oznacza utrzymywanie sesji przez okre\u015blony czas i dopuszcza sp\u00f3\u017anione dane, kt\u00f3re wychodz\u0105 poza okres bezczynno\u015bci sesji, ale wci\u0105\u017c mog\u0105 by\u0107 do\u0142\u0105czane. Ponadto pocz\u0105tek i koniec nowej sesji, kt\u00f3ra powsta\u0142a w wyniku po\u0142\u0105czenia, odpowiadaj\u0105 najwcze\u015bniejszemu i najp\u00f3\u017aniejszemu znacznikowi daty\/czasu.<\/p>\n<p>Przyjrzyjmy si\u0119 kilku zapisom z metody count, aby zobaczy\u0107, jak dzia\u0142aj\u0105 sesje (tab. 5.1).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/ec04aae466d88c2d2349474069c8d541.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nGdy przychodz\u0105 zapisy, szukamy ju\u017c istniej\u0105cych sesji z tym samym kluczem, kt\u00f3rych czas zako\u0144czenia jest wcze\u015bniejszy ni\u017c aktualny znacznik daty\/czasu \u2014 interwa\u0142 bezczynno\u015bci oraz czas rozpocz\u0119cia wi\u0119kszy ni\u017c aktualny znacznik daty\/czasu + interwa\u0142 bezczynno\u015bci. Uwzgl\u0119dniaj\u0105c to, cztery zapisy z tab. 5.1 \u0142\u0105cz\u0105 si\u0119 w jedn\u0105 sesj\u0119 w nast\u0119puj\u0105cy spos\u00f3b.<\/p>\n<p>1. Pierwszy przychodzi zapis 1, wi\u0119c czas rozpocz\u0119cia jest r\u00f3wny czasowi zako\u0144czenia i wynosi 00:00:00.<\/p>\n<p>2. Nast\u0119pnie przychodzi zapis 2, i szukamy sesji ko\u0144cz\u0105cych si\u0119 nie wcze\u015bniej ni\u017c 23:59:55 i rozpoczynaj\u0105cych si\u0119 nie p\u00f3\u017aniej ni\u017c 00:00:35. Znajdujemy zapis 1 i \u0142\u0105czymy sesje 1 i 2. Bierzemy czas rozpocz\u0119cia sesji 1 (wcze\u015bniejszy) oraz czas zako\u0144czenia sesji 2 (p\u00f3\u017aniejszy), wi\u0119c nasza nowa sesja zaczyna si\u0119 o 00:00:00 i ko\u0144czy o 00:00:15.<\/p>\n<p>3. Przychodzi zapis 3, szukamy sesji mi\u0119dzy 00:00:30 a 00:01:10 i nie znajdujemy \u017cadnej. Dodajemy drug\u0105 sesj\u0119 dla klucza 123-345-654,FFBE, rozpoczynaj\u0105c\u0105 si\u0119 i ko\u0144cz\u0105c\u0105 o 00:00:50.<\/p>\n<p>4. Przychodzi zapis 4, i szukamy sesji mi\u0119dzy 23:59:45 a 00:00:25. Tym razem znajdujemy obie sesje \u2014 1 i 2. Wszystkie trzy sesje \u0142\u0105cz\u0105 si\u0119 w jedn\u0105, z czasem rozpocz\u0119cia 00:00:00 i czasem zako\u0144czenia 00:00:15.<\/p>\n<p>Z tego, co zosta\u0142o om\u00f3wione w tej sekcji, warto zapami\u0119ta\u0107 nast\u0119puj\u0105ce istotne szczeg\u00f3\u0142y:<\/p>\n<ul>\n<li>sesje to nie okna o sta\u0142ej wielko\u015bci. Czas trwania sesji jest okre\u015blany przez aktywno\u015b\u0107 w ramach danego interwa\u0142u czasowego;<\/li>\n<li>znaczniki daty\/czasu w danych okre\u015blaj\u0105, czy zdarzenie mie\u015bci si\u0119 w istniej\u0105cej sesji, czy w interwale bezczynno\u015bci.<\/li>\n<\/ul>\n<p>\nNast\u0119pnie om\u00f3wimy nast\u0119pn\u0105 odmian\u0119 okien \u2014 \u00abprzewracaj\u0105ce\u00bb okna.<\/p>\n<h3>\u00abPrzewracaj\u0105ce\u00bb okna<\/h3>\n<p>\n\u201ePrzewracaj\u0105ce si\u0119\u201d okna (tumbling) rejestruj\u0105 zdarzenia, kt\u00f3re wpadaj\u0105 w okre\u015blony interwa\u0142 czasowy. Wyobra\u017a sobie, \u017ce musisz uchwyci\u0107 wszystkie transakcje gie\u0142dowe jakiej\u015b firmy co 20 sekund, wi\u0119c zbierasz wszystkie zdarzenia w tym czasie. Po zako\u0144czeniu 20-sekundowego interwa\u0142u okno \u201eprzewraca si\u0119\u201d i przechodzi do nowego 20-sekundowego okresu obserwacji. Rysunek 5.14 ilustruje t\u0119 sytuacj\u0119.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/e90e560d9ddda5e2e524b7c387ad9874.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nJak wida\u0107, wszystkie zdarzenia, kt\u00f3re mia\u0142y miejsce w ostatnich 20 sekundach, s\u0105 zawarte w oknie. Po zako\u0144czeniu tego okresu tworzone jest nowe okno.<\/p>\n<p>W listing 5.6 znajduje si\u0119 kod demonstruj\u0105cy u\u017cycie \u201eprzewracaj\u0105cych si\u0119\u201d okien do uchwycenia transakcji gie\u0142dowych co 20 sekund (mo\u017cna go znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/e055bb1b288c7d500b64372fe3fbf064.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nDzi\u0119ki tej niewielkiej zmianie wywo\u0142ania metody TimeWindows.of mo\u017cna u\u017cywa\u0107 \u201eprzewracaj\u0105cego si\u0119\u201d okna. W tym przyk\u0142adzie nie ma wywo\u0142ania metody until(), w zwi\u0105zku z czym zostanie u\u017cyty domy\u015blny interwa\u0142 przechowywania, wynosz\u0105cy 24 godziny.<\/p>\n<p>Wreszcie nadszed\u0142 czas, aby przej\u015b\u0107 do ostatniego z typ\u00f3w okien \u2014 \u201eskacz\u0105cych\u201d (hopping) okien.<\/p>\n<h3>Okna skacz\u0105ce (\u201ehopping\u201d)<\/h3>\n<p>\nOkna skacz\u0105ce (\u201esliding\/hopping\u201d) s\u0105 podobne do \u201eprzewracaj\u0105cych si\u0119\u201d, ale istnieje pewna r\u00f3\u017cnica. Okna skacz\u0105ce nie czekaj\u0105 na zako\u0144czenie interwa\u0142u czasu przed utworzeniem nowego okna do przetwarzania niedawnych zdarze\u0144. Rozpoczynaj\u0105 nowe obliczenia po interwale oczekiwania, kt\u00f3ry jest kr\u00f3tszy ni\u017c czas trwania okna.<\/p>\n<p>Aby zobrazowa\u0107 r\u00f3\u017cnice mi\u0119dzy \u201eprzewracaj\u0105cymi si\u0119\u201d a \u201eskacz\u0105cymi\u201d oknami, wr\u00f3\u0107my do przyk\u0142adu z liczeniem transakcji gie\u0142dowych. Naszym celem nadal jest liczenie liczby transakcji, ale nie chcemy czeka\u0107 ca\u0142ego okresu czasu przed zaktualizowaniem licznika. Zamiast tego b\u0119dziemy aktualizowa\u0107 licznik w kr\u00f3tszych odst\u0119pach czasu. Na przyk\u0142ad, b\u0119dziemy nadal liczy\u0107 liczb\u0119 transakcji co 20 sekund, ale b\u0119dziemy aktualizowa\u0107 licznik co 5 sekund, jak pokazano na rys. 5.15. W rezultacie b\u0119dziemy mie\u0107 trzy okna wynikowe z nak\u0142adaj\u0105cymi si\u0119 danymi.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/415c8cd9f2b60d453a1a01c3bc99331f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nW listing 5.7 znajduje si\u0119 kod do zadania skacz\u0105cych okien (mo\u017cna go znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/ab2d1a64380d256d3fb084e16597417c.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nOkno \u201eprzewracaj\u0105ce si\u0119\u201d mo\u017cna przekszta\u0142ci\u0107 w \u201eskacz\u0105ce\u201d przez dodanie wywo\u0142ania metody advanceBy(). W przedstawionym przyk\u0142adzie interwa\u0142 zapisywania wynosi 15 minut.<\/p>\n<p>W tej sekcji zobaczy\u0142e\u015b, jak ogranicza\u0107 wyniki agregacji przez czasowe okna. Chcia\u0142bym, aby\u015b zapami\u0119ta\u0142 z tej sekcji trzy rzeczy:<\/p>\n<ul>\n<li>rozmiar okien sesyjnych jest ograniczony nie przez przedzia\u0142 czasu, ale przez aktywno\u015b\u0107 u\u017cytkownik\u00f3w;<\/li>\n<li>\u201eprzewracaj\u0105ce si\u0119\u201d okna daj\u0105 wgl\u0105d w wydarzenia w danym okresie czasu;<\/li>\n<li>czas trwania \u201eskacz\u0105cych\u201d okien jest sta\u0142y, ale s\u0105 one cz\u0119sto aktualizowane i mog\u0105 zawiera\u0107 w ka\u017cdym oknie nak\u0142adaj\u0105ce si\u0119 rekordy.<\/li>\n<\/ul>\n<p>\nDalej dowiemy si\u0119, jak przekszta\u0142ci\u0107 KTable z powrotem w KStream do \u0142\u0105czenia.<\/p>\n<h3>5.3.3. \u0141\u0105czenie obiekt\u00f3w KStream i KTable<\/h3>\n<p>\nW rozdziale 4 omawiali\u015bmy \u0142\u0105czenie dw\u00f3ch obiekt\u00f3w KStream. Teraz nauczymy si\u0119 \u0142\u0105czy\u0107 KTable i KStream. Mo\u017ce by\u0107 to potrzebne z nast\u0119puj\u0105cego prostego powodu. KStream to strumie\u0144 rekord\u00f3w, a KTable to strumie\u0144 aktualizacji rekord\u00f3w, ale czasami mo\u017ce by\u0107 potrzebne dodanie dodatkowego kontekstu do strumienia rekord\u00f3w za pomoc\u0105 aktualizacji z KTable.<\/p>\n<p>We\u017amy dane o liczbie transakcji gie\u0142dowych i po\u0142\u0105czmy je z informacjami o gie\u0142dowych wiadomo\u015bciach dotycz\u0105cych odpowiednich bran\u017c przemys\u0142owych. Oto, co trzeba zrobi\u0107, aby to osi\u0105gn\u0105\u0107, bior\u0105c pod uwag\u0119 ju\u017c istniej\u0105cy kod.<\/p>\n<ol>\n<li>Przekszta\u0142\u0107 obiekt KTable z danymi o liczbie transakcji gie\u0142dowych w KStream, nast\u0119pnie zamie\u0144 klucz na klucz oznaczaj\u0105cy bran\u017c\u0119 przemys\u0142ow\u0105 odpowiadaj\u0105c\u0105 temu symbolowi akcji.<\/li>\n<li>Utw\u00f3rz obiekt KTable, kt\u00f3ry b\u0119dzie odczytywa\u0142 dane z tematu gie\u0142dowych wiadomo\u015bci. Ten nowy KTable b\u0119dzie kategoryzowany wed\u0142ug bran\u017c przemys\u0142owych.<\/li>\n<li>Po\u0142\u0105cz aktualizacje wiadomo\u015bci z informacjami o liczbie transakcji gie\u0142dowych wed\u0142ug bran\u017c przemys\u0142owych.<\/li>\n<\/ol>\n<p>\nTeraz przyjrzyjmy si\u0119, jak zrealizowa\u0107 ten plan dzia\u0142ania.<\/p>\n<h3>Przekszta\u0142canie KTable w KStream<\/h3>\n<p>\nAby przekszta\u0142ci\u0107 KTable w KStream, nale\u017cy zrobi\u0107 nast\u0119puj\u0105ce kroki.<\/p>\n<ol>\n<li>Wywo\u0142a\u0107 metod\u0119 KTable.toStream().<\/li>\n<li>Za pomoc\u0105 wywo\u0142ania metody KStream.map zamie\u0144 klucz na nazw\u0119 bran\u017cy, a nast\u0119pnie wyodr\u0119bnij z instancji Windowed obiekt TransactionSummary.<\/li>\n<\/ol>\n<p>\nPo\u0142\u0105czymy te operacje w \u0142a\u0144cuch w nast\u0119puj\u0105cy spos\u00f3b (kod mo\u017cna znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listing 5.8).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/0d43c2650f6e66e2816ed383da3a29c2.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nPoniewa\u017c wykonujemy operacj\u0119 KStream.map, ponowne sekcjonowanie dla zwr\u00f3conej instancji KStream odbywa si\u0119 automatycznie podczas jej u\u017cycia w po\u0142\u0105czeniu.<\/p>\n<p>Zako\u0144czyli\u015bmy proces przekszta\u0142cania, teraz musimy stworzy\u0107 obiekt KTable do odczytu wiadomo\u015bci gie\u0142dowych.<\/p>\n<h3>Tworzenie KTable dla wiadomo\u015bci gie\u0142dowych<\/h3>\n<p>\nNa szcz\u0119\u015bcie, aby stworzy\u0107 obiekt KTable, wystarczy jedna linia kodu (ten kod mo\u017cna znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listing 5.9).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/6e83a393fdc9ab74fda4cbdddddb5213.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nNale\u017cy zauwa\u017cy\u0107, \u017ce nie trzeba podawa\u0107 \u017cadnych obiekt\u00f3w Serde, poniewa\u017c w ustawieniach u\u017cywane s\u0105 stringowe Serde. Dzi\u0119ki zastosowaniu enumeracji EARLIEST tabela jest wype\u0142niana zapisami na pocz\u0105tku.<\/p>\n<p>Teraz mo\u017cemy przej\u015b\u0107 do ostatniego kroku \u2014 po\u0142\u0105czenia.<\/p>\n<h3>\u0141\u0105czenie aktualizacji wiadomo\u015bci z danymi o liczbie transakcji<\/h3>\n<p>\nTworzenie po\u0142\u0105czenia nie sprawia trudno\u015bci. Skorzystamy z lewego po\u0142\u0105czenia na wypadek, gdyby nie by\u0142o wiadomo\u015bci gie\u0142dowych dla danej bran\u017cy (potrzebny kod mo\u017cna znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listing 5.10).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/92ed70f98927d2f778ad14dbc2a5aa26.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nOperator leftJoin jest do\u015b\u0107 prosty. W przeciwie\u0144stwie do po\u0142\u0105cze\u0144 z rozdzia\u0142u 4, metoda JoinWindow nie jest u\u017cywana, poniewa\u017c podczas wykonywania po\u0142\u0105czenia KStream-KTable dla ka\u017cdego klucza w KTable znajduje si\u0119 tylko jeden zapis. To po\u0142\u0105czenie nie jest ograniczone czasowo: zapis jest albo w KTable, albo go nie ma. G\u0142\u00f3wna konkluzja: za pomoc\u0105 obiekt\u00f3w KTable mo\u017cna wzbogaca\u0107 KStream danymi referencyjnymi rzadziej aktualizowanymi.<\/p>\n<p>A teraz przyjrzymy si\u0119 bardziej wydajnemu sposobowi wzbogacania zdarze\u0144 z KStream.<\/p>\n<h3>5.3.4. Obiekty GlobalKTable<\/h3>\n<p>\nJak ju\u017c rozumiesz, istnieje potrzeba wzbogacania strumieni zdarze\u0144 lub dodawania do nich kontekstu. W rozdziale 4 widzia\u0142e\u015b po\u0142\u0105czenia dw\u00f3ch obiekt\u00f3w KStream, a w poprzednim rozdziale \u2014 po\u0142\u0105czenie KStream i KTable. We wszystkich tych przypadkach konieczne jest ponowne sekcjonowanie strumienia danych przy mapowaniu kluczy na nowy typ lub warto\u015b\u0107. Czasami ponowne sekcjonowanie jest wykonywane jawnie, a czasami Kafka Streams robi to automatycznie. Ponowne sekcjonowanie jest konieczne, poniewa\u017c klucze si\u0119 zmieni\u0142y i zapisy powinny znale\u017a\u0107 si\u0119 w nowych sekcjach, w przeciwnym razie po\u0142\u0105czenie stanie si\u0119 niemo\u017cliwe (to by\u0142o om\u00f3wione w rozdziale 4, w punkcie 'Ponowne sekcjonowanie danych' sekcji 4.2.4).<\/p>\n<h3>Re-sekcjonowanie wi\u0105\u017ce si\u0119 z kosztami<\/h3>\n<p>\nRe-sekcjonowanie wymaga wydatk\u00f3w \u2014 dodatkowych zasob\u00f3w na tworzenie po\u015brednich temat\u00f3w, przechowywanie powtarzaj\u0105cych si\u0119 danych w jeszcze jednym temacie; oznacza to r\u00f3wnie\u017c zwi\u0119kszenie op\u00f3\u017anienia w wyniku zapisu i odczytu z tego tematu. Dodatkowo, w przypadku konieczno\u015bci wykonania po\u0142\u0105czenia na wi\u0119cej ni\u017c jednym aspekcie lub wymiarze, nale\u017cy zorganizowa\u0107 po\u0142\u0105czenia w \u0142a\u0144cuch, zmapowa\u0107 rekordy z nowymi kluczami i ponownie przeprowadzi\u0107 proces re-sekcjonowania.<\/p>\n<h3>Po\u0142\u0105czenie z zestawami danych mniejszych rozmiar\u00f3w<\/h3>\n<p>\nW niekt\u00f3rych przypadkach obj\u0119to\u015b\u0107 danych referencyjnych, z kt\u00f3rymi planowane jest po\u0142\u0105czenie, jest stosunkowo niewielka, wi\u0119c ich pe\u0142ne kopie mog\u0105 zmie\u015bci\u0107 si\u0119 lokalnie na ka\u017cdym z w\u0119z\u0142\u00f3w. Do takich sytuacji w Kafka Streams przewidziany jest klasa GlobalKTable.<\/p>\n<p>Instancje GlobalKTable s\u0105 unikalne, poniewa\u017c aplikacja replikuj\u0105ca wszystkie dane na ka\u017cdym z w\u0119z\u0142\u00f3w. A poniewa\u017c na ka\u017cdym z w\u0119z\u0142\u00f3w znajduj\u0105 si\u0119 wszystkie dane, nie ma potrzeby sekcjonowania strumienia zdarze\u0144 wed\u0142ug klucza danych referencyjnych, aby by\u0142 dost\u0119pny we wszystkich sekcjach. Dzi\u0119ki obiektom GlobalKTable mo\u017cna r\u00f3wnie\u017c realizowa\u0107 po\u0142\u0105czenia bezkluczowe. Wr\u00f3\u0107my do jednego z poprzednich przyk\u0142ad\u00f3w, aby zobaczy\u0107 t\u0119 mo\u017cliwo\u015b\u0107.<\/p>\n<h3>Po\u0142\u0105czenie obiekt\u00f3w KStream z obiektami GlobalKTable<\/h3>\n<p>\nW podrozdziale 5.3.2 przeprowadzili\u015bmy agregacj\u0119 okienkow\u0105 transakcji gie\u0142dowych wed\u0142ug klient\u00f3w. Wyniki tej agregacji wygl\u0105da\u0142y mniej wi\u0119cej tak:<\/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>\nChocia\u017c te wyniki spe\u0142ni\u0142y postawiony cel, by\u0142oby wygodniej, gdyby wy\u015bwietlane by\u0142y tak\u017ce imi\u0119 klienta i pe\u0142na nazwa firmy. Aby doda\u0107 imi\u0119 kupuj\u0105cego i nazw\u0119 firmy, mo\u017cna wykona\u0107 zwyk\u0142e po\u0142\u0105czenia, ale wymaga\u0142oby to dw\u00f3ch mapowa\u0144 kluczy i ponownego sekcjonowania. Dzi\u0119ki GlobalKTable mo\u017cna unikn\u0105\u0107 koszt\u00f3w zwi\u0105zanych z takimi operacjami.<\/p>\n<p>W tym celu skorzystamy z obiektu countStream z listing 5.11 (odpowiedni kod mo\u017cna znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java), \u0142\u0105cz\u0105c go z dwoma obiektami GlobalKTable.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/fc4d91bbe062ceb94f5650224840b81e.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nJu\u017c to omawiali\u015bmy wcze\u015bniej, wi\u0119c nie b\u0119d\u0119 si\u0119 powtarza\u0107. Zwr\u00f3c\u0119 jednak uwag\u0119, \u017ce kod w funkcji toStream().map jest zorganizowany w obiekt-funkcj\u0119 dla lepszej czytelno\u015bci, zamiast by\u0107 wbudowan\u0105 wyra\u017ceniem lambda.<\/p>\n<p>Nast\u0119pny krok to deklaracja dw\u00f3ch instancji GlobalKTable (przyk\u0142adowy kod znajduje si\u0119 w pliku src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (fragment 5.12).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/db3918895f174c8cfb5fc927b55f87a1.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p>Zauwa\u017c, \u017ce nazwy temat\u00f3w s\u0105 opisywane za pomoc\u0105 typ\u00f3w wyliczanych.<\/p>\n<p>Teraz, gdy przygotowali\u015bmy wszystkie komponenty, pozostaje napisa\u0107 kod do po\u0142\u0105czenia (kt\u00f3ry mo\u017cna znale\u017a\u0107 w pliku src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (fragment 5.13).<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/799360cc99f1920c190a61fd4685d4ff.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nChocia\u017c w tym kodzie znajduj\u0105 si\u0119 dwa po\u0142\u0105czenia, s\u0105 one zorganizowane w formie \u0142a\u0144cucha, poniewa\u017c \u017caden z ich wynik\u00f3w nie jest u\u017cywany osobno. Wyniki s\u0105 wyprowadzane na ko\u0144cu ca\u0142ej operacji.<\/p>\n<p>Podczas uruchamiania powy\u017cszej operacji po\u0142\u0105czenia otrzymasz wyniki w nast\u0119puj\u0105cym formacie:<\/p>\n<pre><code class=\"plaintext\">{customer='Barney, Smith' company=\"Exxon\", transactions= 17}<\/code><\/pre>\n<p>\nIstota si\u0119 nie zmieni\u0142a, ale te wyniki wygl\u0105daj\u0105 bardziej zrozumiale.<\/p>\n<p>Je\u015bli uwzgl\u0119dnimy rozdzia\u0142 4, ju\u017c widzia\u0142e\u015b kilka typ\u00f3w po\u0142\u0105cze\u0144 w dzia\u0142aniu. S\u0105 one wymienione w tabeli 5.2. Ta tabela odzwierciedla mo\u017cliwo\u015bci po\u0142\u0105cze\u0144 aktualne dla wersji 1.0.0 Kafka Streams; w przysz\u0142ych wydaniach mog\u0105 si\u0119 zmieni\u0107.<\/p>\n<p><img decoding=\"async\" alt=\"Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d\" src=\"\/wp-content\/uploads\/8e4cf35c64a8431bda43a5e748de275f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nNa zako\u0144czenie przypomn\u0119 najwa\u017cniejsze: mo\u017cesz \u0142\u0105czy\u0107 strumienie zdarze\u0144 (KStream) i strumienie aktualizacji (KTable) przy u\u017cyciu lokalnego stanu. Ponadto, je\u015bli rozmiar danych referencyjnych nie jest zbyt du\u017cy, mo\u017cna skorzysta\u0107 z obiektu GlobalKTable. GlobalKTable replikuj\u0105 wszystkie sekcje na ka\u017cdym z w\u0119z\u0142\u00f3w aplikacji Kafka Streams, zapewniaj\u0105c dost\u0119pno\u015b\u0107 wszystkich danych, niezale\u017cnie od tego, do kt\u00f3rej sekcji nale\u017cy klucz.<\/p>\n<p>Nast\u0119pnie zobaczymy mo\u017cliwo\u015b\u0107 Kafka Streams, kt\u00f3ra pozwala obserwowa\u0107 zmiany stanu bez konsumowania danych z tematu Kafka.<\/p>\n<h3>5.3.5. Stan dost\u0119pny dla zapyta\u0144<\/h3>\n<p>\nJu\u017c wykona\u0142e\u015b kilka operacji z wykorzystaniem stanu i zawsze wy\u015bwietlali\u015bmy wyniki w konsoli (do cel\u00f3w rozwojowych) lub zapisywali\u015bmy je do tematu (dla cel\u00f3w eksploatacyjnych). Podczas zapisywania wynik\u00f3w do tematu konieczne jest u\u017cycie konsumenta Kafka do ich przegl\u0105dania.<\/p>\n<p>Czytanie danych z tych temat\u00f3w mo\u017cna uzna\u0107 za rodzaj materializowanych widok\u00f3w (materialized views). Dla naszych potrzeb mo\u017cemy u\u017cy\u0107 definicji materializowanego widoku z \u201eWikipedii\u201d: \u201e\u2026fizyczny obiekt bazy danych, kt\u00f3ry zawiera wyniki wykonania zapytania. Na przyk\u0142ad mo\u017ce to by\u0107 lokalna kopia danych zdalnych, lub podzbi\u00f3r wierszy i\/lub kolumn tabeli lub wynik\u00f3w po\u0142\u0105czenia, lub tabela przestawna uzyskana za pomoc\u0105 agregacji\u201d (https:\/\/en.wikipedia.org\/wiki\/Materialized_view).<\/p>\n<p>Kafka Streams umo\u017cliwia r\u00f3wnie\u017c wykonywanie interaktywnych zapyta\u0144 (interactive queries) do magazyn\u00f3w stanu, co daje mo\u017cliwo\u015b\u0107 bezpo\u015bredniego odczytu tych materializowanych widok\u00f3w. Wa\u017cne jest, aby zauwa\u017cy\u0107, \u017ce zapytanie do magazynu stanu ma charakter operacji \u201etylko do odczytu\u201d. Dzi\u0119ki temu mo\u017cesz nie martwi\u0107 si\u0119 przypadkowym wprowadzeniem niezgodno\u015bci w stanie podczas przetwarzania danych przez aplikacj\u0119.<\/p>\n<p>Mo\u017cliwo\u015b\u0107 bezpo\u015brednich zapyta\u0144 do magazyn\u00f3w stanu ma ogromne znaczenie. Oznacza to, \u017ce mo\u017cna tworzy\u0107 aplikacje \u2014 pulpity nawigacyjne, bez potrzeby najpierw pozyskiwania danych od konsumenta Kafka. Zwi\u0119ksza ona r\u00f3wnie\u017c wydajno\u015b\u0107 aplikacji, poniewa\u017c nie ma potrzeby ponownego zapisywania danych:<\/p>\n<ul>\n<li>dzi\u0119ki lokalno\u015bci danych mo\u017cna szybko uzyska\u0107 do nich dost\u0119p;<\/li>\n<li>eliminuje si\u0119 duplikacj\u0119 danych, poniewa\u017c nie s\u0105 one zapisywane w zewn\u0119trznej pami\u0119ci.<\/li>\n<\/ul>\n<p>\nNajwa\u017cniejsze, co chcia\u0142bym, aby\u015b zapami\u0119ta\u0142: mo\u017cna bezpo\u015brednio wykonywa\u0107 zapytania do stanu z aplikacji. Nie mo\u017cna przeceni\u0107 mo\u017cliwo\u015bci, kt\u00f3re to daje. Zamiast konsumowa\u0107 dane z Kafka i zapisywa\u0107 rekordy w bazie danych dla aplikacji, mo\u017cna wykonywa\u0107 zapytania do magazyn\u00f3w stanu z tym samym wynikiem. Bezpo\u015brednie zapytania do magazyn\u00f3w stanu oznaczaj\u0105 mniej kodu (brak konsumenta) i mniej oprogramowania (brak potrzeby posiadania tabeli bazy danych do przechowywania wynik\u00f3w).<\/p>\n<p>W tej rozdziale poruszyli\u015bmy du\u017c\u0105 ilo\u015b\u0107 informacji, dlatego na chwil\u0119 zaprzestaniemy omawiania interaktywnych zapyta\u0144 do magazyn\u00f3w stanu. Ale nie martwcie si\u0119: w rozdziale 9 stworzymy prost\u0105 aplikacj\u0119 \u2014 pulpit nawigacyjny z interaktywnymi zapytaniami. W celu demonstracji interaktywnych zapyta\u0144 i mo\u017cliwo\u015bci ich dodawania do aplikacji Kafka Streams skorzystamy z niekt\u00f3rych przyk\u0142ad\u00f3w z tej i poprzedniej rozdzia\u0142u.<\/p>\n<h3>Podsumowanie<\/h3>\n<p><\/p>\n<ul>\n<li>Obiekty KStream reprezentuj\u0105 strumienie zdarze\u0144, por\u00f3wnywalne z wstawieniami do bazy danych. Obiekty KTable reprezentuj\u0105 strumienie aktualizacji, bardziej przypominaj\u0105 aktualizacje w bazie danych. Rozmiar obiektu KTable nie ro\u015bnie, stare rekordy s\u0105 zast\u0119powane nowymi.<\/li>\n<li>Obiekty KTable s\u0105 niezb\u0119dne do operacji agregacyjnych.<\/li>\n<li>Dzi\u0119ki operacjom okiennym mo\u017cna podzieli\u0107 agregowane dane na przedzia\u0142y czasowe.<\/li>\n<li>Dzi\u0119ki obiektom GlobalKTable mo\u017cemy uzyska\u0107 dost\u0119p do danych referencyjnych w dowolnym punkcie aplikacji, niezale\u017cnie od podzia\u0142u na sekcje.<\/li>\n<li>Mo\u017cliwe s\u0105 po\u0142\u0105czenia mi\u0119dzy obiektami KStream, KTable i GlobalKTable.<\/li>\n<\/ul>\n<p>\nDotychczas koncentrowali\u015bmy si\u0119 na tworzeniu aplikacji Kafka Streams przy u\u017cyciu wysokopoziomowego DSL KStream. Chocia\u017c podej\u015bcie wysokopoziomowe pozwala na pisanie czystych i zwi\u0119z\u0142ych program\u00f3w, jego wykorzystanie wi\u0105\u017ce si\u0119 z pewnym kompromisem. Praca z DSL KStream oznacza wi\u0119ksz\u0105 zwi\u0119z\u0142o\u015b\u0107 kodu kosztem mniejszej kontroli. W nast\u0119pnym rozdziale przyjrzymy si\u0119 niskopoziomowemu API w\u0119z\u0142\u00f3w przetwarzaj\u0105cych i spr\u00f3bujemy innych kompromis\u00f3w. Programy b\u0119d\u0105 d\u0142u\u017csze ni\u017c dotychczas, ale zyskamy mo\u017cliwo\u015b\u0107 tworzenia praktycznie dowolnego w\u0119z\u0142a przetwarzaj\u0105cego, kt\u00f3ry mo\u017ce by\u0107 nam potrzebny.<\/p>\n<p>\u2192 Wi\u0119cej informacji o ksi\u0105\u017cce mo\u017cna znale\u017a\u0107 na <noindex><a rel=\"nofollow\" href=\"https:\/\/www.piter.com\/collection\/best\/product\/kafka-streams-v-deystvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\">stronie wydawcy<\/a><\/noindex><\/p>\n<p>\u2192 Dla u\u017cytkownik\u00f3w Habr zni\u017cka 25% przy u\u017cyciu kuponu \u2014 <b>Kafka Streams<\/b><\/p>\n<p>\u2192 Po dokonaniu p\u0142atno\u015bci za papierow\u0105 wersj\u0119 ksi\u0105\u017cki na e-mail zostanie wys\u0142ana wersja elektroniczna.<br \/>\n<br \/>\u0179r\u00f3d\u0142o: <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\/pl\/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=\"pl_PL\" \/>\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\/pl\/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\udd47Ksi\u0105\u017cka \u201eKafka Streams w dzia\u0142aniu. Aplikacje i mikroserwisy do pracy w czasie rzeczywistym\u201d | ProHoster","description":"","canonical_url":"https:\/\/prohoster.info\/pl\/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":"pl_PL","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\/pl\/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\/pl\/wp-json\/wp\/v2\/posts\/35786","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/comments?post=35786"}],"version-history":[{"count":0,"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/posts\/35786\/revisions"}],"wp:attachment":[{"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/media?parent=35786"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/categories?post=35786"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/prohoster.info\/pl\/wp-json\/wp\/v2\/tags?post=35786"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}