Revenir au contenu principal
Produit

Comment Databricks Feature Store fournit des caractéristiques avec une fraîcheur inférieure à la seconde

Comment Databricks Feature Store fournit des caractéristiques avec une fraîcheur inférieure à la seconde

par Ian Ackerman, Nick Joung et Abhay Bothra

  • Databricks Feature Store apporte une fraîcheur en temps réel aux caractéristiques ML : les agrégations en streaming depuis Kafka peuvent désormais atteindre le Feature Store en ligne avec une latence p99 de 200 ms, réduisant la latence des caractéristiques de quelques minutes ou heures à des millisecondes.
  • Le mode Spark Real-Time (RTM) rend possible le calcul de caractéristiques en quelques millisecondes : le RTM traite les lignes en continu au lieu d'attendre des micro-lots, met à jour les agrégations de fenêtres glissantes par événement et amortit le checkpointing pour maintenir une faible latence de streaming avec état.
  • Lakebase permet des écritures de caractéristiques en ligne à haut débit : la séparation des couches de calcul et de stockage réduit l'amplification d'écriture pour les petits upserts fréquents, rendant les valeurs de caractéristiques fraîches rapidement disponibles pour l'inférence de modèles à faible latence.

Les modèles de machine learning ne valent que par les signaux qu'ils reçoivent. Un cas d'usage de détection de la fraude doit décider en quelques millisecondes après qu'un utilisateur a cliqué sur "acheter" s'il convient d'autoriser la transaction. Prendre la bonne décision dépend de la détection d'une transaction suspecte survenue il y a seulement quelques secondes. Combiner la moyenne des transactions d'un utilisateur sur les 30 derniers jours avec le montant total des transactions des 10 dernières minutes permet de mettre en évidence une fraude potentielle. Les agrégations à long terme établissent un profil de référence de l'utilisateur pour déterminer ce qui est normal, tandis que les données les plus récentes aident à faire ressortir tout comportement anormal au moment même où il se produit. La personnalisation est confrontée à la même pression : les signaux les plus récents sont ceux qui capturent l'intention actuelle de l'utilisateur et stimulent l'engagement.

Les pipelines Spark sont un moyen éprouvé de traiter des données en masse dans le Lakehouse pour obtenir des caractéristiques de référence historiques. L'exécution de ces tâches par lots (batch) à intervalles réguliers est bien maîtrisée, mais introduit un décalage de quelques minutes à plusieurs heures. Pour les signaux de référence concernant les utilisateurs, ce décalage est un prix acceptable pour une infrastructure plus simple. Lorsque les modèles nécessitent des signaux récents, cette infrastructure s'effondre ; descendre à la seconde ou à la milliseconde n'est pas possible avec les plateformes de feature store existantes. Pour exploiter la valeur des caractéristiques récentes, les data scientists sont contraints d'implémenter une logique complexe, spécifique au streaming, pour gérer ces agrégations et de mettre en place une infrastructure hébergée personnalisée.

Databricks Feature Store vous permet de créer une caractéristique une seule fois et de l'utiliser partout : la même définition alimente les flux de traitement par lots hors ligne à grande échelle et les pipelines de caractéristiques en ligne ultra-récentes. Le framework élimine la charge d'infrastructure, en orchestrant Spark Real-Time Mode (RTM) pour le traitement de flux continu, Lakebase pour le stockage en ligne optimisé pour le streaming, et Model Serving pour la récupération à grande échelle. Et une fois créée, cette caractéristique est fournie en quelques millisecondes : une latence p99 de bout en bout de 200 ms, depuis l'arrivée d'un événement dans Kafka jusqu'à sa disponibilité dans le feature store en ligne.

Architecture : de Kafka au Feature Store en 200 ms

image3.png

