Bonjour, membres de Habr ! Ce livre est destinĂ© Ă tout dĂ©veloppeur souhaitant comprendre le traitement de flux. MaĂźtriser la programmation distribuĂ©e facilitera l'apprentissage de Kafka et de Kafka Streams. Il serait utile de connaĂźtre le framework Kafka, mais ce n'est pas obligatoire : je vous expliquerai tout ce qu'il faut savoir. Tant les dĂ©veloppeurs expĂ©rimentĂ©s que les dĂ©butants pourront grĂące Ă ce livre maĂźtriser la crĂ©ation d'applications intĂ©ressantes pour le traitement de flux en utilisant la bibliothĂšque Kafka Streams. Les dĂ©veloppeurs Java de niveau intermĂ©diaire et avancĂ©, dĂ©jĂ familiarisĂ©s avec des concepts tels que la sĂ©rialisation, apprendront Ă appliquer leurs compĂ©tences pour dĂ©velopper des applications Kafka Streams. Le code source du livre est Ă©crit en Java 8 et utilise de maniĂšre significative la syntaxe des expressions lambda de Java 8, donc ĂȘtre capable de travailler avec des fonctions lambda (mĂȘme dans un autre langage de programmation) vous sera utile.
Extrait. 5.3. AgrĂ©gation et opĂ©rations fenĂȘtres
Dans cette section, nous allons explorer les parties les plus prometteuses de Kafka Streams. Jusqu'à présent, nous avons examiné les aspects suivants de Kafka Streams :
- création de topologies de traitement ;
- utilisation des états dans les applications de flux ;
- exécution des jointures de flux de données ;
- différences entre les flux d'événements (KStream) et les flux de mise à jour (KTable).
Dans les exemples suivants, nous allons rassembler tous ces Ă©lĂ©ments. De plus, vous dĂ©couvrirez les opĂ©rations fenĂȘtres â une autre excellente fonctionnalitĂ© des applications de flux. Notre premier exemple sera une agrĂ©gation simple.
5.3.1. Agrégation du volume des ventes d'actions par secteur industriel
L'agrégation et le regroupement sont des outils essentiels lors du travail avec des données de flux. L'examen des enregistrements individuels à mesure qu'ils arrivent s'avÚre souvent insuffisant. Pour extraire des informations supplémentaires à partir des données, il est nécessaire de les regrouper et de les combiner.
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écisément, vous vous intéresserez aux cinq entreprises ayant les volumes de vente d'actions les plus élevés dans chacun des secteurs industriels.
Pour réaliser une telle agrégation, plusieurs étapes seront nécessaires pour convertir les données au format requis (en termes généraux).
- Créer une source basée sur un sujet publiant des informations brutes sur le commerce des actions. Nous devrons mapper l'objet de type StockTransaction à un objet de type ShareVolume. Le fait est que l'objet StockTransaction contient des métadonnées sur les ventes, et nous avons uniquement besoin des données concernant le nombre d'actions vendues.
- Regrouper les données ShareVolume par symboles d'actions. AprÚs avoir regroupé par symboles, il est possible de réduire ces données à des totaux intermédiaires des volumes de ventes d'actions. Il convient de noter que la méthode KStream.groupBy renvoie une instance de type KGroupedStream. Et pour obtenir une instance de KTable, il suffit d'appeler ensuite la méthode KGroupedStream.reduce.
Qu'est-ce que l'interface KGroupedStream ?
Les mĂ©thodes KStream.groupBy et KStream.groupByKey renvoient une instance de KGroupedStream. KGroupedStream est une reprĂ©sentation intermĂ©diaire d'un flux d'Ă©vĂ©nements aprĂšs regroupement par clĂ©s. Il n'est pas du tout destinĂ© Ă ĂȘtre utilisĂ© directement. Au lieu de cela, KGroupedStream est utilisĂ© pour des opĂ©rations d'agrĂ©gation, dont le rĂ©sultat est toujours un KTable. Et puisque le rĂ©sultat des opĂ©rations d'agrĂ©gation est un KTable et que le stockage d'Ă©tat est utilisĂ©, il se peut que toutes les mises Ă jour ne soient pas envoyĂ©es plus loin dans le pipeline.
La mĂ©thode KTable.groupBy renvoie un KGroupedTable similaire â une reprĂ©sentation intermĂ©diaire d'un flux de mises Ă jour regroupĂ©es par clĂ©.
Faisons une petite pause et regardons la fig. 5.9, qui montre ce que nous avons atteint. Cette topologie devrait dĂ©jĂ vous ĂȘtre familiĂšre.

Jetons maintenant un coup d'Ćil au code pour cette topologie (vous pouvez le trouver dans le fichier src/main/java/bbejeck/chapter_5/AggregationsAndReducingExample.java) (listing 5.2).

