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.

:::

Réduire la latence pour le streaming opérationnel

Les charges de travail de streaming opérationnel ingèrent, transforment et agissent sur les données en temps quasi réel. Les exemples courants incluent la détection de la fraude, la détection d'anomalies, la personnalisation, ainsi que le monitoring et les alertes en temps réel, où un traitement différé affecte directement les résultats commerciaux. Une faible latence pour ces charges de travail signifie généralement des dizaines à des centaines de millisecondes, bien que de nombreuses équipes définissent des accords de niveau de service (SLA) de l'ordre de la seconde pour tenir compte de la variabilité aux percentiles les plus élevés.

Pour obtenir la latence de bout en bout la plus faible, utilisez le mode temps réel, qui permet d’atteindre une latence de bout en bout inférieure à une seconde dans les cas extrêmes et d’environ 300 millisecondes dans les cas courants. Voir Concepts du mode temps réel.

Lorsque le mode temps réel ne convient pas à votre workload, les bonnes pratiques suivantes permettent de réduire la latence pour le Structured Streaming en micro-batch :

  • Mode de sortie : utilisez le mode de mise à jour lorsque vos opérateurs de query et votre récepteur le prennent en charge. Le mode de mise à jour émet des lignes mises à jour après chaque trigger et continue de les mettre à jour jusqu'à l'expiration du watermark ; assurez-vous donc que votre récepteur en aval est idempotent pour gérer les résultats mis à jour. Utilisez le mode append pour les workloads que le mode update ne prend pas en charge, comme les jointures stream-stream, ou lorsque vous pouvez supprimer les données arrivant en retard. N'utilisez pas le mode complet pour une faible latence. Consultez Sélectionner un mode de sortie pour Structured Streaming.
  • Trigger : utilisez un processingTime Trigger avec un intervalle 0, qui start le micro-batch suivant dès que le précédent est terminé et que de nouvelles données sont disponibles. Cela permet d’obtenir la latence de micro-batch la plus faible, mais augmente les coûts de l’API de stockage cloud. N’utilisez pas AvailableNow, Once ou Continuous pour les workloads opérationnels. Consultez Configurer les intervalles de trigger Structured Streaming.
  • Watermark : définissez le watermark suffisamment long pour inclure les données arrivées tardivement que votre charge de travail ne doit pas ignorer. Le watermark contrôle la durée pendant laquelle la query accepte des données d’événement hors séquence avant de les supprimer et d’évincer l’état ; ainsi, un watermark trop court rejette silencieusement des enregistrements tardifs valides. Dans cette limite, un watermark plus court réduit la latence et conserve moins d'état, tandis qu'un watermark plus long tolère davantage de données arrivées tardivement au prix d'une latence et d'un état accrus. Un petit multiple de votre SLA de latence, tel que 2x, constitue un point de départ raisonnable pour le réglage. Consultez Appliquer des watermarks pour contrôler les threshold de traitement des données.
  • Sources et récepteurs : lisez à partir de sources à faible latence telles que des bus de messages (Apache Kafka, Amazon Kinesis, Apache Pulsar ou Google Cloud Pub/Sub) ou des flux de données de changement provenant de tables Delta Lake et Apache Iceberg. Écrivez vers des récepteurs à faible latence et à haut throughput tels que des bus de messages, des bases de données opérationnelles ou des récepteurs foreach. Concevez les Opérations de destination de manière à ce qu'elles soient idempotentes, afin que les consommateurs en aval puissent gérer les doublons et les données arrivant en retard.
  • État et création de points de contrôle : pour les queries avec état, utilisez le magasin d’état RocksDB, qui est requis à la fois pour la création de points de contrôle de journal des modifications et pour la création asynchrone de points de contrôle d’état. Activez la création de points de contrôle de journal des modifications pour ne conserver que les changements d’état incrémentiels. Lorsque la création de points de contrôle d’état constitue le goulot d’étranglement de la durée de votre batch, activez la création asynchrone de points de contrôle d’état pour superposer les écritures de points de contrôle au micro-batch suivant, après avoir examiné ses mises en garde concernant la récupération après défaillance et le redimensionnement des clusters. Donnez à chaque query son propre répertoire de points de contrôle dans un stockage cloud durable. Voir Configurer le magasin d’état RocksDB sur Databricks, Création asynchrone de points de contrôle d’état pour les query avec état et Points de contrôle Structured Streaming.
  • Gestion des offsets : pour réduire la latence due à la création de points de contrôle d’offset dans les streams continus, activez le suivi de progression asynchrone, qui met à jour les logs d’offset et de commit sans bloquer le traitement des données. Il n’est pas compatible avec les triggers AvailableNow ou Once. Voir Suivi de progression asynchrone.
  • Sauts de stockage : maintenez le calcul au sein d’un seul pipeline de streaming dans la mesure du possible. Répartir la logique sur plusieurs jobs ou pipelines ajoute des sauts de stockage qui augmentent la latence.

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.