Jetons un coup d'œil sous le capot pour voir comment Databricks Feature Store prend une définition de caractéristique indépendante de l'infrastructure et construit un pipeline pour la calculer de manière cohérente en quelques millisecondes. Le parcours de bout en bout pour une caractéristique en streaming ressemble à ceci :

  1. Les événements arrivent dans Kafka - des données brutes comme des transactions par carte de crédit, des impressions publicitaires ou des événements de parcours de navigation (clickstream)
  2. Un pipeline Spark RTM sur des pipelines serverless Lakeflow Spark Delta traite en continu les événements, calculant des agrégations glissantes en temps réel
  3. Les agrégats mis à jour sont écrits dans Lakebase via un nouveau récepteur (sink) JDBC de streaming, arrivant ainsi dans le feature store en ligne
  4. Les points de terminaison Model Serving récupèrent les dernières caractéristiques depuis Lakebase au moment de l'inférence, les fournissant automatiquement au modèle

Associons cela à notre caractéristique de fraude, à savoir la somme du montant des transactions d'un utilisateur au cours des 10 dernières minutes. Chaque événement entrant contient les détails de la transaction (montant, lieu, identifiant de l'utilisateur, informations sur le commerçant) et est acheminé vers un pipeline avec état (stateful). Le pipeline consulte une instance locale RocksDB contenant le total cumulé des transactions de l'utilisateur, avec des délais d'expiration qui limitent la fenêtre aux 10 dernières minutes. Le pipeline lit et incrémente la valeur localement, puis écrit la valeur de caractéristique mise à jour dans Lakebase. Ainsi, lorsqu'une requête arrive au modèle pour approuver une nouvelle transaction, une somme de transactions à jour est disponible avec une fraîcheur inférieure à la seconde dans le feature store. Cette caractéristique de somme sera récupérée en même temps que la référence d'achat historique de l'utilisateur pour éclairer l'approbation. Une somme bien supérieure à la référence historique est un indicateur fort de fraude potentielle pour le modèle.

Chaque composant de ce pipeline a été optimisé afin que les événements entrants soient acheminés, les agrégations calculées et les caractéristiques écrites dans le magasin en ligne aussi rapidement que possible.

Fenêtre tournante (rolling window) : mettre à jour les agrégations en quelques millisecondes

image4.png

Avant d'entrer plus en détail dans l'infrastructure, parlons des caractéristiques d'agrégation et du passage d'un paradigme de synchronisation par lots (batch) à des mises à jour en temps réel.

Les caractéristiques d'agrégation sur une fenêtre de temps (par exemple, les nombres, les sommes ou les moyennes) sont des signaux puissants et flexibles pour le ML en temps réel. Une caractéristique par lots à long terme établit une référence historique pour l'utilisateur sur une période donnée, ce qui permet au modèle de s'adapter et de comprendre le comportement de chaque utilisateur. Une caractéristique courte et récente réagit rapidement aux situations changeantes pour distinguer un nouvel intérêt de l'utilisateur ou une activité frauduleuse. Les fenêtres temporelles définissent une plage de temps (par exemple, 10 minutes) ainsi que la manière dont ces plages doivent évoluer dans le temps (par exemple, chevauchement ou disjointes).

Databricks Feature Store prend en charge 3 fenêtres temporelles différentes :

  • Les fenêtres fixes (tumbling windows) sont alignées sur des intervalles d'horloge réelle et commencent dès que l'intervalle précédent se termine. Une fenêtre fixe de 10 minutes peut couvrir 12:00–12:10, puis 12:10–12:20. Les événements sont regroupés dans ces intervalles fixes, et une valeur de caractéristique est émise à la fin de l'intervalle. Cela signifie que l'agrégat n'est récent qu'aux limites de l'intervalle.
  • Les fenêtres glissantes (sliding windows) sont également alignées sur des intervalles d'horloge réelle, mais permettent un chevauchement des intervalles. Une fenêtre glissante de 10 minutes avec un intervalle de glissement de 5 minutes peut couvrir 12:00–12:10, puis 12:05–12:15, et enfin 12:10–12:20.
  • Les fenêtres tournantes (rolling windows) ne sont pas alignées sur l'horloge réelle, mais regardent en arrière à partir de l'horodatage de chaque événement avec une résolution à la milliseconde. « La somme des transactions au cours des 10 dernières minutes par rapport à l'heure actuelle » est toujours à jour, car la fenêtre se déplace avec chaque nouvel événement. Cela fait de RollingWindow le choix naturel pour le service en temps réel où le « présent » change constamment.