Le code fourni se distingue par sa concision et le grand nombre d'actions effectuées en quelques lignes. Dans le premier paramÚtre de la méthode builder.stream, vous pouvez remarquer quelque chose de nouveau : la valeur de l'énumération AutoOffsetReset.EARLIEST (il existe également LATEST), définie à l'aide de la méthode Consumed.withOffsetResetPolicy. Avec ce type d'énumération, vous pouvez spécifier la stratégie de réinitialisation des offsets pour chaque KStream ou KTable, qui a la priorité sur le paramÚtre de réinitialisation des offsets de la configuration.
GroupByKey et GroupBy
Dans l'interface KStream, il existe deux méthodes pour regrouper des enregistrements : GroupByKey et GroupBy. Les deux renvoient un KGroupedTable, vous pouvez donc vous poser la question : quelle est la différence entre elles et quand utiliser l'une ou l'autre ?
La méthode GroupByKey est utilisée lorsque les clés dans KStream sont déjà non vides. De plus, le drapeau « nécessite un re-séquençage » n'a jamais été activé.
La méthode GroupBy suppose que vous avez modifié les clés pour le regroupement, donc le drapeau de re-séquençage est défini sur true. L'exécution aprÚs la méthode GroupBy de jointures, d'agrégations, etc. entraßnera un re-séquençage automatique.
Résumé : il est préférable d'utiliser GroupByKey chaque fois que possible, plutÎt que GroupBy.
Les méthodes mapValues et groupBy sont claires, regardons donc la méthode sum() (que vous pouvez trouver dans le fichier src/main/java/bbejeck/model/ShareVolume.java) (listing 5.3).

La méthode ShareVolume.sum retourne la somme intermédiaire du volume des ventes d'actions, et le résultat de toute la chaßne de calculs est un objet KTable. Maintenant, vous comprenez le rÎle de KTable. En réceptionnant des objets ShareVolume, la derniÚre mise à jour valide est conservée dans l'objet KTable correspondant. Il est important de ne pas oublier que toutes les mises à jour sont reflétées dans le shareVolumeKTable précédent, mais pas toutes sont transmises ultérieurement.
Ensuite, à l'aide de ce KTable, nous effectuons une agrégation (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Úre agrégation.
- Effectuer une autre opération groupBy pour regrouper les objets ShareVolume par secteurs d'activité.
- Commencer à sommer les objets ShareVolume. Cette fois, l'objet d'agrégation est une file d'attente de priorité de taille fixe. Dans cette file d'attente de taille fixe, seules cinq entreprises avec le plus grand nombre d'actions vendues sont conservées.
- Convertir les files d'attente du point précédent en une valeur de chaßne et retourner les cinq entreprises les plus vendues en termes d'actions par secteur d'activité.
- Ăcrire les rĂ©sultats sous forme de chaĂźne dans le topic.
La figure 5.10 montre le graphique de la topologie du mouvement des données. Comme vous pouvez le voir, le deuxiÚme cercle de traitement est assez simple.

Maintenant que vous avez bien compris la structure de ce deuxiĂšme 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).
Dans cet initialiseur, il y a une variable fixedQueue. C'est un objet utilisateur â un adaptateur pour java.util.TreeSet, qui est utilisĂ© pour suivre les N rĂ©sultats les plus importants dans l'ordre dĂ©croissant du nombre d'actions vendues.

