Sélectionnez un mode de sortie pour Structured Streaming
Cet article aborde la sélection d'un mode de sortie pour le streaming avec état. Seuls les flux avec état contenant des agrégations nécessitent une configuration du mode de sortie.
Les jointures ne prennent en charge que le mode de sortie d'ajout, et le mode de sortie n'a pas d'impact sur la déduplication. Les opérateurs avec état arbitraires mapGroupsWithState et flatMapGroupsWithState émettent des enregistrements en utilisant leur propre logique personnalisée, de sorte que le mode de sortie du Stream n'affecte pas leur comportement.
Pour le streaming sans état, tous les modes de sortie se comportent de la même manière.
Pour configurer correctement le mode de sortie, vous devez comprendre le streaming avec état, les watermarks et les déclencheurs. Consultez les articles suivants :
- Qu'est-ce que le streaming avec état ?
- Appliquer des filigranes pour contrôler les thresholds de traitement des données
- Configurer les intervalles de Trigger de Structured Streaming
Qu'est-ce que le mode de sortie ?
Le mode de sortie d'une requête Structured Streaming détermine les enregistrements que les opérateurs de la requête émettent lors de chaque Trigger. Les trois types d'enregistrements qui peuvent être émis sont :
- Enregistrements que le traitement futur ne modifie pas.
- Les enregistrements qui ont changé depuis le dernier Trigger.
- Tous les enregistrements de la table d'état.
Savoir quels types d'enregistrements émettre est important pour les opérateurs avec état, car une ligne particulière produite par un opérateur avec état peut changer d'un Trigger à l'autre. Par exemple, lorsqu'un opérateur d'agrégation en streaming reçoit plus de lignes pour une fenêtre particulière, les valeurs d'agrégation de cette fenêtre peuvent changer d'un Trigger à l'autre.
Pour les opérateurs sans état, la distinction entre les types d'enregistrements n'affecte pas le comportement de l'opérateur. Les enregistrements qu'un opérateur sans état émet pendant un Trigger sont toujours les enregistrements source traités pendant ce Trigger.
Modes de sortie disponibles
Il existe trois modes de sortie qui indiquent à un opérateur quels enregistrements émettre pendant un Trigger particulier :
Mode de sortie | Description |
|---|---|
Mode d'ajout (default) | Par default, les requêtes en streaming s'exécutent en mode ajout. Dans ce mode, les opérateurs n'émettent que les lignes qui ne changent pas dans les futurs Trigger. Les opérateurs avec état utilisent le filigrane pour déterminer quand cela se produit. |
Mode de mise à jour | En mode de mise à jour, les opérateurs émettent toutes les lignes qui ont changé pendant le Trigger, même si l'enregistrement émis peut changer lors d'un Trigger ultérieur. |
Mode complet | Le mode complet ne fonctionne qu'avec les agrégations de streaming. En mode complet, toutes les lignes résultantes jamais produites par l'opérateur sont émises en aval. |
Considérations relatives à la production
Pour de nombreuses Opérations de streaming avec état, vous devez choisir entre les modes d'ajout et de mise à jour. Les sections suivantes décrivent les considérations qui pourraient éclairer votre décision.
Le mode complet a quelques applications, mais peut être peu performant à mesure que les données montent en charge. Databricks recommande d'utiliser des vues matérialisées pour obtenir des garanties sémantiques associées au mode complet avec un traitement incrémentiel pour de nombreuses opérations avec état. Consultez Vues matérialisées.
Sémantique d'application
La sémantique des applications décrit la manière dont les applications en aval utilisent les données de streaming.
Si les services en aval doivent effectuer une seule action pour chaque écriture en aval, utilisez le mode d'ajout dans la plupart des cas. Par exemple, si vous disposez d'un service de notification en aval qui envoie des notifications pour chaque nouvel enregistrement écrit dans le récepteur, le mode d'ajout garantit que chaque enregistrement n'est écrit qu'une seule fois. Le mode de mise à jour écrit l'enregistrement chaque fois que l'information d'état change, ce qui entraînerait de nombreuses mises à jour.
Si les services en aval nécessitent des résultats actualisés, le mode de mise à jour garantit que votre cible reste aussi à jour que possible. Les exemples incluent un modèle de machine learning qui lit les fonctionnalités en temps réel ou un tableau de bord analytique qui suit les agrégats en temps réel.
Compatibilité des opérateurs et des sinks
Structured Streaming ne prend pas en charge toutes les opérations disponibles dans Apache Spark, et certaines opérations de streaming ne sont pas prises en charge dans tous les modes de sortie. Pour en savoir plus sur les limitations d'opérateur, consultez la documentation sur le streaming OSS.
Tous les sinks ne prennent pas en charge tous les modes de sortie. Kafka prend en charge tous les modes de sortie. Delta Lake, qui prend en charge toutes les tables gérées par Unity Catalog, prend en charge les modes d'ajout et d'achèvement, mais pas le mode de mise à jour. Pour un comportement similaire au mode de mise à jour avec les récepteurs Delta Lake, consultez Merge en streaming.
Pour en savoir plus sur la compatibilité des récepteurs, consultez la documentation de streaming OSS.
Latence et coût
Le mode de sortie influe sur le délai avant l'écriture d'un enregistrement, et la fréquence ainsi que la quantité de données écrites peuvent avoir une incidence sur les coûts associés aux pipelines de streaming.
Le mode d'ajout oblige les opérateurs avec état à émettre des résultats uniquement après que les résultats avec état sont finalisés, ce qui est au moins aussi long que votre délai de filigrane. Un délai de filigrane de 1 hour en mode de sortie d'ajout signifie que vos enregistrements ont au moins un délai d'une heure avant d'être émis en aval.
Le mode de mise à jour entraîne une écriture par trigger par valeur agrégée. Si votre cible facture par écriture par enregistrement, cela peut être coûteux si les enregistrements sont mis à jour de nombreuses fois avant que le délai du filigrane ne s'écoule.
Exemples de configuration
Les exemples de code suivants montrent comment configurer le mode de sortie pour les mises à jour en streaming des tables Unity Catalog :
- Python
- Scala
# Append output mode (default)
(df.writeStream
.toTable("target_table")
)
# Append output mode (same as default behavior)
(df.writeStream
.outputMode("append")
.toTable("target_table")
)
# Update output mode
(df.writeStream
.outputMode("update")
.toTable("target_table")
)
# Complete output mode
(df.writeStream
.outputMode("complete")
.toTable("target_table")
)
// Append output mode (default)
df.writeStream
.toTable("target_table")
// Append output mode (same as default behavior)
df.writeStream
.outputMode("append")
.toTable("target_table")
// Update output mode
df.writeStream
.outputMode("update")
.toTable("target_table")
// Complete output mode
df.writeStream
.outputMode("complete")
.toTable("target_table")
Consultez la documentation OSS pour PySpark DataStreamWriter.outputMode ou Scala DataStreamWriter.outputMode.
Exemple de streaming avec état et de modes de sortie
L'exemple suivant vise à vous aider à comprendre comment le mode de sortie interagit avec les filigranes pour le streaming avec état.
Considérez une agrégation de streaming qui calcule le chiffre d'affaires total généré chaque heure dans un magasin avec un délai de filigrane de 15 minutes. Le premier micro-lot traite les enregistrements suivants :
- 15 $ à 14 h 40
- 10 $ à 14 h 30
- 30 $ à 15 h 10
À ce stade, le filigrane du moteur est 14 h 55, car il soustrait 15 minutes (le délai) de l'heure maximale observée (15 h 10). L'opérateur d'agrégation de streaming contient les éléments suivants dans son état :
[2pm, 3pm]: 25 $[3pm, 4pm]: 30 $
Le tableau suivant décrit ce qui se passerait dans chaque mode de sortie :
Mode de résultat | Résultat et motif |
|---|---|
Ajouter | L'opérateur d'agrégation de streaming n'émet rien en aval. Ceci est dû au fait que ces deux fenêtres pourraient changer à mesure que de nouvelles valeurs apparaissent avec un Trigger ultérieur : le filigrane de 14 h 55 indique que les enregistrements après 14 h 55 pourraient toujours arriver, et ces enregistrements pourraient se trouver dans la fenêtre |
Mettre à jour | L'opérateur émet les deux enregistrements, car les deux enregistrements ont reçu des mises à jour. |
Terminé | L'opérateur émet tous les enregistrements. |
Maintenant, supposons que le stream reçoive un enregistrement supplémentaire :
- 20 $ à 15 h 20
Le filigrane se met à jour à 15 h 05 parce que le moteur soustrait 15 minutes de 15 h 20. À ce stade, l'opérateur d'agrégation de streaming contient les éléments suivants dans son état :
[2pm, 3pm]: 25 $[3pm, 4pm]: 50 $
Le tableau suivant décrit ce qui se passerait dans chaque mode de sortie :
Mode de résultat | Résultat et motif |
|---|---|
Ajouter | L'opérateur d'agrégation de streaming observe que le filigrane de 15 h 05 est supérieur à la fin de la fenêtre |
Mettre à jour | L'opérateur d'agrégation de streaming émet la fenêtre |
Terminé | L'opérateur émet tous les enregistrements. |
Voici un résumé du comportement des opérateurs avec état dans chaque mode d’ajout :
- En mode ajout, écrivez les enregistrements une seule fois après le délai du filigrane.
- En mode mise à jour, écrivez les enregistrements qui ont changé depuis le Trigger précédent.
- En mode complet, écrivez tous les enregistrements jamais produits par l'opérateur avec état.