Les fenêtres fixes et glissantes restent utiles lorsqu'une caractéristique ne change pas fréquemment : elles émettent moins de mises à jour, sont moins coûteuses à maintenir et s'intègrent naturellement dans des pipelines planifiés plus simples. Les fenêtres tournantes sacrifient cette efficacité au profit d'une fraîcheur maximale, ce qui est particulièrement précieux pour les signaux où chaque nouvel événement doit immédiatement affecter la valeur fournie au modèle.

Voici à quel point il est simple de définir une caractéristique de fenêtre tournante avec l'API déclarative du Feature Store :

Spark Real-Time Mode : le moteur de calcul des caractéristiques

Pour en venir à l'infrastructure sous-jacente, le pipeline de streaming est ce qui permet d'obtenir des caractéristiques récentes avec un débit élevé. Ce pipeline achemine les données de Kafka jusqu'au feature store en ligne. Le pipeline de streaming est alimenté par Spark Real-Time Mode (RTM), un mode d'exécution fondamentalement nouveau pour Spark Structured Streaming. Le RTM est l'innovation architecturale clé qui rend possible une fraîcheur à la milliseconde.

Étapes simultanées et traitement avec état (stateful)

Dans le mode micro-lot traditionnel (MBM), Spark traite les données en streaming par lots distincts. Chaque lot collecte des événements sur un intervalle configurable, les traite de manière séquentielle à chaque étape, crée un point de contrôle (checkpoint), puis démarre le lot suivant. Cela crée un seuil minimal de latence : même avec un réglage agressif, les pipelines MBM pour les agrégations avec état fonctionnent généralement de l'ordre de la seconde à la minute. Le RTM, quant à lui, exécute les étapes de manière simultanée. Les opérateurs d'agrégation traitent activement les lignes dès qu'elles sont disponibles, sans attendre que l'étape en amont ait fini de traiter toutes les lignes.

Pour les agrégations tournantes, il y a deux étapes importantes. La première étape est le traitement des données, la validation du schéma, la fusion des données (coalescing), le transtypage (type casting). Cela exécute la logique métier qui convertit les événements d'action génériques au format requis pour l'agrégation de vos caractéristiques. La seconde étape consiste à agréger les données par entité pour calculer les agrégations de fenêtres tournantes. Chaque ligne entrante met immédiatement à jour l'agrégat dans un magasin d'état local RocksDB et émet la nouvelle valeur en aval. L'expiration de la fenêtre se produit également par ligne : lorsque la durée de la fenêtre s'écoule pour un événement donné, le pipeline supprime la contribution de cet événement et émet l'agrégat corrigé vers Lakebase. RocksDB s'exécute localement sur chaque exécuteur, ce qui permet d'avoir des tailles d'état qui dépassent la capacité mémoire du cluster.

Gestion de l'état du pipeline dans le RTM serverless

La création de points de contrôle (checkpointing) est essentielle pour la tolérance aux pannes dans le streaming avec état, car elle permet au pipeline de se rétablir en cas de défaillance d'un nœud de calcul (worker) individuel. Mais le checkpointing a un coût. En mode micro-lot, Spark crée des points de contrôle à chaque limite de lot, et chaque point de contrôle ajoute de la latence au pipeline car il interagit avec les magasins d'objets cloud.

