Aller au contenu principal

Configurer les intervalles de trigger du Structured Streaming

Apache Spark Structured Streaming traite les données de manière incrémentielle. Les intervalles de Trigger contrôlent la fréquence à laquelle Structured Streaming vérifie les nouvelles données. Vous pouvez configurer des intervalles de Trigger pour le traitement quasi en temps réel, pour les refresh planifiés de bases de données, ou le traitement batch de toutes les nouvelles données pour une journée ou une semaine.

Parce que Qu'est-ce que l'Auto Loader ? utilise Structured Streaming pour charger des données, comprendre le fonctionnement des triggers vous offre la plus grande flexibilité pour contrôler les coûts tout en ingérant des données avec la fréquence souhaitée.

important

Databricks vous recommande de définir un mode de Trigger qui équilibre la latence et le coût pour votre cas d'utilisation. Dans le cas contraire, vous pourriez voir des coûts de stockage inattendus de la part de votre fournisseur cloud. Consultez Contrôler le coût du stockage cloud pour plus de détails.

Aperçu des modes de trigger

Le tableau suivant récapitule les modes de Trigger disponibles dans Structured Streaming :

Mode Trigger

Exemple de syntaxe (Python)

Idéal pour

Non spécifié (Default)

N/A

Streaming à usage général avec une latence de 3 à 5 secondes. Équivalent à un Trigger processingTime avec des intervalles de 0 ms. Le traitement en Stream s'exécute en continu tant que de nouvelles données arrivent.

Temps de traitement

.trigger(processingTime='10 seconds')

Équilibrer les coûts et les performances. Réduit la surcharge en empêchant le système de vérifier les données trop fréquemment.

Désormais disponible

.trigger(availableNow=True)

Traitement par batch incrémentiel planifié. Traite autant de données que disponible au moment où le Job de streaming est déclenché.

Mode temps réel

.trigger(realTime='5 minutes')

Charges de travail opérationnelles à très faible latence nécessitant un traitement en moins d'une seconde, telles que la détection de la fraude ou la personnalisation en temps réel. Aperçu public. « 5 minutes » indique la durée d'un micro-batch. Utilisez 5 minutes pour minimiser la surcharge par batch, telle que la compilation de query.

Continu

.trigger(continuous='1 second')

Non pris en charge. Il s'agit d'une fonctionnalité expérimentale incluse dans Spark OSS. Utilisez le mode temps réel à la place.

Mode Trigger

Exemple de syntaxe (Python)

Idéal pour

Non spécifié (Default)

N/A

Streaming à usage général avec une latence de 3 à 5 secondes. Équivalent à un Trigger processingTime avec des intervalles de 0 ms. Le traitement en Stream s'exécute en continu tant que de nouvelles données arrivent.

Temps de traitement

.trigger(processingTime='10 seconds')

Équilibrer les coûts et les performances. Réduit la surcharge en empêchant le système de vérifier les données trop fréquemment.

Désormais disponible

.trigger(availableNow=True)

Traitement par batch incrémentiel planifié. Traite autant de données que disponible au moment où le Job de streaming est déclenché.

Mode temps réel

.trigger(realTime='5 minutes')

Charges de travail opérationnelles à très faible latence nécessitant un traitement en moins d'une seconde, telles que la détection de la fraude ou la personnalisation en temps réel. Aperçu public. « 5 minutes » indique la durée d'un micro-batch. Utilisez 5 minutes pour minimiser la surcharge par batch, telle que la compilation de query.

Continu

.trigger(continuous='1 second')

Non pris en charge. Il s'agit d'une fonctionnalité expérimentale incluse dans Spark OSS. Utilisez le mode temps réel à la place.

:::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.

:::

processingTime : intervalles de Trigger basés sur le temps

Structured Streaming fait référence aux intervalles de trigger basés sur le temps en tant que « micro-batchs à intervalle fixe ». En utilisant le mot-clé processingTime, spécifiez une durée sous forme de chaîne, telle que .trigger(processingTime='10 seconds').

La configuration de cet intervalle détermine la fréquence à laquelle le système effectue des vérifications pour voir si de nouvelles données sont arrivées. Configurez votre temps de traitement pour équilibrer les exigences de latence et le débit d'arrivée des données dans la source.

AvailableNow: Traitement par batch incrémentiel

important

