Aller au contenu principal

Considérations de production pour le Structured Streaming

Exécutez les charges de travail Structured Streaming en production sous forme de Lakeflow Jobs planifiés sur Databricks. See Lakeflow Jobs.

Databricks vous recommande de toujours configurer les éléments suivants :

  • Supprimez le code inutile des Notebooks qui renverraient des résultats, tels que display et count.
  • Ne lancez pas les charges de travail Structured Streaming à l'aide d'un compute polyvalent. Planifiez toujours les Stream comme des Lakeflow Jobs en utilisant le compute des Job.
  • Planifier les Lakeflow Jobs en utilisant le mode Continuous. Cela fait référence à la fonctionnalité de planification de Databricks Jobs, et non à l'intervalle de Trigger de Structured Streaming.
  • N'activez pas l'autoscaling pour le compute pour les Jobs de Structured Streaming.

Certaines charges de travail bénéficient des éléments suivants :

Databricks a introduit les Lakeflow pipelines pour réduire les complexités de la gestion de l'infrastructure de production pour les charges de travail Structured Streaming. Databricks recommande d'utiliser les LakeFlow Pipelines pour les nouveaux pipelines Structured Streaming. See Spark Declarative Pipelines.

remarque

La mise à l’échelle automatique de compute présente des limites pour réduire la taille du cluster pour les charges de travail Structured Streaming. Databricks recommande d'utiliser Spark Declarative Pipelines sur Lakeflow avec une mise à l'échelle automatique améliorée pour les charges de travail de streaming. Voir Optimiser l’utilisation des clusters de LakeFlow Pipelines avec l’autoscaling.

:::note Compute serverless

Sur le compute serverless, seuls Trigger.AvailableNow() et Trigger.Once() sont pris en charge. Databricks recommande Trigger.AvailableNow().

Pour le streaming continu sur compute serverless, utilisez le mode de pipeline déclenché ou continu en mode continu.

Consultez les limitations du streaming.

:::

Concevez les charges de travail en streaming pour anticiper les pannes

Databricks vous recommande de toujours configurer les jobs de streaming pour qu'ils redémarrent automatiquement en cas d'échec. Certaines fonctionnalités, y compris l'évolution des schémas, exigent que les tâches Structured Streaming effectuent des nouvelles tentatives automatiquement. Consultez Configurer les Job Structured Streaming pour relancer les query de streaming en cas d'échec.

Certaines opérations comme foreachBatch offrent des garanties au moins une fois plutôt qu'exactement une fois. Pour ces Opérations, assurez-vous que votre pipeline de traitement est idempotent. Consultez Utiliser foreachBatch pour écrire dans des puits de données arbitraires.

remarque

Lorsqu'une query redémarre, le micro-batch prévu lors de l'exécution précédente est traité. Si votre Job a échoué en raison d’une erreur de mémoire insuffisante ou si vous avez annulé manuellement un Job en raison d’un micro-batch surdimensionné, vous devrez peut-être monter en charge le compute pour traiter le micro-batch avec succès.

Si vous modifiez les configurations entre les exécutions, ces configurations s’appliquent au premier nouveau batch planifié. Consultez l’article dédié à la récupération après des modifications dans une query Structured Streaming.

Lorsqu'un Job fait l'objet d'une nouvelle tentative

Vous pouvez planifier plusieurs tâches dans le cadre d'un Job Databricks. Lorsque vous configurez un Job à l'aide du Trigger continu, vous ne pouvez pas définir de dépendances entre les tâches.

Vous pouvez choisir de planifier plusieurs Stream dans un seul Job en utilisant l'une des approches suivantes :

  • Plusieurs tâches : définissez un Job avec plusieurs tâches qui exécutent des charges de travail en streaming à l'aide du Trigger continu.
  • Requêtes multiples : définissez plusieurs requêtes de streaming dans le code source pour une seule tâche.

Vous pouvez également combiner ces stratégies. Le tableau suivant compare ces approches.

Stratégie

Plusieurs tâches

Plusieurs requêtes

Comment le compute est partagé ?

Databricks recommande de déployer le compute pour les jobs avec une taille appropriée à chaque tâche de streaming. Vous pouvez éventuellement partager le compute entre les tâches.

Toutes les requêtes partagent le même compute. Vous pouvez éventuellement attribuer des query à des scheduler Pool.