Vous avez déjà rencontré les appels groupBy et mapValues, donc nous ne les aborderons pas ici (nous appelons la méthode KTable.toStream, car la méthode KTable.print est considérée comme obsolÚte). Mais vous n'avez pas encore vu la version KTable de la méthode aggregate(), alors nous allons consacrer un peu de temps à sa discussion.
Comme vous le savez, KTable se distingue par le fait que les enregistrements avec des clĂ©s identiques sont considĂ©rĂ©s comme des mises Ă jour. KTable remplace l'ancien enregistrement par un nouveau. L'agrĂ©gation se produit de la mĂȘme maniĂšre : les derniers enregistrements avec une clĂ© unique sont agrĂ©gĂ©s. Lorsqu'un enregistrement arrive, il est ajoutĂ© Ă une instance de la classe FixedSizePriorityQueue Ă l'aide d'un accumulateur (le deuxiĂšme paramĂštre dans l'appel de la mĂ©thode aggregate), mais si un autre enregistrement avec la mĂȘme clĂ© existe dĂ©jĂ , l'ancien enregistrement est supprimĂ© Ă l'aide d'un dĂ©duiseur (le troisiĂšme paramĂštre dans l'appel de la mĂ©thode aggregate).
Tout cela signifie que notre agrĂ©gateur, FixedSizePriorityQueue, n'agrĂšge pas toutes les valeurs avec une mĂȘme clĂ©, mais conserve la somme glissante des quantitĂ©s des N types d'actions les plus vendus. Chaque enregistrement entrant contient le nombre total d'actions vendues jusqu'Ă prĂ©sent. KTable vous donnera des informations sur les actions des entreprises qui se vendent le plus en ce moment, l'agrĂ©gation glissante de chaque mise Ă jour n'est pas nĂ©cessaire.
Nous avons appris Ă faire deux choses importantes :
- grouper les valeurs dans KTable par une clé commune ;
- effectuer sur ces valeurs regroupées des opérations utiles, telles que la réduction et l'agrégation.
Ătre capable d'effectuer ces opĂ©rations est important pour comprendre le sens des donnĂ©es circulant Ă travers l'application Kafka Streams et dĂ©terminer quelle information elles transportent.
Nous avons également rassemblé certains des concepts clés discutés précédemment dans ce livre. Au chapitre 4, nous avons évoqué l'importance d'un état local résilient pour une application de streaming. Le premier exemple de ce chapitre a démontré pourquoi l'état local est si important : il permet de suivre les informations que vous avez déjà vues. L'accÚs local permet d'éviter les latences réseau, ce qui rend l'application plus performante et plus résiliente aux erreurs.
Lors de l'exécution de toute opération de réduction ou d'agrégation, il est nécessaire d'indiquer le nom du magasin d'état. Les opérations de réduction et d'agrégation retournent une instance de KTable, et KTable utilise le magasin d'état pour remplacer les anciens résultats par les nouveaux. Comme vous l'avez vu, toutes les mises à jour ne sont pas envoyées plus loin dans le pipeline, et c'est important car les opérations d'agrégation sont destinées à obtenir des informations finales. Si vous n'utilisez pas l'état local, KTable enverra toutes les résultats d'agrégation et de réduction.
Voyons maintenant l'exĂ©cution d'opĂ©rations telles que l'agrĂ©gation sur une pĂ©riode de temps spĂ©cifique â ce que l'on appelle les opĂ©rations de fenĂȘtrage (windowing operations).
5.3.2. OpĂ©rations de fenĂȘtrage
Dans la section précédente, nous avons découvert les opérations de réduction et d'agrégation « glissantes ». L'application effectuait une réduction continue du volume des ventes d'actions, suivie de l'agrégation des cinq actions les plus vendues sur le marché.
Parfois, de telles agrĂ©gations et rĂ©ductions continues des rĂ©sultats sont nĂ©cessaires. D'autres fois, il est requis d'effectuer des opĂ©rations uniquement sur une pĂ©riode de temps dĂ©finie. Par exemple, calculer combien d'opĂ©rations boursiĂšres ont Ă©tĂ© effectuĂ©es sur les actions d'une entreprise spĂ©cifique au cours des 10 derniĂšres minutes. Ou combien d'utilisateurs ont cliquĂ© sur la nouvelle banniĂšre publicitaire au cours des 15 derniĂšres minutes. L'application peut effectuer ces opĂ©rations plusieurs fois, mais avec des rĂ©sultats relatifs uniquement aux pĂ©riodes de temps spĂ©cifiĂ©es (fenĂȘtres temporelles).
Comptage des transactions boursiĂšres par acheteur
Dans l'exemple suivant, nous allons suivre les transactions boursiĂšres de plusieurs traders â soit de grandes organisations, soit des financiers indĂ©pendants astucieux.
Il existe deux raisons possibles pour un tel suivi. L'une d'elles est la nécessité de savoir ce que les leaders du marché achÚtent/vendent. Si ces grands acteurs et investisseurs avertis voient des opportunités se présenter, il est judicieux de suivre leur stratégie. La seconde raison concerne le désir de déceler tout signe potentiel de transactions illégales utilisant des informations internes. Pour cela, il vous faudra analyser la corrélation entre de grands pics de vente et d'importants communiqués de presse.
Ce suivi comprend les étapes suivantes :
- la création d'un flux de lecture à partir du topic stock-transactions ;
- le regroupement des enregistrements entrants par identifiant d'acheteur et symbole boursier. L'appel de la méthode groupBy retourne une instance de la classe KGroupedStream ;
- le retour de donnĂ©es par KGroupedStream.windowedBy, un flux de donnĂ©es limitĂ© par une fenĂȘtre temporelle, permettant d'effectuer une agrĂ©gation par fenĂȘtres. Selon le type de fenĂȘtre, cela retourne soit un TimeWindowedKStream, soit un SessionWindowedKStream ;
- le comptage des transactions pour l'opĂ©ration d'agrĂ©gation. Le flux de donnĂ©es par fenĂȘtre dĂ©termine si un enregistrement spĂ©cifique est pris en compte dans ce comptage ;
- l'enregistrement des résultats dans un topic ou leur affichage dans la console pendant le développement.
La topologie de cette application est simple, mais une illustration claire n'est pas superflue. Regardons la figure 5.11.
Nous examinerons ensuite la fonctionnalitĂ© des opĂ©rations par fenĂȘtres et le code associĂ©.

