Adaptez la taille de vos flux avec état les plus exigeants, de la détection des fraudes à la surveillance en temps réel, sans jamais reconstruire l'état de votre checkpoint.
par Thangam Vaiyapuri, Jay Palaniappan, B. Micheal Okutubo et Zifei Feng
Quiconque exécute des requêtes de streaming avec état Apache Spark™ Structured Streaming en production finit par se heurter au même mur.
Vous avez lancé la requête il y a plusieurs mois. À l'époque, le volume de données était modeste, vous avez donc accepté la valeur par défaut de 200 partitions de shuffle et vous êtes passé à autre chose. Le pipeline fonctionnait sans problème. Puis l'activité a grandi, le trafic a triplé, et le magasin d'états a explosé. Soudain, ces 200 partitions ne sont plus à la bonne taille. Certaines partitions sont déséquilibrées et surchargées, le cluster est mis à rude épreuve, et chaque micro-batch prend plus de temps qu'il ne le devrait.
Vous faites donc ce qui semble logique : vous augmentez spark.sql.shuffle.partitions et redémarrez la requête. Rien ne change.
La requête ignore silencieusement votre nouvelle valeur car le nombre de partitions a été figé dans le checkpoint lors du premier démarrage du flux. Historiquement, la seule façon d'appliquer un nouveau nombre consistait à abandonner le checkpoint existant et à recommencer à zéro, ce qui, pour une requête avec état, signifie perdre tout l'état accumulé que vous aviez soigneusement conservé. Pour un modèle de fraude qui suit des millions de comptes, ou une tâche de sessionisation conservant des jours de fenêtres, « recommencer à zéro » n'est pas une phrase que l'on souhaite prononcer lors d'une revue d'incident en production.
Le repartitionnement d'état à la demande (Public Preview), disponible dans Databricks Runtime 18 et versions ultérieures, élimine cet obstacle. Vous pouvez désormais redimensionner le nombre de partitions pour une requête de streaming avec état tout en conservant l'état de votre checkpoint intact.
Cela s'applique à toute requête de streaming avec état, que vous exécutiez des agrégations, des jointures flux-flux, de la déduplication, de la sessionisation ou transformWithState, et à toute charge de travail, de la détection des fraudes à la surveillance en temps réel.
Pour les premiers utilisateurs comme Coveo, la possibilité d'adapter à la demande la taille de leur infrastructure de streaming s'est immédiatement traduite par d'importantes économies opérationnelles.
« Chez Coveo, nous exécutons des pipelines de streaming avec état à grande échelle où les volumes de données fluctuent considérablement au fil du temps. Grâce à Databricks et à la fonctionnalité de repartitionnement d'état (State Repartitioning), nous avons réduit de 40 % nos coûts d'API Amazon S3 associés. Auparavant, chaque décision de mise à l'échelle imposait un compromis : soit surprovisionner, soit reconstruire à partir de nouveaux checkpoints, ce qui entraînait des coûts d'API de stockage presque identiques aux coûts de calcul. Désormais, nous adaptons l'échelle librement au gré des variations de la demande, sans perturber l'état existant ni déclencher de coûteuses migrations de checkpoints. » —Alexis Chicoine, développeur logiciel senior, Coveo
Pour comprendre pourquoi les résultats de Coveo représentent un bond en avant significatif pour Structured Streaming, nous devons examiner pourquoi le nombre de partitions était figé à l'origine.
Une requête de streaming avec état conserve son état dans un magasin d'états, et cet état est physiquement partitionné. Chaque clé de votre flux (un ID utilisateur, un numéro de compte, une fenêtre) est hachée vers une partition spécifique, et les données de chaque partition sont stockées dans leur propre instance RocksDB distincte au sein du checkpoint. Le nombre de partitions définit la structure de l'ensemble du magasin d'états sur le disque.
Si vous modifiez simplement le nombre de partitions entre deux redémarrages, le hachage ne correspondrait plus. Une clé qui résidait auparavant dans une partition (par exemple, la partition 47) pourrait désormais être hachée vers une autre (la partition 12), mais son état accumulé se trouverait toujours dans les fichiers de la partition d'origine. La requête perdrait en fait la trace de sa propre mémoire. Pour éviter précisément ce type de corruption silencieuse, Structured Streaming verrouillait le nombre de partitions lors de la création du checkpoint et ignorait toute modification ultérieure de spark.sql.shuffle.partitions.
Sûr, mais rigide. Les deux inconvénients étaient :
Le repartitionnement d'état à la demande résout ces deux problèmes en faisant la seule chose que l'ancienne architecture refusait de faire, mais de manière sécurisée, en redistribuant physiquement l'état pour correspondre au nouveau nombre de partitions.
Les prérequis sont simples :
C'est tout pour les prérequis. Si vous utilisez DBR 18 avec le magasin d'états par défaut, vous disposez déjà de tout ce dont vous avez besoin.
Le mécanisme est simple et réutilise un modèle que tout développeur de streaming connaît déjà : arrêter, reconfigurer, redémarrer.
Au lieu de spark.sql.shuffle.partitions, vous définissez une configuration dédiée, spark.sql.streaming.stateStore.partitions, et redémarrez la requête :
Le détail clé réside dans la nouvelle configuration elle-même. Pour les requêtes avec état, spark.sql.streaming.stateStore.partitions prévaudra sur spark.sql.shuffle.partitions. C'est ce qui permet de pérenniser le changement, là où l'ancienne approche échouait.
Lorsque la requête redémarre, elle ne reprend pas immédiatement son traitement normal. Tout d'abord, elle termine le dernier micro-batch planifié, s'il y en a un en attente. Ensuite, elle effectue une opération de repartitionnement unique : elle redistribue physiquement les données d'état sur le nouveau nombre de partitions, en re-hachant les clés vers leurs nouveaux emplacements corrects afin que rien ne soit perdu ou égaré. Une fois cette redistribution terminée, la requête reprend son traitement habituel, en utilisant désormais le nombre de partitions que vous avez demandé.
Cette étape de repartitionnement est le cœur de la fonctionnalité. C'est toute la différence entre « nous avons changé un nombre » et « nous avons déplacé votre état en toute sécurité vers une nouvelle structure ».
Le repartitionnement étant une opération réelle dont le temps d'exécution est proportionnel au volume de l'état, vous aurez besoin de visibilité. Structured Streaming expose cette information via ses rapports de progression standard.
Une fois le micro-batch suivant terminé, les événements StreamingQueryProgress incluent la durée de l'opération de repartitionnement. Recherchez dans les métriques durationMs de l'événement le champ controlBatch.REPARTITION, qui indique la durée du repartitionnement en millisecondes.
Une empreinte d'état plus importante signifie un repartitionnement plus long, mais nous prévoyons que cela ne prendra que quelques secondes pour la plupart des charges de travail. Ainsi, sur les tâches volumineuses, il est utile de capturer cette métrique pour en comprendre la durée. Pour en savoir plus sur la lecture de ces événements, consultez Surveiller les requêtes Structured Streaming sur Databricks.
Rendons cela concret avec une agrégation simple, à savoir un décompte d'événements par ID sur une fenêtre basculante (tumbling window). Nous allons la démarrer avec la valeur par défaut de 200 partitions, décider que c'est plus que nécessaire pour cette charge de travail, et réduire l'échelle à 100.
Tout d'abord, la requête telle qu'elle s'exécute aujourd'hui, avec le nombre de partitions par défaut :
À présent, après avoir observé ce flux pendant un certain temps, nous avons conclu que 200 partitions étaient excessives. Nous payons une surcharge de coordination pour un parallélisme dont nous n'avons pas besoin. Nous arrêtons la requête, définissons le nouveau nombre de partitions et la redémarrons avec les mêmes options et le même checkpoint :
Lorsque la requête redémarrée se lance, elle finalise le dernier micro-batch planifié, s'il y en a un en attente, exécute le repartitionnement pour redistribuer l'état de 200 à 100 partitions, puis continue le décompte en préservant intégralement chaque fenêtre et chaque total cumulé. La même procédure fonctionne en sens inverse : pour augmenter l'échelle sous une charge plus lourde, il vous suffit de définir un nombre plus élevé.
La même approche s'applique aux Spark Declarative Pipelines (SDP). Consultez l'exemple de SDP dans la documentation pour un guide complet.
Le repartitionnement d'état à la demande est un outil d'ajustement et de mise à l'échelle plutôt qu'une opération de routine. Il s'avère précieux dans quelques situations clés :
Chaque modification nécessitant un arrêt et un redémarrage avec une pause de repartitionnement unique, considérez cela comme une action de maintenance délibérée. Planifiez le redimensionnement lors d'une fenêtre d'intervention où une brève pause de traitement est acceptable, surveillez controlBatch.REPARTITION pour confirmer le temps que cela a pris, et laissez la requête reprendre son rythme normal.
Pendant des années, le nombre de partitions d'une requête de streaming avec état était une décision prise une fois pour toutes, au tout début, sur laquelle on ne revenait jamais sous peine de devoir reconstruire l'état à partir de zéro à un coût élevé. Le repartitionnement d'état à la demande élimine ces contraintes. La redistribution sécurisée de l'état sur un nouveau nombre de partitions transforme une décision initialement définitive en un choix que vous pouvez réévaluer chaque fois que votre charge de travail l'exige.
Le résultat correspond exactement à ce que les exploitants de flux de longue durée attendaient : la liberté de dimensionner correctement une requête en fonction de ses besoins d'évolution, par un simple arrêt, une modification de configuration et un redémarrage, sans perdre son état.
Le repartitionnement d'état à la demande est disponible dans Databricks Runtime 18 et versions ultérieures, en utilisant le fournisseur de stockage d'état RocksDB. Pour obtenir la référence complète, consultez Le repartitionnement d'état à la demande pour les requêtes de streaming avec état.
(Cet article de blog a été traduit à l'aide d'outils basés sur l'intelligence artificielle) Article original
Abonnez-vous à notre blog et recevez les derniers articles directement dans votre boîte mail.