Comment les tentatives sont-elles gérées ?

Toutes les tâches doivent échouer avant la relance du job.

La tâche est relancée si une query échoue.

Stratégie

Plusieurs tâches

Plusieurs requêtes

Comment le compute est partagé ?

Databricks recommande de déployer le compute pour les jobs avec une taille appropriée à chaque tâche de streaming. Vous pouvez éventuellement partager le compute entre les tâches.

Toutes les requêtes partagent le même compute. Vous pouvez éventuellement attribuer des query à des scheduler Pool.

Comment les tentatives sont-elles gérées ?

Toutes les tâches doivent échouer avant la relance du job.

La tâche est relancée si une query échoue.

Pour plus de détails sur l'utilisation de plusieurs tâches ou queries, consultez Exécuter plusieurs Structured Streaming queries sur le même cluster.

Configurez les Jobs Structured Streaming pour redémarrer les requêtes de streaming en cas d'échec

Databricks recommande de configurer toutes les charges de travail en streaming à l’aide du Trigger continu. Consultez Exécuter les Jobs en continu.

Le Trigger continu a le comportement suivant par default :

  • Empêche plus d'une exécution simultanée du job.
  • Démarre une nouvelle exécution lorsqu’une exécution précédente échoue.
  • Utilise un backoff exponentiel pour les nouvelles tentatives.

Databricks recommande d'utiliser toujours le compute de jobs plutôt que le compute polyvalent lors de la planification des workflows. En cas d'échec et de nouvelle tentative du job, de nouvelles ressources de compute sont déployées.

remarque

Databricks vous recommande de ne pas utiliser streamingQuery.awaitTermination() ou spark.streams.awaitAnyTermination(). Voir Quand utiliser awaitTermination().

Quand utiliser awaitTermination()

streamingQuery.awaitTermination() et spark.streams.awaitAnyTermination() bloquent le thread actuel jusqu'à ce qu'une query de streaming se termine. L'utilisation de ces fonctions dépend de votre environnement d'exécution.

Pour les Lakeflow Jobs, n'utilisez pas streamingQuery.awaitTermination() ou spark.streams.awaitAnyTermination(). Ces fonctions ne sont pas nécessaires, car le service Jobs empêche automatiquement l'exécution de se terminer lorsqu'une query de streaming est active. Les deux fonctions empêchent les cellules de Notebook de se terminer et empêchent le service Jobs de suivre la query streaming, ce qui perturbe les métriques de backlog et les notifications de Jobs.

Utilisez awaitTermination() dans les cas suivants :

Cas d'usage

Comportement

Notebooks interactifs sur compute multifonction

awaitTermination() garde la cellule en cours d'exécution, vous permet d'observer l'état de la query et garantit que les échecs apparaissent dans la sortie du Notebook.

Environnements locaux et de développement

Lors de l'exécution locale d'un programme Spark, le processus se termine lorsque le thread principal est terminé. Appelez awaitTermination() pour maintenir le programme en vie jusqu'à ce que la query streaming se termine ou échoue.

Propagation de l'échec au Driver

Sans awaitTermination(), un échec de la query de streaming dans un contexte autre que celui d’un Job peut ne pas se propager au thread appelant. La query peut échouer silencieusement, ce qui rend les échecs plus difficiles à détecter et à diagnostiquer. L’appel de awaitTermination() déclenche à nouveau l’exception de query sur le Driver.

Cas d'usage

Comportement

Notebooks interactifs sur compute multifonction

awaitTermination() garde la cellule en cours d'exécution, vous permet d'observer l'état de la query et garantit que les échecs apparaissent dans la sortie du Notebook.

Environnements locaux et de développement

Lors de l'exécution locale d'un programme Spark, le processus se termine lorsque le thread principal est terminé. Appelez awaitTermination() pour maintenir le programme en vie jusqu'à ce que la query streaming se termine ou échoue.

Propagation de l'échec au Driver

Sans awaitTermination(), un échec de la query de streaming dans un contexte autre que celui d’un Job peut ne pas se propager au thread appelant. La query peut échouer silencieusement, ce qui rend les échecs plus difficiles à détecter et à diagnostiquer. L’appel de awaitTermination() déclenche à nouveau l’exception de query sur le Driver.