Types de fenĂȘtres
Dans Kafka Streams, il existe trois types de fenĂȘtres :
- de session ;
- «tumbling» ;
- glissantes/«hopping».
Le choix dĂ©pend des exigences mĂ©tiers. Les fenĂȘtres «tumbling» et «hopping» sont limitĂ©es dans le temps, tandis que les limites des fenĂȘtres de session sont liĂ©es aux actions des utilisateurs â la durĂ©e de la ou des sessions est dĂ©finie uniquement par l'activitĂ© de l'utilisateur. Il faut garder Ă l'esprit que tous les types de fenĂȘtres se basent sur les marques de date/heure des enregistrements, et non sur l'heure systĂšme.
Nous allons ensuite rĂ©aliser notre topologie avec chacun des types de fenĂȘtres. Le code complet sera fourni uniquement dans le premier exemple, pour les autres types de fenĂȘtres, rien ne changera, sauf le type de l'opĂ©ration de fenĂȘtre.
FenĂȘtres de session
Les fenĂȘtres de session diffĂšrent considĂ©rablement de tous les autres types de fenĂȘtres. Elles ne sont pas limitĂ©es tant par le temps que par l'activitĂ© de l'utilisateur (ou l'activitĂ© de l'entitĂ© que vous souhaitez suivre). Les fenĂȘtres de session sont dĂ©limitĂ©es par des pĂ©riodes d'inactivitĂ©.
La figure 5.12 illustre le concept des fenĂȘtres de session. Une session plus courte fusionnera avec la session Ă sa gauche. En revanche, la session Ă sa droite sera distincte, car elle succĂšde Ă une longue pĂ©riode d'inactivitĂ©. Les fenĂȘtres de session reposent sur les actions des utilisateurs, mais utilisent des horodatages des enregistrements pour dĂ©terminer Ă quelle session appartient un enregistrement.

Utilisation des fenĂȘtres de session pour le suivi des transactions boursiĂšres
Nous utiliserons les fenĂȘtres de session pour capturer des informations sur les transactions boursiĂšres. La mise en Ćuvre des fenĂȘtres de session est prĂ©sentĂ©e dans le listing 5.5 (qui peut ĂȘtre trouvĂ© dans le fichier src/main/java/bbejeck/chapter_5/CountingWindowingAndKTableJoinExample.java).

La plupart des opérations de cette topologie vous sont déjà familiÚres, donc il n'est pas nécessaire de les examiner à nouveau ici. Mais il y a quelques nouveaux éléments que nous allons discuter.
Lors de toute opĂ©ration groupBy, une opĂ©ration d'agrĂ©gation (agrĂ©gation, rĂ©duction ou comptage) est gĂ©nĂ©ralement effectuĂ©e. On peut effectuer une agrĂ©gation cumulative avec un total accumulĂ©, ou une agrĂ©gation de fenĂȘtre, qui prend en compte les enregistrements dans une fenĂȘtre temporelle dĂ©finie.
Le code du listing 5.5 rĂ©alise un comptage des transactions dans les fenĂȘtres de session. Ă la figure 5.13, ces actions sont analysĂ©es Ă©tape par Ă©tape.
Avec l'appel windowedBy(SessionWindows.with(twentySeconds).until(fifteenMinutes)), nous crĂ©ons une fenĂȘtre de session avec un intervalle d'inactivitĂ© de 20 secondes et un intervalle de conservation de 15 minutes. Un intervalle d'inactivitĂ© de 20 secondes signifie que l'application inclura tout enregistrement reçu dans les 20 secondes suivant la fin ou le dĂ©but de la session actuelle dans la session actuelle (active).

Ensuite, nous prĂ©cisons quelle opĂ©ration d'agrĂ©gation doit ĂȘtre effectuĂ©e dans la fenĂȘtre de session â dans ce cas, le comptage. Si l'enregistrement entrant dĂ©passe le dĂ©lai d'inactivitĂ© (de chaque cĂŽtĂ© de l'horodatage), l'application crĂ©e une nouvelle session. L'intervalle de conservation signifie maintenir la session pour une certaine pĂ©riode et permet des donnĂ©es tardives dĂ©passant la pĂ©riode d'inactivitĂ© de la session, mais qui peuvent encore ĂȘtre jointes. De plus, le dĂ©but et la fin de la nouvelle session rĂ©sultant de la fusion correspondent Ă l'horodatage le plus ancien et le plus rĂ©cent.
Examinons quelques enregistrements de la méthode de comptage pour voir comment les sessions fonctionnent (tableau 5.1).