Dans Databricks Runtime 11.3 LTS et versions ultérieures, Trigger.Once est déprécié. Utilisez Trigger.AvailableNow pour toutes les charges de travail de traitement par batch incrémentiel.

L'option de Trigger AvailableNow consomme tous les enregistrements disponibles sous forme de batch incrémentiel avec la possibilité de configurer la taille du batch avec des options telles que maxBytesPerTrigger. Les options de dimensionnement varient selon la source de données.

Sources de données prises en charge

Databricks prend en charge l'utilisation de Trigger.AvailableNow pour le traitement par batch incrémental à partir de nombreuses sources Structured Streaming. Le tableau suivant inclut la version minimale prise en charge de Databricks Runtime requise pour chaque source de données :

Source

Version minimale de Databricks Runtime

Sources de fichiers (JSON, Parquet, etc.)

9.1 LTS

Delta Lake

10.4 LTS

Auto Loader

10.4 LTS

Apache Kafka

10.4 LTS

Kinesis

13,1

OpenSharing (responseFormat=delta; responseFormat=parquet nécessite delta-sharing-client 1.4.0 ou version ultérieure)

18,0

Source

Version minimale de Databricks Runtime

Sources de fichiers (JSON, Parquet, etc.)

9.1 LTS

Delta Lake

10.4 LTS

Auto Loader

10.4 LTS

Apache Kafka

10.4 LTS

Kinesis

13,1

OpenSharing (responseFormat=delta; responseFormat=parquet nécessite delta-sharing-client 1.4.0 ou version ultérieure)

18,0

realTime : Charges de travail opérationnelles à ultra-faible latence

Le mode temps réel pour Structured Streaming atteint une latence de bout en bout inférieure à 1 seconde à la fin, et dans les cas courants, d'environ 300 ms. Pour plus de détails sur la façon de configurer et d'utiliser efficacement le mode temps réel, consultez Mode temps réel dans Structured Streaming.

Apache Spark dispose d'un intervalle de trigger supplémentaire connu sous le nom de Continuous Processing. Ce mode est classé comme expérimental depuis Spark 2.3. Databricks ne prend pas en charge ni ne recommande ce mode. Utilisez plutôt le mode temps réel pour les cas d'utilisation à faible latence.

remarque

Le mode de traitement continu de cette page n’est pas lié au traitement continu dans les Spark Declarative Pipelines.

Contrôlez le coût du stockage cloud

Par défaut, si vous ne définissez pas de mode de Trigger, Structured Streaming définit le mode de Trigger sur processingTime et l'intervalle sur 0, ce qui vérifie les nouvelles données toutes les quelques millisecondes. Cela peut générer un volume élevé d'appels d'API de stockage cloud par jour et entraîner des frais inattendus de la part de votre fournisseur de cloud.

Databricks vous recommande de configurer un mode de trigger approprié à vos exigences en matière de latence et de coût. Voir processingTime pour plus d’informations sur la configuration d’un intervalle de trigger temporel.

Modifier les intervalles de Trigger entre les exécutions

Vous pouvez modifier l'intervalle de Trigger entre les exécutions tout en utilisant le même point de contrôle.

Comportement lors de la modification des intervalles

Si une query Structured Streaming s'arrête pendant qu'un micro-batch est en cours de traitement, ce micro-batch doit être terminé avant que le nouvel intervalle de Trigger ne s'applique. Après avoir modifié l'intervalle de Trigger, vous pourriez observer qu'un micro-batch est traité avec la configuration précédemment spécifiée. Ce qui suit décrit le comportement attendu après une transition :

  • D'intervalle temporel à AvailableNow : un micro-batch peut être traité comme un batch incrémentiel avant que tous les enregistrements disponibles ne soient traités.
  • De AvailableNow à intervalle basé sur le temps : Le traitement pourrait continuer pour tous les enregistrements qui étaient disponibles lorsque le dernier job AvailableNow Trigger.

Récupérer après des échecs de query

Si vous tentez de récupérer après un échec de query avec un batch incrémentiel, une modification de l'intervalle de Trigger ne résout pas le problème. Le batch infructueux précédent doit être terminé car Structured Streaming nécessite des micro-batchs idempotents. Voir la sémantique de tolérance aux pannes pour Apache Spark.

Pour résoudre l'échec, montez en charge la capacité de compute, telle que l'augmentation de la taille des nœuds worker. Dans de rares cas, vous devrez peut-être redémarrer le Stream avec un nouveau point de contrôle.