Exécuter plusieurs queries Structured Streaming sur le même cluster
De nombreux clients exécutent plusieurs query Structured Streaming sur le même cluster Databricks. Bien que ce modèle soit pris en charge, Databricks recommande de limiter le nombre de query par cluster afin d’éviter les problèmes de mise à l’échelle et les goulots d’étranglement en termes de performances. Sur le compute serverless, Databricks gère la mise à l'échelle automatiquement, de sorte que ces considérations sont prises en charge pour vous. Si vous utilisez le compute classique, où vous contrôlez le dimensionnement du driver et des exécuteurs, cette page décrit les principaux goulots d'étranglement à prendre en compte et les moyens d'y remédier.
Databricks recommande d'utiliser les LakeFlow Pipelines pour les nouvelles charges de travail de streaming, qui gèrent automatiquement la complexité de l'infrastructure. See Spark Declarative Pipelines.
Quand utiliser plusieurs requêtes sur le même cluster
L'exécution de plusieurs requêtes de streaming sur le même cluster réduit les coûts d'infrastructure, en particulier lorsque vous avez de nombreux petits Stream qui ne nécessitent pas chacun un compute dédié. Le compromis clé est la défaillance partagée : si le cluster échoue, chaque Stream sur celui-ci échoue. Pour les pipelines stratégiques, ce mode de défaillance partagé est souvent inacceptable.
Pour les charges de travail qui mélangent des Stream critiques et non critiques, Databricks recommande ce qui suit :
- Attribuez à chaque Stream une priorité basée sur son impact commercial.
- Placez les flux critiques sur des clusters dédiés, même à un coût plus élevé.
- Co-localisez les Streams à priorité inférieure pour partager le compute et réduire les coûts.
Dimensionnement du Driver
Le driver est une ressource partagée. Plusieurs requêtes partagent le même CPU, la même mémoire, le même planificateur DAG, le même planificateur de tâches et l'exécution des UDF côté driver (par exemple, foreachBatch). Lorsque vous exécutez de nombreux streams concurrents, surveillez ces goulots d'étranglement spécifiques au-delà du provisionnement standard du CPU et de la mémoire :
- **Surcoût Auto Loader** : si vos Streams utilisent Auto Loader, la découverte de fichiers et la liste des répertoires augmentent la charge du Driver.
- **Limites de Ressources au niveau du système d'exploitation (fichiers ouverts)** : L'exécution d'un volume élevé de Stream basés sur des fichiers (tels que
FileStreamSourceou Auto Loader) simultanément sur un seul Driver peut épuiser les limites de descripteurs de fichiers ouverts au niveau de l'utilisateur, ce qui peut provoquer des échecs de Stream aléatoires. - Contre-pression du bus d'écouteur : Un nombre élevé de queries streaming concurrentes peut provoquer une contre-pression sur le
StreamingQueryListenerbus de la session Spark unique. Tous les événements (y comprisonQueryIdle) sont envoyés à ce bus unique, et un important arriéré d'événements peut retarder considérablement les gestionnaires asynchronesonQueryProgresset affecter la stabilité du cluster. - Opérations de driver coûteuses : Évitez d'appeler
collect()ou d'autres opérations coûteuses de DataFrame sur le driver, à moins que ce ne soit absolument nécessaire, pour éviter de matérialiser de grands ensembles de résultats et de provoquer des erreurs de mémoire insuffisante (OOM).
Dépanner le conflit de Driver.
Si vous rencontrez des plantages de Driver dus à des problèmes OOM ou de contention :
- Surveiller les métriques du Driver dans la Spark UI. Si vous constatez une utilisation élevée du CPU, de la mémoire ou du disque, ajustez le dimensionnement du driver dans les paramètres de compute du cluster.
- Si les problèmes persistent, vérifiez que votre code n'exécute pas d'opérations gourmandes en mémoire ou d'UDF sur le driver.
- Si vous ne pouvez pas monter en charge le Driver verticalement plus loin, Databricks vous recommande fortement de diviser vos Jobs sur plusieurs clusters afin de contourner ces goulots d'étranglement de mise à l'échelle des nœuds partagés.
Dimensionnement des exécuteurs
Avec plusieurs query exécutées sur le même cluster, toutes les query partagent des emplacements de tâche sur les exécuteurs. Les étapes d'une query peuvent occuper les emplacements disponibles, ce qui entraîne des retards et une famine des ressources pour les autres query. Spark utilise un mappage 1:1 entre les emplacements de tâches et les cœurs disponibles. Assurez-vous que suffisamment de cœurs sont disponibles si les query doivent s’exécuter simultanément.
En général, les exécuteurs peuvent effectuer des Opérations plus gourmandes en mémoire que le nœud du Driver. Ajustez les parameter d'allocation de mémoire hors tas et de la JVM de l'exécuteur si nécessaire pour gérer la charge de votre application. Assurez-vous que les nœuds d'exécuteur sont dimensionnés de manière appropriée en termes de CPU, de mémoire et d'espace disque et Montera en charge verticalement si nécessaire. Si la mise à l'échelle verticale n'est pas possible, vous pouvez envisager d'ajouter des nœuds Worker supplémentaires au cluster.
Certains de ces changements peuvent nécessiter le redémarrage du cluster pour prendre effet.
Utiliser les Pools du planificateur
Vous pouvez configurer des pools de planificateurs pour attribuer une capacité de compute aux requêtes lors de l'exécution de plusieurs requêtes de streaming à partir du même code source.
Par default, toutes les queries start dans un notebook s'exécutent dans le même Pool de planification équitable. Les jobs Apache Spark générés par des triggers à partir de toutes les queries en streaming dans un notebook s'exécutent les uns après les autres selon l'ordre « premier entré, premier sorti » (FIFO). Cela peut entraîner des retards inutiles dans les queries, car elles ne partagent pas efficacement les ressources du cluster.
Les Pools de planificateur vous permettent de déclarer quelles queries Structured Streaming partagent des Ressources de compute.
L'exemple suivant attribue query1 à un pool dédié, tandis que query2 et query3 partagent un pool de planificateur.
:::note Compatibilité Serverless
Databricks recommande de s’éloigner de spark.sparkContext, car il n’est pas compatible avec l’architecture de compute serverless de Databricks. Utilisez spark (SparkSession) directement à la place. Les pools de planificateur sont un concept de compute classique ; sur le Serverless, Databricks gère automatiquement la mise à l'échelle et l'allocation des ressources.
:::
# Run streaming query1 in scheduler pool1
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "pool1")
df.writeStream.queryName("query1").toTable("table1")
# Run streaming query2 in scheduler pool2
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "pool2")
df.writeStream.queryName("query2").toTable("table2")
# Run streaming query3 in scheduler pool2
spark.sparkContext.setLocalProperty("spark.scheduler.pool", "pool2")
df.writeStream.queryName("query3").toTable("table3")
La configuration des propriétés locales doit se trouver dans la même cellule de Notebook où vous start votre query streaming.
Pour plus d'informations sur les pools de planificateur équitable, consultez la documentation Apache Spark sur le planificateur équitable.
Considérations relatives aux requêtes avec état
Pour les queries avec état s'exécutant sur le même cluster, gardez les points suivants à l'esprit :
- Utilisez RocksDB comme fournisseur de magasin d’état pour éviter les problèmes d’OOM et les pauses du GC. RocksDB est le fournisseur de magasin d’état default dans Databricks Runtime 17,3 et versions ultérieures. Consultez Configurer le magasin d’état RocksDB sur Databricks.
- Optimiser les partitions de shuffle en fonction des exigences de votre application. Pour les étapes avec état, Spark planifie les tâches proportionnellement au nombre de partitions de shuffle.
- Limitez l'utilisation de la mémoire RocksDB par nœud pour éviter les erreurs OOM liées à l'utilisation de la mémoire hors tas. Ceci est géré automatiquement dans Databricks Runtime 17,3 et versions ultérieures, mais nécessite une configuration manuelle sur les versions antérieures. Consultez Limiter l'utilisation de la mémoire RocksDB.
- Évitez de surcharger trop de partitions sur le même nœud d'exécuteur. Les opérations de maintenance sur le magasin d'état, y compris l'upload et le nettoyage des instantanés, s'exécutent par nœud. L'affectation d'un trop grand nombre de partitions à un seul nœud d'exécuteur peut entraîner une insuffisance de maintenance et des temps de récupération plus longs en raison du nombre réduit d'instantanés complets disponibles.