Comment Glean utilise Apache Beam Streaming pour lire efficacement un vaste corpus depuis MySQL

0
minutes de lecture
Comment Glean utilise Apache Beam Streaming pour lire efficacement un vaste corpus depuis MySQL

Table des matières

Vous avez des questions ou souhaitez une démo ?

Nous sommes là pour vous aider ! Cliquez sur le bouton ci-dessous et nous vous recontacterons.

Demander une démo
Partager cet article :

Chez Glean, nous explorons chaque jour des corpus de plusieurs centaines de millions de documents. Une fois tous ces documents explorés, nous devons nous assurer qu’ils sont en permanence analysés et traités afin de maintenir à jour le contenu, les autorisations et les statistiques, pour offrir la meilleure expérience de recherche possible.  ;

Exécuter un job batch pour cet usage était coûteux, fragile et difficile à intégrer avec MySQL. Avec un job en streaming, en revanche, nous avons pu retraiter chaque document au moins une fois par jour. Dans cet article de blog, nous expliquons comment cela a été rendu possible grâce à Cloud SQL (MySQL) et Apache Beam sur Google Cloud Dataflow. Tous les pipelines décrits ci-dessous utilisent un seul worker, en configuration multi-cœurs, sur Google Cloud Platform (GCP), ce qui offre une scalabilité verticale massive, suffisante pour nos cas d’usage.

Lorsqu’un nouveau document est exploré, le contenu brut est écrit dans une table MySQL. Nous maintenons également une file Pub/Sub qui est notifiée avec l’ID unique du document nouvellement exploré. Ce document doit ensuite être traité pour analyser et extraire des informations structurées, lesquelles sont réécrites dans la même table MySQL et indexées par notre système de recherche. La file Pub/Sub garantit que cela se fait en temps réel. Les documents doivent aussi être retraités périodiquement afin de gérer les événements éventuellement perdus par la file Pub/Sub et de recalculer les statistiques qui évoluent au fil du temps, à mesure que le corpus change. Cela est réalisé en scannant et en retraitant en continu les données de la table MySQL. Dans la suite de cet article, nous nous concentrerons principalement sur la manière dont le système lit et traite les données depuis MySQL.

Pour lire les documents depuis la table MySQL, nous utilisons un pipeline en streaming. Deux facteurs principaux ont influencé notre choix du streaming plutôt que du batch : la scalabilité et le coût. Avec un corpus en croissance continue, un job batch hebdomadaire n’était pas scalable, et l’utilisation d’un pipeline en streaming nous a permis de traiter chaque document au moins une fois par jour. Une part importante de ce gain de coûts vient du fait que nous devions déjà maintenir un pipeline en streaming pour gérer la file Pub/Sub. Ces machines étaient sous-utilisées et pouvaient donc être exploitées pour effectuer un scan en streaming de la base MySQL, sans coût supplémentaire.

Illustration du produit
Figure 1 : Architecture du pipeline

Illustration du produit
Figure 2 : Cadence de traitement des documents sur un seul cœur

Pour mettre en place un pipeline en streaming (voir Figure 1), nous utilisons Apache Beam avec le runner Google Cloud Dataflow. Nous définissons une UnboundedSource personnalisée, dont chaque instance crée un UnboundedReader personnalisé qui interroge un cache de documents afin de récupérer le prochain document à traiter. Ce cache de documents est maintenu par une instance statique d’un objet scanner. Plusieurs lectures parallèles de la base sont émises par ce scanner vers la table MySQL, en fonction du nombre de cœurs du worker utilisé. Ces lectures récupèrent les documents par lots, ordonnés par ID de document, afin de garantir que chaque document est traité au plus une fois lors de chaque scan complet du corpus. Avec cette configuration, nous atteignons des cadences de traitement de l’ordre de plusieurs centaines de documents par seconde sur un seul cœur (voir Figure 2).  ;

Le scanner suit également l’état du cache : l’ID du dernier document lu, le nombre de documents restant dans le cache, etc. Il s’appuie sur ces informations pour déterminer quand récupérer des documents supplémentaires, ainsi que pour déterminer quand un scan complet du corpus est terminé. Il réinitialise ensuite son état et recommence à scanner le corpus depuis le début. Cela met ainsi en place une file de données infinie pour notre pipeline en streaming.  ;

Cette conception et ce travail d’optimisation nous ont permis de garantir une solution à la fois scalable et rentable, du plus petit au plus grand de nos clients. Nous itérons en permanence pour que nos utilisateurs puissent trouver les documents dont ils ont besoin, au moment où ils en ont besoin, afin d’avancer dans leur travail.  ;

Nous entrerons dans le détail de requêtes SQL supplémentaires et d’optimisations dans un prochain article, alors restez à l’écoute ! Si vous avez trouvé cet article intéressant et que vous souhaitez travailler sur ce type de systèmes, contactez-nous !

Work AI qui fonctionne.

CTA Background Gradient 3CTA Background Gradient 3CTA Background Mobile