Lors de la rĂ©ception des enregistrements, nous recherchons les sessions existantes avec la mĂȘme clĂ©, dont le temps de fin est infĂ©rieur Ă l'horodatage actuel â le dĂ©lai d'inactivitĂ© â et le temps de dĂ©but est supĂ©rieur Ă l'horodatage actuel + le dĂ©lai d'inactivitĂ©. En tenant compte de cela, quatre enregistrements du tableau 5.1 sont fusionnĂ©s en une seule session comme suit.
1. Le premier enregistrement est reçu, donc le temps de début est égal à l'heure de fin et est 00:00:00.
2. Ensuite, l'enregistrement 2 arrive, et nous recherchons des sessions se terminant au plus tard à 23:59:55 et commençant au plus tard à 00:00:35. Nous trouvons l'enregistrement 1 et fusionnons les sessions 1 et 2. Nous prenons le temps de début de la session 1 (le plus ancien) et le temps de fin de la session 2 (le plus récent), de sorte que notre nouvelle session commence à 00:00:00 et se termine à 00:00:15.
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Úme session pour la clé 123-345-654,FFBE, commençant et se terminant à 00:00:50.
4. L'enregistrement 4 arrive, et nous recherchons des sessions entre 23:59:45 et 00:00:25. Cette fois, les deux sessions â 1 et 2 â sont trouvĂ©es. Les trois sessions sont toutes fusionnĂ©es en une seule, avec un temps de dĂ©but Ă 00:00:00 et un temps de fin Ă 00:00:15.
Il est important de retenir les points suivants tirés de cette section :
- les sessions ne sont pas des fenĂȘtres de taille fixe. La durĂ©e d'une session est dĂ©terminĂ©e par l'activitĂ© dans la pĂ©riode de temps donnĂ©e ;
- les horodatages dans les données déterminent si un événement appartient à une session existante ou à une période d'inactivité.
Nous discuterons ensuite d'une autre variante des fenĂȘtres â les fenĂȘtres « roulantes ».
Les fenĂȘtres « roulantes »
Les fenĂȘtres « tumbling » capturent des Ă©vĂ©nements qui se produisent au cours d'un intervalle de temps dĂ©fini. Imaginez que vous devez capturer toutes les transactions boursiĂšres d'une entreprise toutes les 20 secondes, de sorte que vous collectez tous les Ă©vĂ©nements pendant cette pĂ©riode. Ă la fin de l'intervalle de 20 secondes, la fenĂȘtre « tumble » et passe Ă un nouvel intervalle d'observation de 20 secondes. La figure 5.14 illustre cette situation.

Comme vous pouvez le voir, tous les Ă©vĂ©nements reçus au cours des 20 derniĂšres secondes sont inclus dans la fenĂȘtre. Ă la fin de cette pĂ©riode, une nouvelle fenĂȘtre est créée.
Le listing 5.6 montre le code illustrant l'utilisation des fenĂȘtres « tumbling » pour capturer toutes les 20 secondes les transactions boursiĂšres (vous pouvez le trouver dans le fichier src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

GrĂące Ă ce petit changement dans l'appel de la mĂ©thode TimeWindows.of, il est possible d'utiliser une fenĂȘtre « tumbling ». Dans cet exemple, il n'y a pas d'appel Ă la mĂ©thode until(), ce qui signifie que l'intervalle de conservation par dĂ©faut, qui est de 24 heures, sera utilisĂ©.
Enfin, il est temps de passer Ă la derniĂšre variante de fenĂȘtres : les fenĂȘtres « hopping ».
FenĂȘtres glissantes (« hopping »)
Les fenĂȘtres glissantes « hopping » sont similaires aux fenĂȘtres « tumbling », mais avec une lĂ©gĂšre diffĂ©rence. Les fenĂȘtres glissantes ne attendent pas la fin du temps imparti avant de crĂ©er une nouvelle fenĂȘtre pour traiter les Ă©vĂ©nements rĂ©cents. Elles lancent de nouveaux calculs aprĂšs un intervalle d'attente qui est plus court que la durĂ©e de la fenĂȘtre.
Pour illustrer les diffĂ©rences entre les fenĂȘtres « tumbling » et « hopping », revenons Ă l'exemple de comptage des transactions boursiĂšres. Notre objectif reste de compter le nombre de transactions, mais nous ne voulons pas attendre toute la durĂ©e avant de mettre Ă jour le compteur. Au lieu de cela, nous allons mettre Ă 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 Ă jour le compteur toutes les 5 secondes, comme le montre la fig. 5.15. Cela nous donne trois fenĂȘtres de rĂ©sultats avec des donnĂ©es qui se chevauchent.

Le listing 5.7 montre le code pour dĂ©finir les fenĂȘtres glissantes (vous pouvez le trouver dans le fichier src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java).