RTM adopte une approche différente : le coût de la planification et des points de contrôle (checkpointing) est amorti sur des intervalles plus longs. Le coût des points de contrôle est réparti sur toutes les lignes traitées au cours de cet intervalle, plutôt que de bloquer le pipeline à chaque limite de lot (batch). Cela ne sacrifie en rien la tolérance aux pannes. Les garanties de traitement exactement une fois sont maintenues : en cas de défaillance, le pipeline rejoue au maximum 5 minutes de données provenant de la source Kafka. Le compromis est une légère augmentation du volume de rejeu pour une réduction significative de la latence de traitement en régime permanent.

Feature Store exécute des pipelines RTM serverless sur Lakeflow Spark Delta Pipelines (SDP), éliminant ainsi complètement la gestion des clusters et la planification de la capacité. Vous n'avez pas besoin de provisionner des machines, d'ajuster le nombre d'exécuteurs ou de vous soucier de la maintenance des clusters. Lorsque les mises à jour de l'infrastructure nécessitent un redémarrage du pipeline, SDP coordonne la transition : le nouveau cluster serverless est provisionné et entièrement prêt avant l'arrêt de l'ancien. Cette coordination est synchronisée aux intervalles de points de contrôle de 5 minutes, ce qui minimise les temps d'arrêt et évite les interruptions de retraitement. Cela se traduit par une interruption quasi nulle de la fraîcheur des fonctionnalités pendant les fenêtres de maintenance.

Lakebase : minimiser la surcharge pour les écritures en streaming

Databricks Feature Store utilise Lakebase pour stocker les valeurs de fonctionnalités en ligne pour l'inférence. L'architecture de Lakebase, qui sépare le calcul et le stockage, permet une mise à l'échelle automatique (autoscaling) pour gérer la charge variable de l'inférence de modèles. Le Feature Store en ligne exploite cette capacité pour monter en charge jusqu'à des dizaines de milliers de lectures par seconde avec une latence de l'ordre de la dizaine de millisecondes.

Les écritures en streaming sont particulièrement complexes car elles consistent en un grand nombre de petits upserts à mesure que de nouvelles valeurs de fenêtres glissantes sont émises pour chaque ligne Kafka reçue. Dans un système Postgres standard, ce modèle peut générer un volume important de journaux d'écriture anticipée (WAL), car Postgres écrit des pages entières pour faciliter la récupération. Après chaque point de contrôle, la première modification d'une page écrit l'image complète de la page de 8 Ko dans le journal WAL, et non pas seulement la petite modification logique. Pour les lignes d'entités très sollicitées ("hot") qui sont fréquemment mises à jour, l'amplification WAL devient le goulot d'étranglement pour le débit d'écriture, la réplication et la surcharge de récupération.

Lakebase exploite désormais la séparation du calcul et du stockage distribué pour minimiser l'amplification des écritures en streaming par rapport à un Postgres standard. L'architecture de Lakebase permet à Postgres d'écrire des enregistrements de modifications petits et compacts au lieu d'écrire de manière répétée des instantanés de pages entières de 8 Ko dans le WAL. La durabilité est toujours préservée car ces enregistrements compacts sont validés par un quorum de nœuds de protection (safekeepers) distribués. Des instantanés de pages entières restent nécessaires pour la récupération après un certain nombre d'enregistrements de modifications, mais ils sont générés plus tard dans la couche de stockage plutôt que d'encombrer le chemin d'écriture. Pour le Feature Store, le résultat est que RTM peut publier en continu de nouvelles valeurs de fonctionnalités dans Lakebase avec beaucoup moins d'amplification WAL et une latence supplémentaire minimale.

Model Serving : récupération de fonctionnalités à faible latence et à grande échelle

La dernière étape du processus consiste à récupérer les fonctionnalités fraîches depuis Lakebase et à les fournir au modèle au moment de l'inférence. Cela est géré par Databricks Model Serving, une infrastructure de service entièrement gérée et optimisée pour les charges de travail à fort volume de QPS et à faible latence.

