{"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\/fr\/blog\/administrirovanie\/kniga-kafka-streams-v-dejstvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni","title":{"rendered":"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb","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=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/6f8bd2b31b87b0c760c1148515893c43.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/a><\/noindex> Bonjour, membres de Habr ! Ce livre est destin\u00e9 \u00e0 tout d\u00e9veloppeur souhaitant comprendre le traitement de flux. Ma\u00eetriser la programmation distribu\u00e9e facilitera l'apprentissage de Kafka et de Kafka Streams. Il serait utile de conna\u00eetre le framework Kafka, mais ce n'est pas obligatoire : je vous expliquerai tout ce qu'il faut savoir. Tant les d\u00e9veloppeurs exp\u00e9riment\u00e9s que les d\u00e9butants pourront gr\u00e2ce \u00e0 ce livre ma\u00eetriser la cr\u00e9ation d'applications int\u00e9ressantes pour le traitement de flux en utilisant la biblioth\u00e8que Kafka Streams. Les d\u00e9veloppeurs Java de niveau interm\u00e9diaire et avanc\u00e9, d\u00e9j\u00e0 familiaris\u00e9s avec des concepts tels que la s\u00e9rialisation, apprendront \u00e0 appliquer leurs comp\u00e9tences pour d\u00e9velopper des applications Kafka Streams. Le code source du livre est \u00e9crit en Java 8 et utilise de mani\u00e8re significative la syntaxe des expressions lambda de Java 8, donc \u00eatre capable de travailler avec des fonctions lambda (m\u00eame dans un autre langage de programmation) vous sera utile.<br \/>\n<noindex><a rel=\"nofollow\" name=\"habracut\"><\/a><\/noindex><\/p>\n<h3>Extrait. 5.3. Agr\u00e9gation et op\u00e9rations fen\u00eatres<\/h3>\n<p>\nDans cette section, nous allons explorer les parties les plus prometteuses de Kafka Streams. Jusqu'\u00e0 pr\u00e9sent, nous avons examin\u00e9 les aspects suivants de Kafka Streams :<\/p>\n<ul>\n<li>cr\u00e9ation de topologies de traitement ;<\/li>\n<li>utilisation des \u00e9tats dans les applications de flux ;<\/li>\n<li>ex\u00e9cution des jointures de flux de donn\u00e9es ;<\/li>\n<li>diff\u00e9rences entre les flux d'\u00e9v\u00e9nements (KStream) et les flux de mise \u00e0 jour (KTable).<\/li>\n<\/ul>\n<p>\nDans les exemples suivants, nous allons rassembler tous ces \u00e9l\u00e9ments. De plus, vous d\u00e9couvrirez les op\u00e9rations fen\u00eatres \u2014 une autre excellente fonctionnalit\u00e9 des applications de flux. Notre premier exemple sera une agr\u00e9gation simple.<\/p>\n<h3>5.3.1. Agr\u00e9gation du volume des ventes d'actions par secteur industriel<\/h3>\n<p>\nL'agr\u00e9gation et le regroupement sont des outils essentiels lors du travail avec des donn\u00e9es de flux. L'examen des enregistrements individuels \u00e0 mesure qu'ils arrivent s'av\u00e8re souvent insuffisant. Pour extraire des informations suppl\u00e9mentaires \u00e0 partir des donn\u00e9es, il est n\u00e9cessaire de les regrouper et de les combiner.<\/p>\n<p>Dans cet exemple, vous serez dans la peau d'un trader intraday qui doit suivre le volume des ventes d'actions des entreprises dans plusieurs secteurs industriels. Plus pr\u00e9cis\u00e9ment, vous vous int\u00e9resserez aux cinq entreprises ayant les volumes de vente d'actions les plus \u00e9lev\u00e9s dans chacun des secteurs industriels.<\/p>\n<p>Pour r\u00e9aliser une telle agr\u00e9gation, plusieurs \u00e9tapes seront n\u00e9cessaires pour convertir les donn\u00e9es au format requis (en termes g\u00e9n\u00e9raux).<\/p>\n<ol>\n<li>Cr\u00e9er une source bas\u00e9e sur un sujet publiant des informations brutes sur le commerce des actions. Nous devrons mapper l'objet de type StockTransaction \u00e0 un objet de type ShareVolume. Le fait est que l'objet StockTransaction contient des m\u00e9tadonn\u00e9es sur les ventes, et nous avons uniquement besoin des donn\u00e9es concernant le nombre d'actions vendues.<\/li>\n<li>Regrouper les donn\u00e9es ShareVolume par symboles d'actions. Apr\u00e8s avoir regroup\u00e9 par symboles, il est possible de r\u00e9duire ces donn\u00e9es \u00e0 des totaux interm\u00e9diaires des volumes de ventes d'actions. Il convient de noter que la m\u00e9thode KStream.groupBy renvoie une instance de type KGroupedStream. Et pour obtenir une instance de KTable, il suffit d'appeler ensuite la m\u00e9thode KGroupedStream.reduce.<\/li>\n<\/ol>\n<p><\/p>\n<blockquote><p><b>Qu'est-ce que l'interface KGroupedStream ?<\/b><\/p>\n<p>Les m\u00e9thodes KStream.groupBy et KStream.groupByKey renvoient une instance de KGroupedStream. KGroupedStream est une repr\u00e9sentation interm\u00e9diaire d'un flux d'\u00e9v\u00e9nements apr\u00e8s regroupement par cl\u00e9s. Il n'est pas du tout destin\u00e9 \u00e0 \u00eatre utilis\u00e9 directement. Au lieu de cela, KGroupedStream est utilis\u00e9 pour des op\u00e9rations d'agr\u00e9gation, dont le r\u00e9sultat est toujours un KTable. Et puisque le r\u00e9sultat des op\u00e9rations d'agr\u00e9gation est un KTable et que le stockage d'\u00e9tat est utilis\u00e9, il se peut que toutes les mises \u00e0 jour ne soient pas envoy\u00e9es plus loin dans le pipeline.<\/p>\n<p>La m\u00e9thode KTable.groupBy renvoie un KGroupedTable similaire \u2014 une repr\u00e9sentation interm\u00e9diaire d'un flux de mises \u00e0 jour regroup\u00e9es par cl\u00e9.<\/p><\/blockquote>\n<p>\nFaisons une petite pause et regardons la fig. 5.9, qui montre ce que nous avons atteint. Cette topologie devrait d\u00e9j\u00e0 vous \u00eatre famili\u00e8re.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/9fd61317cde376362adcaeec72908919.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nJetons maintenant un coup d'\u0153il au code pour cette topologie (vous pouvez le trouver dans le fichier src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (listing 5.2).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/f937287e448295fbd467c283ceca316a.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLe code fourni se distingue par sa concision et le grand nombre d'actions effectu\u00e9es en quelques lignes. Dans le premier param\u00e8tre de la m\u00e9thode builder.stream, vous pouvez remarquer quelque chose de nouveau : la valeur de l'\u00e9num\u00e9ration AutoOffsetReset.EARLIEST (il existe \u00e9galement LATEST), d\u00e9finie \u00e0 l'aide de la m\u00e9thode Consumed.withOffsetResetPolicy. Avec ce type d'\u00e9num\u00e9ration, vous pouvez sp\u00e9cifier la strat\u00e9gie de r\u00e9initialisation des offsets pour chaque KStream ou KTable, qui a la priorit\u00e9 sur le param\u00e8tre de r\u00e9initialisation des offsets de la configuration.<\/p>\n<blockquote><p><b>GroupByKey et GroupBy<\/b><\/p>\n<p>Dans l'interface KStream, il existe deux m\u00e9thodes pour regrouper des enregistrements : GroupByKey et GroupBy. Les deux renvoient un KGroupedTable, vous pouvez donc vous poser la question : quelle est la diff\u00e9rence entre elles et quand utiliser l'une ou l'autre ?<\/p>\n<p>La m\u00e9thode GroupByKey est utilis\u00e9e lorsque les cl\u00e9s dans KStream sont d\u00e9j\u00e0 non vides. De plus, le drapeau \u00ab n\u00e9cessite un re-s\u00e9quen\u00e7age \u00bb n'a jamais \u00e9t\u00e9 activ\u00e9.<\/p>\n<p>La m\u00e9thode GroupBy suppose que vous avez modifi\u00e9 les cl\u00e9s pour le regroupement, donc le drapeau de re-s\u00e9quen\u00e7age est d\u00e9fini sur true. L'ex\u00e9cution apr\u00e8s la m\u00e9thode GroupBy de jointures, d'agr\u00e9gations, etc. entra\u00eenera un re-s\u00e9quen\u00e7age automatique.<br \/>\nR\u00e9sum\u00e9 : il est pr\u00e9f\u00e9rable d'utiliser GroupByKey chaque fois que possible, plut\u00f4t que GroupBy.<\/p><\/blockquote>\n<p>\nLes m\u00e9thodes mapValues et groupBy sont claires, regardons donc la m\u00e9thode sum() (que vous pouvez trouver dans le fichier src\/main\/java\/bbejeck\/model\/ShareVolume.java) (listing 5.3).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/8e9a6f873594f9b5fef9a96c42353a61.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLa m\u00e9thode ShareVolume.sum retourne la somme interm\u00e9diaire du volume des ventes d'actions, et le r\u00e9sultat de toute la cha\u00eene de calculs est un objet KTable. Maintenant, vous comprenez le r\u00f4le de KTable. En r\u00e9ceptionnant des objets ShareVolume, la derni\u00e8re mise \u00e0 jour valide est conserv\u00e9e dans l'objet KTable correspondant. Il est important de ne pas oublier que toutes les mises \u00e0 jour sont refl\u00e9t\u00e9es dans le shareVolumeKTable pr\u00e9c\u00e9dent, mais pas toutes sont transmises ult\u00e9rieurement.<\/p>\n<p>Ensuite, \u00e0 l'aide de ce KTable, nous effectuons une agr\u00e9gation (par le nombre d'actions vendues) pour obtenir les cinq entreprises avec le plus grand volume de ventes d'actions dans chaque secteur. Nos actions seront similaires aux actions lors de la premi\u00e8re agr\u00e9gation.<\/p>\n<ol>\n<li>Effectuer une autre op\u00e9ration groupBy pour regrouper les objets ShareVolume par secteurs d'activit\u00e9.<\/li>\n<li>Commencer \u00e0 sommer les objets ShareVolume. Cette fois, l'objet d'agr\u00e9gation est une file d'attente de priorit\u00e9 de taille fixe. Dans cette file d'attente de taille fixe, seules cinq entreprises avec le plus grand nombre d'actions vendues sont conserv\u00e9es.<\/li>\n<li>Convertir les files d'attente du point pr\u00e9c\u00e9dent en une valeur de cha\u00eene et retourner les cinq entreprises les plus vendues en termes d'actions par secteur d'activit\u00e9.<\/li>\n<li>\u00c9crire les r\u00e9sultats sous forme de cha\u00eene dans le topic.<\/li>\n<\/ol>\n<p>\nLa figure 5.10 montre le graphique de la topologie du mouvement des donn\u00e9es. Comme vous pouvez le voir, le deuxi\u00e8me cercle de traitement est assez simple.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/dabd1507eee267038edb7f8d76d8d8ae.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nMaintenant que vous avez bien compris la structure de ce deuxi\u00e8me cercle de traitement, vous pouvez consulter son code source (vous le trouverez dans le fichier src\/main\/java\/bbejeck\/chapter_5\/AggregationsAndReducingExample.java) (listing 5.4).<\/p>\n<p>Dans cet initialiseur, il y a une variable fixedQueue. C'est un objet utilisateur \u2014 un adaptateur pour java.util.TreeSet, qui est utilis\u00e9 pour suivre les N r\u00e9sultats les plus importants dans l'ordre d\u00e9croissant du nombre d'actions vendues.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/488b072b0d91b81c925ca72291e69e48.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nVous avez d\u00e9j\u00e0 rencontr\u00e9 les appels groupBy et mapValues, donc nous ne les aborderons pas ici (nous appelons la m\u00e9thode KTable.toStream, car la m\u00e9thode KTable.print est consid\u00e9r\u00e9e comme obsol\u00e8te). Mais vous n'avez pas encore vu la version KTable de la m\u00e9thode aggregate(), alors nous allons consacrer un peu de temps \u00e0 sa discussion.<\/p>\n<p>Comme vous le savez, KTable se distingue par le fait que les enregistrements avec des cl\u00e9s identiques sont consid\u00e9r\u00e9s comme des mises \u00e0 jour. KTable remplace l'ancien enregistrement par un nouveau. L'agr\u00e9gation se produit de la m\u00eame mani\u00e8re : les derniers enregistrements avec une cl\u00e9 unique sont agr\u00e9g\u00e9s. Lorsqu'un enregistrement arrive, il est ajout\u00e9 \u00e0 une instance de la classe FixedSizePriorityQueue \u00e0 l'aide d'un accumulateur (le deuxi\u00e8me param\u00e8tre dans l'appel de la m\u00e9thode aggregate), mais si un autre enregistrement avec la m\u00eame cl\u00e9 existe d\u00e9j\u00e0, l'ancien enregistrement est supprim\u00e9 \u00e0 l'aide d'un d\u00e9duiseur (le troisi\u00e8me param\u00e8tre dans l'appel de la m\u00e9thode aggregate).<\/p>\n<p>Tout cela signifie que notre agr\u00e9gateur, FixedSizePriorityQueue, n'agr\u00e8ge pas toutes les valeurs avec une m\u00eame cl\u00e9, mais conserve la somme glissante des quantit\u00e9s des N types d'actions les plus vendus. Chaque enregistrement entrant contient le nombre total d'actions vendues jusqu'\u00e0 pr\u00e9sent. KTable vous donnera des informations sur les actions des entreprises qui se vendent le plus en ce moment, l'agr\u00e9gation glissante de chaque mise \u00e0 jour n'est pas n\u00e9cessaire.<\/p>\n<p>Nous avons appris \u00e0 faire deux choses importantes :<\/p>\n<ul>\n<li>grouper les valeurs dans KTable par une cl\u00e9 commune ;<\/li>\n<li>effectuer sur ces valeurs regroup\u00e9es des op\u00e9rations utiles, telles que la r\u00e9duction et l'agr\u00e9gation.<\/li>\n<\/ul>\n<p>\n\u00catre capable d'effectuer ces op\u00e9rations est important pour comprendre le sens des donn\u00e9es circulant \u00e0 travers l'application Kafka Streams et d\u00e9terminer quelle information elles transportent.<\/p>\n<p>Nous avons \u00e9galement rassembl\u00e9 certains des concepts cl\u00e9s discut\u00e9s pr\u00e9c\u00e9demment dans ce livre. Au chapitre 4, nous avons \u00e9voqu\u00e9 l'importance d'un \u00e9tat local r\u00e9silient pour une application de streaming. Le premier exemple de ce chapitre a d\u00e9montr\u00e9 pourquoi l'\u00e9tat local est si important : il permet de suivre les informations que vous avez d\u00e9j\u00e0 vues. L'acc\u00e8s local permet d'\u00e9viter les latences r\u00e9seau, ce qui rend l'application plus performante et plus r\u00e9siliente aux erreurs.<\/p>\n<p>Lors de l'ex\u00e9cution de toute op\u00e9ration de r\u00e9duction ou d'agr\u00e9gation, il est n\u00e9cessaire d'indiquer le nom du magasin d'\u00e9tat. Les op\u00e9rations de r\u00e9duction et d'agr\u00e9gation retournent une instance de KTable, et KTable utilise le magasin d'\u00e9tat pour remplacer les anciens r\u00e9sultats par les nouveaux. Comme vous l'avez vu, toutes les mises \u00e0 jour ne sont pas envoy\u00e9es plus loin dans le pipeline, et c'est important car les op\u00e9rations d'agr\u00e9gation sont destin\u00e9es \u00e0 obtenir des informations finales. Si vous n'utilisez pas l'\u00e9tat local, KTable enverra toutes les r\u00e9sultats d'agr\u00e9gation et de r\u00e9duction.<\/p>\n<p>Voyons maintenant l'ex\u00e9cution d'op\u00e9rations telles que l'agr\u00e9gation sur une p\u00e9riode de temps sp\u00e9cifique \u2014 ce que l'on appelle les op\u00e9rations de fen\u00eatrage (windowing operations).<\/p>\n<h3>5.3.2. Op\u00e9rations de fen\u00eatrage<\/h3>\n<p>\nDans la section pr\u00e9c\u00e9dente, nous avons d\u00e9couvert les op\u00e9rations de r\u00e9duction et d'agr\u00e9gation \u00ab glissantes \u00bb. L'application effectuait une r\u00e9duction continue du volume des ventes d'actions, suivie de l'agr\u00e9gation des cinq actions les plus vendues sur le march\u00e9.<\/p>\n<p>Parfois, de telles agr\u00e9gations et r\u00e9ductions continues des r\u00e9sultats sont n\u00e9cessaires. D'autres fois, il est requis d'effectuer des op\u00e9rations uniquement sur une p\u00e9riode de temps d\u00e9finie. Par exemple, calculer combien d'op\u00e9rations boursi\u00e8res ont \u00e9t\u00e9 effectu\u00e9es sur les actions d'une entreprise sp\u00e9cifique au cours des 10 derni\u00e8res minutes. Ou combien d'utilisateurs ont cliqu\u00e9 sur la nouvelle banni\u00e8re publicitaire au cours des 15 derni\u00e8res minutes. L'application peut effectuer ces op\u00e9rations plusieurs fois, mais avec des r\u00e9sultats relatifs uniquement aux p\u00e9riodes de temps sp\u00e9cifi\u00e9es (fen\u00eatres temporelles).<\/p>\n<h3>Comptage des transactions boursi\u00e8res par acheteur<\/h3>\n<p>\nDans l'exemple suivant, nous allons suivre les transactions boursi\u00e8res de plusieurs traders \u2014 soit de grandes organisations, soit des financiers ind\u00e9pendants astucieux.<\/p>\n<p>Il existe deux raisons possibles pour un tel suivi. L'une d'elles est la n\u00e9cessit\u00e9 de savoir ce que les leaders du march\u00e9 ach\u00e8tent\/vendent. Si ces grands acteurs et investisseurs avertis voient des opportunit\u00e9s se pr\u00e9senter, il est judicieux de suivre leur strat\u00e9gie. La seconde raison concerne le d\u00e9sir de d\u00e9celer tout signe potentiel de transactions ill\u00e9gales utilisant des informations internes. Pour cela, il vous faudra analyser la corr\u00e9lation entre de grands pics de vente et d'importants communiqu\u00e9s de presse.<\/p>\n<p>Ce suivi comprend les \u00e9tapes suivantes :<\/p>\n<ul>\n<li>la cr\u00e9ation d'un flux de lecture \u00e0 partir du topic stock-transactions ;<\/li>\n<li>le regroupement des enregistrements entrants par identifiant d'acheteur et symbole boursier. L'appel de la m\u00e9thode groupBy retourne une instance de la classe KGroupedStream ;<\/li>\n<li>le retour de donn\u00e9es par KGroupedStream.windowedBy, un flux de donn\u00e9es limit\u00e9 par une fen\u00eatre temporelle, permettant d'effectuer une agr\u00e9gation par fen\u00eatres. Selon le type de fen\u00eatre, cela retourne soit un TimeWindowedKStream, soit un SessionWindowedKStream ;<\/li>\n<li>le comptage des transactions pour l'op\u00e9ration d'agr\u00e9gation. Le flux de donn\u00e9es par fen\u00eatre d\u00e9termine si un enregistrement sp\u00e9cifique est pris en compte dans ce comptage ;<\/li>\n<li>l'enregistrement des r\u00e9sultats dans un topic ou leur affichage dans la console pendant le d\u00e9veloppement.<\/li>\n<\/ul>\n<p>\nLa topologie de cette application est simple, mais une illustration claire n'est pas superflue. Regardons la figure 5.11.<\/p>\n<p>Nous examinerons ensuite la fonctionnalit\u00e9 des op\u00e9rations par fen\u00eatres et le code associ\u00e9.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/67d9d8d8acb1370a7c7fe5cd9a8b6aa7.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Types de fen\u00eatres<\/h3>\n<p>\nDans Kafka Streams, il existe trois types de fen\u00eatres :<\/p>\n<ul>\n<li>de session ;<\/li>\n<li>\u00abtumbling\u00bb ;<\/li>\n<li>glissantes\/\u00abhopping\u00bb.<\/li>\n<\/ul>\n<p>\nLe choix d\u00e9pend des exigences m\u00e9tiers. Les fen\u00eatres \u00abtumbling\u00bb et \u00abhopping\u00bb sont limit\u00e9es dans le temps, tandis que les limites des fen\u00eatres de session sont li\u00e9es aux actions des utilisateurs \u2014 la dur\u00e9e de la ou des sessions est d\u00e9finie uniquement par l'activit\u00e9 de l'utilisateur. Il faut garder \u00e0 l'esprit que tous les types de fen\u00eatres se basent sur les marques de date\/heure des enregistrements, et non sur l'heure syst\u00e8me.<\/p>\n<p>Nous allons ensuite r\u00e9aliser notre topologie avec chacun des types de fen\u00eatres. Le code complet sera fourni uniquement dans le premier exemple, pour les autres types de fen\u00eatres, rien ne changera, sauf le type de l'op\u00e9ration de fen\u00eatre.<\/p>\n<h3>Fen\u00eatres de session<\/h3>\n<p>\nLes fen\u00eatres de session diff\u00e8rent consid\u00e9rablement de tous les autres types de fen\u00eatres. Elles ne sont pas limit\u00e9es tant par le temps que par l'activit\u00e9 de l'utilisateur (ou l'activit\u00e9 de l'entit\u00e9 que vous souhaitez suivre). Les fen\u00eatres de session sont d\u00e9limit\u00e9es par des p\u00e9riodes d'inactivit\u00e9.<\/p>\n<p>La figure 5.12 illustre le concept des fen\u00eatres de session. Une session plus courte fusionnera avec la session \u00e0 sa gauche. En revanche, la session \u00e0 sa droite sera distincte, car elle succ\u00e8de \u00e0 une longue p\u00e9riode d'inactivit\u00e9. Les fen\u00eatres de session reposent sur les actions des utilisateurs, mais utilisent des horodatages des enregistrements pour d\u00e9terminer \u00e0 quelle session appartient un enregistrement.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/d39a5db7af7d6aa2194802622b2b47fd.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<h3>Utilisation des fen\u00eatres de session pour le suivi des transactions boursi\u00e8res<\/h3>\n<p>\nNous utiliserons les fen\u00eatres de session pour capturer des informations sur les transactions boursi\u00e8res. La mise en \u0153uvre des fen\u00eatres de session est pr\u00e9sent\u00e9e dans le listing 5.5 (qui peut \u00eatre trouv\u00e9 dans le fichier src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKTableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/2dcbd9a36baec0e746aad165121451b3.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLa plupart des op\u00e9rations de cette topologie vous sont d\u00e9j\u00e0 famili\u00e8res, donc il n'est pas n\u00e9cessaire de les examiner \u00e0 nouveau ici. Mais il y a quelques nouveaux \u00e9l\u00e9ments que nous allons discuter.<\/p>\n<p>Lors de toute op\u00e9ration groupBy, une op\u00e9ration d'agr\u00e9gation (agr\u00e9gation, r\u00e9duction ou comptage) est g\u00e9n\u00e9ralement effectu\u00e9e. On peut effectuer une agr\u00e9gation cumulative avec un total accumul\u00e9, ou une agr\u00e9gation de fen\u00eatre, qui prend en compte les enregistrements dans une fen\u00eatre temporelle d\u00e9finie.<\/p>\n<p>Le code du listing 5.5 r\u00e9alise un comptage des transactions dans les fen\u00eatres de session. \u00c0 la figure 5.13, ces actions sont analys\u00e9es \u00e9tape par \u00e9tape.<\/p>\n<p>Avec l'appel windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)), nous cr\u00e9ons une fen\u00eatre de session avec un intervalle d'inactivit\u00e9 de 20 secondes et un intervalle de conservation de 15 minutes. Un intervalle d'inactivit\u00e9 de 20 secondes signifie que l'application inclura tout enregistrement re\u00e7u dans les 20 secondes suivant la fin ou le d\u00e9but de la session actuelle dans la session actuelle (active).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/9bd47b04698086872fd135b4c67eb938.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEnsuite, nous pr\u00e9cisons quelle op\u00e9ration d'agr\u00e9gation doit \u00eatre effectu\u00e9e dans la fen\u00eatre de session \u2014 dans ce cas, le comptage. Si l'enregistrement entrant d\u00e9passe le d\u00e9lai d'inactivit\u00e9 (de chaque c\u00f4t\u00e9 de l'horodatage), l'application cr\u00e9e une nouvelle session. L'intervalle de conservation signifie maintenir la session pour une certaine p\u00e9riode et permet des donn\u00e9es tardives d\u00e9passant la p\u00e9riode d'inactivit\u00e9 de la session, mais qui peuvent encore \u00eatre jointes. De plus, le d\u00e9but et la fin de la nouvelle session r\u00e9sultant de la fusion correspondent \u00e0 l'horodatage le plus ancien et le plus r\u00e9cent.<\/p>\n<p>Examinons quelques enregistrements de la m\u00e9thode de comptage pour voir comment les sessions fonctionnent (tableau 5.1).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/ec04aae466d88c2d2349474069c8d541.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLors de la r\u00e9ception des enregistrements, nous recherchons les sessions existantes avec la m\u00eame cl\u00e9, dont le temps de fin est inf\u00e9rieur \u00e0 l'horodatage actuel \u2014 le d\u00e9lai d'inactivit\u00e9 \u2014 et le temps de d\u00e9but est sup\u00e9rieur \u00e0 l'horodatage actuel + le d\u00e9lai d'inactivit\u00e9. En tenant compte de cela, quatre enregistrements du tableau 5.1 sont fusionn\u00e9s en une seule session comme suit.<\/p>\n<p>1. Le premier enregistrement est re\u00e7u, donc le temps de d\u00e9but est \u00e9gal \u00e0 l'heure de fin et est 00:00:00.<\/p>\n<p>2. Ensuite, l'enregistrement 2 arrive, et nous recherchons des sessions se terminant au plus tard \u00e0 23:59:55 et commen\u00e7ant au plus tard \u00e0 00:00:35. Nous trouvons l'enregistrement 1 et fusionnons les sessions 1 et 2. Nous prenons le temps de d\u00e9but de la session 1 (le plus ancien) et le temps de fin de la session 2 (le plus r\u00e9cent), de sorte que notre nouvelle session commence \u00e0 00:00:00 et se termine \u00e0 00:00:15.<\/p>\n<p>3. L'enregistrement 3 arrive, nous recherchons des sessions entre 00:00:30 et 00:01:10 et ne trouvons aucune. Nous ajoutons la deuxi\u00e8me session pour la cl\u00e9 123-345-654,FFBE, commen\u00e7ant et se terminant \u00e0 00:00:50.<\/p>\n<p>4. L'enregistrement 4 arrive, et nous recherchons des sessions entre 23:59:45 et 00:00:25. Cette fois, les deux sessions \u2014 1 et 2 \u2014 sont trouv\u00e9es. Les trois sessions sont toutes fusionn\u00e9es en une seule, avec un temps de d\u00e9but \u00e0 00:00:00 et un temps de fin \u00e0 00:00:15.<\/p>\n<p>Il est important de retenir les points suivants tir\u00e9s de cette section :<\/p>\n<ul>\n<li>les sessions ne sont pas des fen\u00eatres de taille fixe. La dur\u00e9e d'une session est d\u00e9termin\u00e9e par l'activit\u00e9 dans la p\u00e9riode de temps donn\u00e9e ;<\/li>\n<li>les horodatages dans les donn\u00e9es d\u00e9terminent si un \u00e9v\u00e9nement appartient \u00e0 une session existante ou \u00e0 une p\u00e9riode d'inactivit\u00e9.<\/li>\n<\/ul>\n<p>\nNous discuterons ensuite d'une autre variante des fen\u00eatres \u2014 les fen\u00eatres \u00ab roulantes \u00bb.<\/p>\n<h3>Les fen\u00eatres \u00ab roulantes \u00bb<\/h3>\n<p>\nLes fen\u00eatres \u00ab tumbling \u00bb capturent des \u00e9v\u00e9nements qui se produisent au cours d'un intervalle de temps d\u00e9fini. Imaginez que vous devez capturer toutes les transactions boursi\u00e8res d'une entreprise toutes les 20 secondes, de sorte que vous collectez tous les \u00e9v\u00e9nements pendant cette p\u00e9riode. \u00c0 la fin de l'intervalle de 20 secondes, la fen\u00eatre \u00ab tumble \u00bb et passe \u00e0 un nouvel intervalle d'observation de 20 secondes. La figure 5.14 illustre cette situation.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/e90e560d9ddda5e2e524b7c387ad9874.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nComme vous pouvez le voir, tous les \u00e9v\u00e9nements re\u00e7us au cours des 20 derni\u00e8res secondes sont inclus dans la fen\u00eatre. \u00c0 la fin de cette p\u00e9riode, une nouvelle fen\u00eatre est cr\u00e9\u00e9e.<\/p>\n<p>Le listing 5.6 montre le code illustrant l'utilisation des fen\u00eatres \u00ab tumbling \u00bb pour capturer toutes les 20 secondes les transactions boursi\u00e8res (vous pouvez le trouver dans le fichier src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/e055bb1b288c7d500b64372fe3fbf064.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nGr\u00e2ce \u00e0 ce petit changement dans l'appel de la m\u00e9thode TimeWindows.of, il est possible d'utiliser une fen\u00eatre \u00ab tumbling \u00bb. Dans cet exemple, il n'y a pas d'appel \u00e0 la m\u00e9thode until(), ce qui signifie que l'intervalle de conservation par d\u00e9faut, qui est de 24 heures, sera utilis\u00e9.<\/p>\n<p>Enfin, il est temps de passer \u00e0 la derni\u00e8re variante de fen\u00eatres : les fen\u00eatres \u00ab hopping \u00bb.<\/p>\n<h3>Fen\u00eatres glissantes (\u00ab hopping \u00bb)<\/h3>\n<p>\nLes fen\u00eatres glissantes \u00ab hopping \u00bb sont similaires aux fen\u00eatres \u00ab tumbling \u00bb, mais avec une l\u00e9g\u00e8re diff\u00e9rence. Les fen\u00eatres glissantes ne attendent pas la fin du temps imparti avant de cr\u00e9er une nouvelle fen\u00eatre pour traiter les \u00e9v\u00e9nements r\u00e9cents. Elles lancent de nouveaux calculs apr\u00e8s un intervalle d'attente qui est plus court que la dur\u00e9e de la fen\u00eatre.<\/p>\n<p>Pour illustrer les diff\u00e9rences entre les fen\u00eatres \u00ab tumbling \u00bb et \u00ab hopping \u00bb, revenons \u00e0 l'exemple de comptage des transactions boursi\u00e8res. Notre objectif reste de compter le nombre de transactions, mais nous ne voulons pas attendre toute la dur\u00e9e avant de mettre \u00e0 jour le compteur. Au lieu de cela, nous allons mettre \u00e0 jour le compteur par intervalles de temps plus courts. Par exemple, nous continuerons de compter le nombre de transactions toutes les 20 secondes, mais nous mettrons \u00e0 jour le compteur toutes les 5 secondes, comme le montre la fig. 5.15. Cela nous donne trois fen\u00eatres de r\u00e9sultats avec des donn\u00e9es qui se chevauchent.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/415c8cd9f2b60d453a1a01c3bc99331f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nLe listing 5.7 montre le code pour d\u00e9finir les fen\u00eatres glissantes (vous pouvez le trouver dans le fichier src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/ab2d1a64380d256d3fb084e16597417c.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nUne fen\u00eatre \u00ab titubante \u00bb peut \u00eatre transform\u00e9e en fen\u00eatre \u00ab sautante \u00bb en ajoutant l'appel de la m\u00e9thode advanceBy(). Dans l'exemple donn\u00e9, l'intervalle de conservation est de 15 minutes.<\/p>\n<p>Vous avez vu dans cette section comment limiter les r\u00e9sultats d'agr\u00e9gation par des fen\u00eatres temporelles. En particulier, voici les trois choses que vous devriez retenir de cette section :<\/p>\n<ul>\n<li>la taille des fen\u00eatres de session n'est pas limit\u00e9e par une p\u00e9riode de temps, mais par l'activit\u00e9 des utilisateurs ;<\/li>\n<li>les fen\u00eatres \u00ab titubantes \u00bb fournissent un aper\u00e7u des \u00e9v\u00e9nements sur une p\u00e9riode donn\u00e9e ;<\/li>\n<li>la dur\u00e9e de vie des fen\u00eatres \u00ab sautantes \u00bb est fixe, mais elles sont souvent mises \u00e0 jour et peuvent contenir des enregistrements qui se chevauchent dans toutes les fen\u00eatres.<\/li>\n<\/ul>\n<p>\nNous allons maintenant voir comment transformer un KTable en KStream pour effectuer une jointure.<\/p>\n<h3>5.3.3. Jointure d'objets KStream et KTable<\/h3>\n<p>\nDans le chapitre 4, nous avons discut\u00e9 de la jointure de deux objets KStream. Maintenant, nous allons apprendre \u00e0 joindre KTable et KStream. Cela peut \u00eatre n\u00e9cessaire pour la simple raison suivante. KStream est un flux d\u2019enregistrements, et KTable est un flux de mises \u00e0 jour d\u2019enregistrements, mais parfois, il peut \u00eatre n\u00e9cessaire d'ajouter un contexte suppl\u00e9mentaire au flux d'enregistrements \u00e0 l'aide de mises \u00e0 jour de KTable.<\/p>\n<p>Prenons des donn\u00e9es sur le nombre de transactions boursi\u00e8res et joignons-les aux actualit\u00e9s boursi\u00e8res selon les secteurs d'activit\u00e9 concern\u00e9s. Voici ce que nous devons faire pour atteindre cet objectif en tenant compte du code existant.<\/p>\n<ol>\n<li>Transformer l'objet KTable contenant les donn\u00e9es sur le nombre de transactions boursi\u00e8res en KStream tout en rempla\u00e7ant la cl\u00e9 par une cl\u00e9 indiquant le secteur d'activit\u00e9 correspondant \u00e0 ce symbole boursier.<\/li>\n<li>Cr\u00e9er un objet KTable en lisant les donn\u00e9es d'un sujet sur les actualit\u00e9s boursi\u00e8res. Ce nouveau KTable sera cat\u00e9goris\u00e9 par secteurs d'activit\u00e9.<\/li>\n<li>Joindre les mises \u00e0 jour des actualit\u00e9s avec les informations sur le nombre de transactions boursi\u00e8res par secteur d'activit\u00e9.<\/li>\n<\/ol>\n<p>\nMaintenant, voyons comment mettre en \u0153uvre ce plan d'action.<\/p>\n<h3>Transformation de KTable en KStream<\/h3>\n<p>\nPour transformer un KTable en KStream, il est n\u00e9cessaire de proc\u00e9der comme suit.<\/p>\n<ol>\n<li>Appeler la m\u00e9thode KTable.toStream().<\/li>\n<li>En appelant la m\u00e9thode KStream.map, remplacer la cl\u00e9 par le nom du secteur d'activit\u00e9, puis extraire \u00e0 partir de l'instance Windowed l'objet TransactionSummary.<\/li>\n<\/ol>\n<p>\nNous allons encha\u00eener ces op\u00e9rations de la mani\u00e8re suivante (le code peut \u00eatre trouv\u00e9 dans le fichier src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listing 5.8).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/0d43c2650f6e66e2816ed383da3a29c2.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nPuisque nous effectuons l'op\u00e9ration KStream.map, le re-partitionnement pour l'instance KStream retourn\u00e9e est automatiquement effectu\u00e9 lors de son utilisation dans une jointure.<\/p>\n<p>Nous avons termin\u00e9 le processus de transformation, maintenant nous devons cr\u00e9er un objet KTable pour lire les nouvelles du march\u00e9.<\/p>\n<h3>Cr\u00e9er un KTable pour les nouvelles du march\u00e9.<\/h3>\n<p>\nHeureusement, cr\u00e9er un objet KTable n\u00e9cessite juste une ligne de code (ce code peut \u00eatre trouv\u00e9 dans le fichier src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listing 5.9).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/6e83a393fdc9ab74fda4cbdddddb5213.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nIl convient de noter qu'aucun objet Serde n'est n\u00e9cessaire car des Serde de cha\u00eenes sont utilis\u00e9s dans les configurations. De plus, gr\u00e2ce \u00e0 l'utilisation de l'\u00e9num\u00e9ration EARLIEST, la table est remplie avec des enregistrements d\u00e8s le d\u00e9but.<\/p>\n<p>Nous pouvons maintenant passer \u00e0 l'\u00e9tape finale \u2014 la jointure.<\/p>\n<h3>Jointure des mises \u00e0 jour d'actualit\u00e9s avec les donn\u00e9es sur le nombre de transactions.<\/h3>\n<p>\nLa cr\u00e9ation de la jointure ne pose pas de difficult\u00e9s. Nous utiliserons une jointure gauche au cas o\u00f9 il n'y aurait pas de nouvelles du march\u00e9 pour l'industrie correspondante (le code requis peut \u00eatre trouv\u00e9 dans le fichier src\/main\/java\/bbejeck\/chapter_5\/CountingWindowingAndKtableJoinExample.java) (listing 5.10).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/92ed70f98927d2f778ad14dbc2a5aa26.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nCette op\u00e9ration leftJoin est assez simple. Contrairement aux jointures du chapitre 4, la m\u00e9thode JoinWindow n'est pas utilis\u00e9e car lors de la jointure KStream-KTable, il n'y a qu'un seul enregistrement pour chaque cl\u00e9 dans le KTable. Cette jointure n'est pas limit\u00e9e dans le temps : l'enregistrement est soit dans le KTable, soit il n'y est pas. En r\u00e9sum\u00e9, avec des objets KTable, on peut enrichir KStream avec des donn\u00e9es de r\u00e9f\u00e9rence moins fr\u00e9quemment mises \u00e0 jour.<\/p>\n<p>Passons maintenant \u00e0 une mani\u00e8re plus efficace d'enrichir les \u00e9v\u00e9nements provenant de KStream.<\/p>\n<h3>5.3.4. Objets GlobalKTable.<\/h3>\n<p>\nComme vous l'avez compris, il existe un besoin d'enrichir les flux d'\u00e9v\u00e9nements ou d'y ajouter du contexte. Dans le chapitre 4, vous avez vu des jointures entre deux objets KStream, et dans la section pr\u00e9c\u00e9dente, la jointure entre KStream et KTable. Dans tous ces cas, il est n\u00e9cessaire de re-partitionner le flux de donn\u00e9es lors de la mappage des cl\u00e9s \u00e0 un nouveau type ou valeur. Parfois, le re-partitionnement est effectu\u00e9 explicitement, et parfois Kafka Streams le fait automatiquement. Le re-partitionnement est n\u00e9cessaire car les cl\u00e9s ont chang\u00e9 et les enregistrements doivent se trouver dans de nouvelles partitions, sinon la jointure sera impossible (ce qui a \u00e9t\u00e9 discut\u00e9 au chapitre 4, section \"Re-partitionnement des donn\u00e9es\", sous-section 4.2.4).<\/p>\n<h3>La re-s\u00e9curisation a un co\u00fbt<\/h3>\n<p>\nLa re-s\u00e9curisation n\u00e9cessite des co\u00fbts suppl\u00e9mentaires \u2014 des ressources suppl\u00e9mentaires pour cr\u00e9er des sujets interm\u00e9diaires, stocker des donn\u00e9es dupliqu\u00e9es dans un autre sujet ; cela implique \u00e9galement une augmentation de la latence due aux op\u00e9rations d'\u00e9criture et de lecture de ce sujet. De plus, si vous devez effectuer des jointures sur plus d'un aspect ou dimension, il faut organiser ces jointures en cha\u00eene, mapper les enregistrements avec de nouvelles cl\u00e9s et r\u00e9p\u00e9ter le processus de re-s\u00e9curisation.<\/p>\n<h3>Jointure avec des ensembles de donn\u00e9es de plus petite taille<\/h3>\n<p>\nDans certains cas, le volume de donn\u00e9es de r\u00e9f\u00e9rence pr\u00e9vu pour la jointure est relativement faible, de sorte que les copies compl\u00e8tes peuvent tenir localement sur chacun des n\u0153uds. Pour de telles situations, Kafka Streams pr\u00e9voit une classe GlobalKTable.<\/p>\n<p>Les instances de GlobalKTable sont uniques, car l'application r\u00e9plique toutes les donn\u00e9es sur chaque n\u0153ud. \u00c9tant donn\u00e9 que toutes les donn\u00e9es sont pr\u00e9sentes sur chaque n\u0153ud, il n'est pas n\u00e9cessaire de partitionner le flux d'\u00e9v\u00e9nements par la cl\u00e9 des donn\u00e9es de r\u00e9f\u00e9rence pour qu'il soit accessible \u00e0 toutes les partitions. Avec les objets GlobalKTable, il est \u00e9galement possible d'effectuer des jointures sans cl\u00e9. Revenons \u00e0 l'un des exemples pr\u00e9c\u00e9dents pour illustrer cette capacit\u00e9.<\/p>\n<h3>Jointure d'objets KStream avec des objets GlobalKTable<\/h3>\n<p>\nDans la sous-section 5.3.2, nous avons effectu\u00e9 une agr\u00e9gation temporelle des transactions boursi\u00e8res par client. Les r\u00e9sultats de cette agr\u00e9gation ressemblaient \u00e0 ceci :<\/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>\nBien que ces r\u00e9sultats correspondent \u00e0 l'objectif vis\u00e9, il serait plus pratique d'afficher \u00e9galement le nom du client et le nom complet de l'entreprise. Pour ajouter le nom de l'acheteur et le nom de l'entreprise, des jointures standard peuvent \u00eatre effectu\u00e9es, mais cela n\u00e9cessitera deux mappings de cl\u00e9s et une re-partition. Avec GlobalKTable, ces op\u00e9rations peuvent \u00eatre \u00e9vit\u00e9es.<\/p>\n<p>Pour ce faire, nous allons utiliser l'objet countStream de l'extrait 5.11 (le code correspondant peut \u00eatre trouv\u00e9 dans le fichier src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java), en le joignant \u00e0 deux objets GlobalKTable.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/fc4d91bbe062ceb94f5650224840b81e.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nNous en avons d\u00e9j\u00e0 discut\u00e9 pr\u00e9c\u00e9demment, donc je ne vais pas me r\u00e9p\u00e9ter. Mais je tiens \u00e0 noter que le code dans la fonction toStream().map a \u00e9t\u00e9 abstrait dans un objet-fonction pour des raisons de lisibilit\u00e9, au lieu d'utiliser une expression lambda int\u00e9gr\u00e9e.<\/p>\n<p>L'\u00e9tape suivante consiste \u00e0 d\u00e9clarer deux instances de GlobalKTable (le code correspondant se trouve dans le fichier src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (listing 5.12).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/db3918895f174c8cfb5fc927b55f87a1.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p>Notez que les noms des topics sont d\u00e9crits \u00e0 l'aide de types \u00e9num\u00e9r\u00e9s.<\/p>\n<p>Maintenant que nous avons pr\u00e9par\u00e9 tous les composants, il ne reste plus qu'\u00e0 \u00e9crire le code pour la jointure (que vous pouvez trouver dans le fichier src\/main\/java\/bbejeck\/chapter_5\/GlobalKTableExample.java) (listing 5.13).<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/799360cc99f1920c190a61fd4685d4ff.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nBien que ce code contienne deux jointures, elles sont organis\u00e9es sous forme de cha\u00eene, car aucun des r\u00e9sultats n'est utilis\u00e9 individuellement. Les r\u00e9sultats sont affich\u00e9s \u00e0 la fin de l'op\u00e9ration enti\u00e8re.<\/p>\n<p>Lorsque vous ex\u00e9cuterez l'op\u00e9ration de jointure ci-dessus, vous obtiendrez des r\u00e9sultats de ce type :<\/p>\n<pre><code class=\"plaintext\">{customer='Barney, Smith' company=\"Exxon\", transactions= 17}<\/code><\/pre>\n<p>\nL'essence n'a pas chang\u00e9, mais ces r\u00e9sultats sont plus lisibles.<\/p>\n<p>Si l'on consid\u00e8re le chapitre 4, vous avez d\u00e9j\u00e0 vu plusieurs types de jointures en action. Elles sont r\u00e9pertori\u00e9es dans le tableau 5.2. Ce tableau refl\u00e8te les capacit\u00e9s de jointure pertinentes pour la version 1.0.0 de Kafka Streams ; dans les prochaines versions, il pourrait y avoir des changements.<\/p>\n<p><img decoding=\"async\" alt=\"Le livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb\" src=\"\/wp-content\/uploads\/8e4cf35c64a8431bda43a5e748de275f.jpg\" style=\"display:block;margin: 0 auto;\" \/><br \/>\nEn conclusion, rappelons l'essentiel : vous pouvez joindre des flux d'\u00e9v\u00e9nements (KStream) et des flux de mises \u00e0 jour (KTable) \u00e0 l'aide d'un \u00e9tat local. De plus, si la taille des donn\u00e9es de r\u00e9f\u00e9rence n'est pas trop grande, vous pouvez utiliser l'objet GlobalKTable. GlobalKTable r\u00e9plique toutes les partitions sur chacun des n\u0153uds de l'application Kafka Streams, garantissant ainsi l'accessibilit\u00e9 de toutes les donn\u00e9es, quel que soit le partitionnement associ\u00e9 \u00e0 la cl\u00e9.<\/p>\n<p>Nous allons maintenant d\u00e9couvrir une fonctionnalit\u00e9 de Kafka Streams qui permet d'observer les changements d'\u00e9tat sans consommer de donn\u00e9es depuis un topic Kafka.<\/p>\n<h3>5.3.5. \u00c9tat accessible aux requ\u00eates<\/h3>\n<p>\nNous avons d\u00e9j\u00e0 r\u00e9alis\u00e9 plusieurs op\u00e9rations impliquant un \u00e9tat et avons toujours affich\u00e9 les r\u00e9sultats dans la console (\u00e0 des fins de d\u00e9veloppement) ou les avons enregistr\u00e9s dans un topic (pour des fins de production). Lors de l'enregistrement des r\u00e9sultats dans un topic, il est n\u00e9cessaire d'utiliser un consommateur Kafka pour les consulter.<\/p>\n<p>La lecture des donn\u00e9es \u00e0 partir de ces sujets peut \u00eatre consid\u00e9r\u00e9e comme une forme de vues mat\u00e9rialis\u00e9es. Pour nos besoins, nous pouvons utiliser la d\u00e9finition d'une vue mat\u00e9rialis\u00e9e provenant de \u00ab Wikip\u00e9dia \u00bb : \u00ab ... un objet physique de base de donn\u00e9es contenant les r\u00e9sultats de l'ex\u00e9cution d'une requ\u00eate. Par exemple, il peut s'agir d'une copie locale de donn\u00e9es distantes, ou d'un sous-ensemble de lignes et\/ou de colonnes d'une table ou des r\u00e9sultats d'une jointure, ou d'un tableau crois\u00e9 obtenu par agr\u00e9gation \u00bb (https:\/\/fr.wikipedia.org\/wiki\/Vue_matrialise).<\/p>\n<p>Kafka Streams permet \u00e9galement d'effectuer des requ\u00eates interactives sur les magasins d'\u00e9tat, ce qui permet de lire directement ces vues mat\u00e9rialis\u00e9es. Il est important de noter que les requ\u00eates sur les magasins d'\u00e9tat sont des op\u00e9rations en lecture seule. Cela signifie que vous n'avez pas \u00e0 craindre de rendre l'\u00e9tat incoh\u00e9rent lors du traitement des donn\u00e9es par l'application.<\/p>\n<p>La possibilit\u00e9 d'effectuer des requ\u00eates directement sur les magasins d'\u00e9tat est tr\u00e8s importante. Cela signifie qu'il est possible de cr\u00e9er des applications \u2014 des tableaux de bord \u2014 sans avoir \u00e0 d'abord r\u00e9cup\u00e9rer les donn\u00e9es du consommateur Kafka. Cela am\u00e9liore \u00e9galement l'efficacit\u00e9 de l'application, car il n'est pas n\u00e9cessaire de r\u00e9\u00e9crire les donn\u00e9es :<\/p>\n<ul>\n<li>en raison de la localit\u00e9 des donn\u00e9es, on peut y acc\u00e9der rapidement ;<\/li>\n<li>la duplication des donn\u00e9es est \u00e9vit\u00e9e, car elles ne sont pas enregistr\u00e9es dans un stockage externe.<\/li>\n<\/ul>\n<p>\nLe principal point que je voudrais que vous reteniez est qu'il est possible d'ex\u00e9cuter des requ\u00eates directement sur l'\u00e9tat depuis l'application. On ne saurait trop souligner les capacit\u00e9s que cela vous offre. Plut\u00f4t que de consommer des donn\u00e9es \u00e0 partir de Kafka et de stocker des enregistrements dans une base de donn\u00e9es pour l'application, on peut effectuer des requ\u00eates sur les magasins d'\u00e9tat avec le m\u00eame r\u00e9sultat. Les requ\u00eates directes sur les magasins d'\u00e9tat signifient moins de code (absence de consommateur) et moins de logiciels (absence de besoin d'une table de base de donn\u00e9es pour stocker les r\u00e9sultats).<\/p>\n<p>Nous avons couvert une quantit\u00e9 importante d'informations dans ce chapitre, nous allons donc suspendre notre discussion sur les requ\u00eates interactives aux magasins d'\u00e9tat pour le moment. Mais ne vous inqui\u00e9tez pas : au chapitre 9, nous allons cr\u00e9er une simple application - un tableau de bord avec des requ\u00eates interactives. Pour d\u00e9montrer les requ\u00eates interactives et les possibilit\u00e9s de les int\u00e9grer dans les applications Kafka Streams, nous utiliserons certains des exemples de ce chapitre et du pr\u00e9c\u00e9dent.<\/p>\n<h3>R\u00e9sum\u00e9<\/h3>\n<p><\/p>\n<ul>\n<li>Les objets KStream repr\u00e9sentent des flux d'\u00e9v\u00e9nements, comparables \u00e0 des inserts dans une base de donn\u00e9es. Les objets KTable repr\u00e9sentent des flux de mises \u00e0 jour, ils ressemblent davantage \u00e0 des mises \u00e0 jour dans une base de donn\u00e9es. La taille de l'objet KTable ne cro\u00eet pas, les anciennes entr\u00e9es sont remplac\u00e9es par les nouvelles.<\/li>\n<li>Les objets KTable sont n\u00e9cessaires pour les op\u00e9rations d'agr\u00e9gation.<\/li>\n<li>Avec les op\u00e9rations de fen\u00eatres, vous pouvez segmenter les donn\u00e9es agr\u00e9g\u00e9es en fonction de plages de temps.<\/li>\n<li>Gr\u00e2ce aux objets GlobalKTable, vous pouvez acc\u00e9der aux donn\u00e9es de r\u00e9f\u00e9rence \u00e0 n'importe quel endroit de l'application, ind\u00e9pendamment du partitionnement par sections.<\/li>\n<li>Des jointures peuvent avoir lieu entre les objets KStream, KTable et GlobalKTable.<\/li>\n<\/ul>\n<p>\nJusqu'\u00e0 pr\u00e9sent, nous avons concentr\u00e9 notre attention sur la cr\u00e9ation d'applications Kafka Streams en utilisant le DSL KStream de haut niveau. Bien que l'approche de haut niveau permette de cr\u00e9er des programmes propres et concis, son utilisation repr\u00e9sente un compromis certain. Travailler avec le DSL KStream signifie r\u00e9duire la concision du code au profit d'un contr\u00f4le moindre. Dans le chapitre suivant, nous examinerons l'API bas niveau des n\u0153uds de traitement et essayerons d'autres compromis. Les programmes deviendront plus longs que ce qu'ils \u00e9taient jusqu'\u00e0 pr\u00e9sent, mais nous aurons la possibilit\u00e9 de cr\u00e9er pratiquement n'importe quel n\u0153ud de traitement dont nous pourrions avoir besoin.<\/p>\n<p>\u2192 Vous pouvez en savoir plus sur le livre sur <noindex><a rel=\"nofollow\" href=\"https:\/\/www.piter.com\/collection\/best\/product\/kafka-streams-v-deystvii-prilozheniya-i-mikroservisy-dlya-raboty-v-realnom-vremeni\">le site de l'\u00e9diteur<\/a><\/noindex><\/p>\n<p>\u2192 Pour les membres de Habr, une r\u00e9duction de 25% avec le coupon - <b>Kafka Streams<\/b><\/p>\n<p>\u2192 Apr\u00e8s le paiement de la version papier du livre, un livre \u00e9lectronique est envoy\u00e9 par e-mail.<br \/>\n<br \/>Source : <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.2 - 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\/fr\/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.2\" \/>\n\t\t<meta property=\"og:locale\" content=\"fr_FR\" \/>\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\/fr\/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\udd47Livre \u00ab Kafka Streams en action. Applications et microservices pour le traitement en temps r\u00e9el \u00bb | ProHoster","description":"","canonical_url":"https:\/\/prohoster.info\/fr\/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":"fr_FR","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\/fr\/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\/fr\/wp-json\/wp\/v2\/posts\/35786","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/comments?post=35786"}],"version-history":[{"count":0,"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/posts\/35786\/revisions"}],"wp:attachment":[{"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/media?parent=35786"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/categories?post=35786"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/prohoster.info\/fr\/wp-json\/wp\/v2\/tags?post=35786"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}