Une fenĂȘtre « titubante » peut ĂȘtre transformĂ©e en fenĂȘtre « sautante » en ajoutant l'appel de la mĂ©thode advanceBy(). Dans l'exemple donnĂ©, l'intervalle de conservation est de 15 minutes.
Vous avez vu dans cette section comment limiter les rĂ©sultats d'agrĂ©gation par des fenĂȘtres temporelles. En particulier, voici les trois choses que vous devriez retenir de cette section :
- la taille des fenĂȘtres de session n'est pas limitĂ©e par une pĂ©riode de temps, mais par l'activitĂ© des utilisateurs ;
- les fenĂȘtres « titubantes » fournissent un aperçu des Ă©vĂ©nements sur une pĂ©riode donnĂ©e ;
- la durĂ©e de vie des fenĂȘtres « sautantes » est fixe, mais elles sont souvent mises Ă jour et peuvent contenir des enregistrements qui se chevauchent dans toutes les fenĂȘtres.
Nous allons maintenant voir comment transformer un KTable en KStream pour effectuer une jointure.
5.3.3. Jointure d'objets KStream et KTable
Dans le chapitre 4, nous avons discutĂ© de la jointure de deux objets KStream. Maintenant, nous allons apprendre Ă joindre KTable et KStream. Cela peut ĂȘtre nĂ©cessaire pour la simple raison suivante. KStream est un flux dâenregistrements, et KTable est un flux de mises Ă jour dâenregistrements, mais parfois, il peut ĂȘtre nĂ©cessaire d'ajouter un contexte supplĂ©mentaire au flux d'enregistrements Ă l'aide de mises Ă jour de KTable.
Prenons des données sur le nombre de transactions boursiÚres et joignons-les aux actualités boursiÚres selon les secteurs d'activité concernés. Voici ce que nous devons faire pour atteindre cet objectif en tenant compte du code existant.
- Transformer l'objet KTable contenant les données sur le nombre de transactions boursiÚres en KStream tout en remplaçant la clé par une clé indiquant le secteur d'activité correspondant à ce symbole boursier.
- Créer un objet KTable en lisant les données d'un sujet sur les actualités boursiÚres. Ce nouveau KTable sera catégorisé par secteurs d'activité.
- Joindre les mises à jour des actualités avec les informations sur le nombre de transactions boursiÚres par secteur d'activité.
Maintenant, voyons comment mettre en Ćuvre ce plan d'action.
Transformation de KTable en KStream
Pour transformer un KTable en KStream, il est nécessaire de procéder comme suit.
- Appeler la méthode KTable.toStream().
- En appelant la méthode KStream.map, remplacer la clé par le nom du secteur d'activité, puis extraire à partir de l'instance Windowed l'objet TransactionSummary.
Nous allons enchaĂźner ces opĂ©rations de la maniĂšre suivante (le code peut ĂȘtre trouvĂ© dans le fichier src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (listing 5.8).

Puisque nous effectuons l'opération KStream.map, le re-partitionnement pour l'instance KStream retournée est automatiquement effectué lors de son utilisation dans une jointure.
Nous avons terminé le processus de transformation, maintenant nous devons créer un objet KTable pour lire les nouvelles du marché.
Créer un KTable pour les nouvelles du marché.
Heureusement, crĂ©er un objet KTable nĂ©cessite juste une ligne de code (ce code peut ĂȘtre trouvĂ© dans le fichier src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (listing 5.9).

Il convient de noter qu'aucun objet Serde n'est nécessaire car des Serde de chaßnes sont utilisés dans les configurations. De plus, grùce à l'utilisation de l'énumération EARLIEST, la table est remplie avec des enregistrements dÚs le début.
Nous pouvons maintenant passer Ă l'Ă©tape finale â la jointure.
Jointure des mises à jour d'actualités avec les données sur le nombre de transactions.
La crĂ©ation de la jointure ne pose pas de difficultĂ©s. Nous utiliserons une jointure gauche au cas oĂč il n'y aurait pas de nouvelles du marchĂ© pour l'industrie correspondante (le code requis peut ĂȘtre trouvĂ© dans le fichier src/main/java/bbejeck/chapter_5/CountingWindowingAndKtableJoinExample.java) (listing 5.10).

Cette opération leftJoin est assez simple. Contrairement aux jointures du chapitre 4, la méthode JoinWindow n'est pas utilisée car lors de la jointure KStream-KTable, il n'y a qu'un seul enregistrement pour chaque clé dans le KTable. Cette jointure n'est pas limitée dans le temps : l'enregistrement est soit dans le KTable, soit il n'y est pas. En résumé, avec des objets KTable, on peut enrichir KStream avec des données de référence moins fréquemment mises à jour.
Passons maintenant à une maniÚre plus efficace d'enrichir les événements provenant de KStream.
5.3.4. Objets GlobalKTable.
Comme vous l'avez compris, il existe un besoin d'enrichir les flux d'événements ou d'y ajouter du contexte. Dans le chapitre 4, vous avez vu des jointures entre deux objets KStream, et dans la section précédente, la jointure entre KStream et KTable. Dans tous ces cas, il est nécessaire de re-partitionner le flux de données lors de la mappage des clés à un nouveau type ou valeur. Parfois, le re-partitionnement est effectué explicitement, et parfois Kafka Streams le fait automatiquement. Le re-partitionnement est nécessaire car les clés ont changé et les enregistrements doivent se trouver dans de nouvelles partitions, sinon la jointure sera impossible (ce qui a été discuté au chapitre 4, section "Re-partitionnement des données", sous-section 4.2.4).
La re-sécurisation a un coût
La re-sĂ©curisation nĂ©cessite des coĂ»ts supplĂ©mentaires â des ressources supplĂ©mentaires pour crĂ©er des sujets intermĂ©diaires, stocker des donnĂ©es dupliquĂ©es dans un autre sujet ; cela implique Ă©galement une augmentation de la latence due aux opĂ©rations d'Ă©criture 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Ăźne, mapper les enregistrements avec de nouvelles clĂ©s et rĂ©pĂ©ter le processus de re-sĂ©curisation.
Jointure avec des ensembles de données de plus petite taille
Dans certains cas, le volume de donnĂ©es de rĂ©fĂ©rence prĂ©vu pour la jointure est relativement faible, de sorte que les copies complĂštes peuvent tenir localement sur chacun des nĆuds. Pour de telles situations, Kafka Streams prĂ©voit une classe GlobalKTable.
Les instances de GlobalKTable sont uniques, car l'application rĂ©plique toutes les donnĂ©es sur chaque nĆud. Ătant donnĂ© que toutes les donnĂ©es sont prĂ©sentes sur chaque nĆud, il n'est pas nĂ©cessaire de partitionner le flux d'Ă©vĂ©nements par la clĂ© des donnĂ©es de rĂ©fĂ©rence pour qu'il soit accessible Ă toutes les partitions. Avec les objets GlobalKTable, il est Ă©galement possible d'effectuer des jointures sans clĂ©. Revenons Ă l'un des exemples prĂ©cĂ©dents pour illustrer cette capacitĂ©.
Jointure d'objets KStream avec des objets GlobalKTable
Dans la sous-section 5.3.2, nous avons effectué une agrégation temporelle des transactions boursiÚres par client. Les résultats de cette agrégation ressemblaient à ceci :
{customerId='074-09-3705', stockTicker='GUTM'}, 17
{customerId='037-34-5184', stockTicker='CORK'}, 16Bien que ces rĂ©sultats correspondent Ă l'objectif visĂ©, il serait plus pratique d'afficher Ă©galement 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 ĂȘtre effectuĂ©es, mais cela nĂ©cessitera deux mappings de clĂ©s et une re-partition. Avec GlobalKTable, ces opĂ©rations peuvent ĂȘtre Ă©vitĂ©es.
Pour ce faire, nous allons utiliser l'objet countStream de l'extrait 5.11 (le code correspondant peut ĂȘtre trouvĂ© dans le fichier src/main/java/bbejeck/chapter_5/GlobalKTableExample.java), en le joignant Ă deux objets GlobalKTable.

Nous en avons déjà discuté précédemment, donc je ne vais pas me répéter. Mais je tiens à noter que le code dans la fonction toStream().map a été abstrait dans un objet-fonction pour des raisons de lisibilité, au lieu d'utiliser une expression lambda intégrée.
L'étape suivante consiste à déclarer deux instances de GlobalKTable (le code correspondant se trouve dans le fichier src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (listing 5.12).

Notez que les noms des topics sont décrits à l'aide de types énumérés.
Maintenant que nous avons préparé tous les composants, il ne reste plus qu'à écrire le code pour la jointure (que vous pouvez trouver dans le fichier src/main/java/bbejeck/chapter_5/GlobalKTableExample.java) (listing 5.13).

Bien que ce code contienne deux jointures, elles sont organisées sous forme de chaßne, car aucun des résultats n'est utilisé individuellement. Les résultats sont affichés à la fin de l'opération entiÚre.
Lorsque vous exécuterez l'opération de jointure ci-dessus, vous obtiendrez des résultats de ce type :
{customer='Barney, Smith' company="Exxon", transactions= 17}L'essence n'a pas changé, mais ces résultats sont plus lisibles.
Si l'on considÚre le chapitre 4, vous avez déjà vu plusieurs types de jointures en action. Elles sont répertoriées dans le tableau 5.2. Ce tableau reflÚte les capacités de jointure pertinentes pour la version 1.0.0 de Kafka Streams ; dans les prochaines versions, il pourrait y avoir des changements.

En conclusion, rappelons l'essentiel : vous pouvez joindre des flux d'Ă©vĂ©nements (KStream) et des flux de mises Ă jour (KTable) Ă l'aide d'un Ă©tat local. De plus, si la taille des donnĂ©es de rĂ©fĂ©rence n'est pas trop grande, vous pouvez utiliser l'objet GlobalKTable. GlobalKTable rĂ©plique toutes les partitions sur chacun des nĆuds de l'application Kafka Streams, garantissant ainsi l'accessibilitĂ© de toutes les donnĂ©es, quel que soit le partitionnement associĂ© Ă la clĂ©.
Nous allons maintenant découvrir une fonctionnalité de Kafka Streams qui permet d'observer les changements d'état sans consommer de données depuis un topic Kafka.
5.3.5. Ătat accessible aux requĂȘtes
Nous avons déjà réalisé plusieurs opérations impliquant un état et avons toujours affiché les résultats dans la console (à des fins de développement) ou les avons enregistrés dans un topic (pour des fins de production). Lors de l'enregistrement des résultats dans un topic, il est nécessaire d'utiliser un consommateur Kafka pour les consulter.
La lecture des donnĂ©es Ă partir de ces sujets peut ĂȘtre considĂ©rĂ©e comme une forme de vues matĂ©rialisĂ©es. Pour nos besoins, nous pouvons utiliser la dĂ©finition d'une vue matĂ©rialisĂ©e provenant de « WikipĂ©dia » : « ... un objet physique de base de donnĂ©es contenant les rĂ©sultats de l'exĂ©cution d'une requĂȘte. Par exemple, il peut s'agir d'une copie locale de donnĂ©es distantes, ou d'un sous-ensemble de lignes et/ou de colonnes d'une table ou des rĂ©sultats d'une jointure, ou d'un tableau croisĂ© obtenu par agrĂ©gation » (https://fr.wikipedia.org/wiki/Vue_matrialise).
Kafka Streams permet Ă©galement d'effectuer des requĂȘtes interactives sur les magasins d'Ă©tat, ce qui permet de lire directement ces vues matĂ©rialisĂ©es. Il est important de noter que les requĂȘtes sur les magasins d'Ă©tat sont des opĂ©rations en lecture seule. Cela signifie que vous n'avez pas Ă craindre de rendre l'Ă©tat incohĂ©rent lors du traitement des donnĂ©es par l'application.
La possibilitĂ© d'effectuer des requĂȘtes directement sur les magasins d'Ă©tat est trĂšs importante. Cela signifie qu'il est possible de crĂ©er des applications â des tableaux de bord â sans avoir Ă d'abord rĂ©cupĂ©rer les donnĂ©es du consommateur Kafka. Cela amĂ©liore Ă©galement l'efficacitĂ© de l'application, car il n'est pas nĂ©cessaire de réécrire les donnĂ©es :
- en raison de la localité des données, on peut y accéder rapidement ;
- la duplication des données est évitée, car elles ne sont pas enregistrées dans un stockage externe.
Le principal point que je voudrais que vous reteniez est qu'il est possible d'exĂ©cuter des requĂȘtes directement sur l'Ă©tat depuis l'application. On ne saurait trop souligner les capacitĂ©s que cela vous offre. PlutĂŽt que de consommer des donnĂ©es Ă partir de Kafka et de stocker des enregistrements dans une base de donnĂ©es pour l'application, on peut effectuer des requĂȘtes sur les magasins d'Ă©tat avec le mĂȘme rĂ©sultat. Les requĂȘtes directes sur les magasins d'Ă©tat signifient moins de code (absence de consommateur) et moins de logiciels (absence de besoin d'une table de base de donnĂ©es pour stocker les rĂ©sultats).
Nous avons couvert une quantitĂ© importante d'informations dans ce chapitre, nous allons donc suspendre notre discussion sur les requĂȘtes interactives aux magasins d'Ă©tat pour le moment. Mais ne vous inquiĂ©tez pas : au chapitre 9, nous allons crĂ©er une simple application - un tableau de bord avec des requĂȘtes interactives. Pour dĂ©montrer les requĂȘtes interactives et les possibilitĂ©s de les intĂ©grer dans les applications Kafka Streams, nous utiliserons certains des exemples de ce chapitre et du prĂ©cĂ©dent.
Résumé
- Les objets KStream représentent des flux d'événements, comparables à des inserts dans une base de données. Les objets KTable représentent des flux de mises à jour, ils ressemblent davantage à des mises à jour dans une base de données. La taille de l'objet KTable ne croßt pas, les anciennes entrées sont remplacées par les nouvelles.
- Les objets KTable sont nécessaires pour les opérations d'agrégation.
- Avec les opĂ©rations de fenĂȘtres, vous pouvez segmenter les donnĂ©es agrĂ©gĂ©es en fonction de plages de temps.
- Grùce aux objets GlobalKTable, vous pouvez accéder aux données de référence à n'importe quel endroit de l'application, indépendamment du partitionnement par sections.
- Des jointures peuvent avoir lieu entre les objets KStream, KTable et GlobalKTable.
Jusqu'Ă prĂ©sent, nous avons concentrĂ© notre attention sur la crĂ©ation d'applications Kafka Streams en utilisant le DSL KStream de haut niveau. Bien que l'approche de haut niveau permette de crĂ©er des programmes propres et concis, son utilisation reprĂ©sente un compromis certain. Travailler avec le DSL KStream signifie rĂ©duire la concision du code au profit d'un contrĂŽle moindre. Dans le chapitre suivant, nous examinerons l'API bas niveau des nĆuds de traitement et essayerons d'autres compromis. Les programmes deviendront plus longs que ce qu'ils Ă©taient jusqu'Ă prĂ©sent, mais nous aurons la possibilitĂ© de crĂ©er pratiquement n'importe quel nĆud de traitement dont nous pourrions avoir besoin.
â Vous pouvez en savoir plus sur le livre sur
â Pour les membres de Habr, une rĂ©duction de 25% avec le coupon - Kafka Streams
â AprĂšs le paiement de la version papier du livre, un livre Ă©lectronique est envoyĂ© par e-mail.
Source : habr.com