Model Serving est conçu pour répondre aux exigences de débit du ML en temps réel :

  • Architecture entièrement évolutive horizontalement : le serveur d'inférence, la couche d'authentification, le proxy et le limiteur de débit évoluent tous de manière indépendante, supportant plus de 100 000 QPS sur les points de terminaison CPU
  • Mise à l'échelle élastique et rapide : le système s'adapte aux pics et aux baisses de trafic sans surprovisionnement, alignant ainsi les coûts sur la demande réelle
  • Gouverner et surveiller les modèles : gérez l'accès au réseau, gérez les autorisations pour les points de terminaison de modèles et surveillez la qualité à l'aide d'AI Gateway.

Pour le Feature Store, l'intégration est transparente. Lorsqu'un modèle est enregistré avec MLflow, ses dépendances de fonctionnalités sont consignées. Au moment de l'inférence, Model Serving recherche automatiquement les fonctionnalités requises dans Lakebase : pas de code de recherche personnalisé, pas d'intégration manuelle complexe. L'agrégat frais calculé par RTM et stocké dans Lakebase est récupéré et associé à la requête d'inférence de manière transparente.

Le Feature Store au-delà du streaming

Les capacités performantes en temps réel ne représentent qu'une partie de ce qu'un Feature Store peut résoudre. Deux autres défis méritent d'être brièvement examinés :

Données d'entraînement pour les fonctionnalités en streaming

La génération de données d'entraînement peut être difficile pour les fonctionnalités en streaming, car les courtes fenêtres de rétention sur les flux nécessitent de maintenir un stockage hors ligne (offline) distinct. Databricks Feature Store résout ce problème en stockant une copie hors ligne des données Kafka ingérées. Pour l'entraînement des modèles, le Feature Store calcule les mêmes valeurs de fonctionnalités que les pipelines de streaming pour les valeurs historiques et effectue des jointures temporellement exactes. Cette même capacité est utilisée pour rétrocharger les fonctionnalités de streaming en ligne afin de permettre un lancement rapide en production.

Intégrations

Comme illustré ci-dessus, les Feature Stores orchestrent plusieurs composants d'infrastructure complexes. Cette fragmentation peut rendre difficiles la gouvernance, le lignage et la réutilisation des fonctionnalités. Elle ralentit également le développement, car les ingénieurs doivent coordonner les modifications à travers différentes limites de systèmes.

Dans Databricks, les fonctionnalités sont des objets de premier ordre dans Unity Catalog : découvrables, régies par des contrôles d'accès et suivies avec un lignage complet. Les transformations de fonctionnalités sont packagées avec le modèle, MLflow capture les fonctionnalités utilisées et le lignage de déploiement connecte les modèles à leurs dépendances de fonctionnalités. La plateforme est un guichet unique pour développer, déployer et gouverner l'ensemble de votre pile ML.

Le Feature Store de Databricks orchestre des briques clés comme Spark RTM, Lakebase et Model Serving afin que vous bénéficiiez d'une latence et d'une échelle de premier ordre sans avoir à gérer l'infrastructure vous-même. Chacun de ces systèmes a été finement optimisé pour les charges de travail en streaming afin de faire d'une fraîcheur de 200 ms une réalité pour les fonctionnalités de machine learning.

Veuillez consulter la documentation sur les pipelines de streaming pour savoir comment définir des fonctionnalités de streaming. Expérimentez avec les fonctionnalités existantes pour voir à quel point elles fourniraient un signal plus fort avec une fraîcheur de l'ordre de la milliseconde.

Si vous souhaitez mieux comprendre la technologie sous-jacente, consultez le blog Lakebase sur les écritures plus rapides et l'analyse de l'architecture RTM.

Si ce sont le genre de défis sur lesquels vous souhaitez travailler, nous recrutons !

(Cet article de blog a été traduit à l'aide d'outils basés sur l'intelligence artificielle) Article original

Recevez les derniers articles dans votre boîte mail

Abonnez-vous à notre blog et recevez les derniers articles directement dans votre boîte mail.