Aller au contenu principal

Appliquer des filigranes pour contrôler les threshold de traitement des données

Cette page décrit les concepts de filigrane et propose des recommandations pour l'utilisation des filigranes dans les opérations de streaming avec état courantes.

Les queries de streaming accumulent les données d'état au fil du temps. Les filigranes suppriment automatiquement les anciennes données d'état pour éviter les erreurs de mémoire et une latence de traitement accrue.

Qu'est-ce qu'un filigrane ?

Pendant le traitement, Structured Streaming conserve l'état entre les micro-batchs. Les queries en streaming utilisent un état pour mettre à jour les résultats de manière incrémentielle au lieu de tout recalculer après chaque micro-batch. Les watermarks contrôlent le threshold à partir duquel une query cesse de traiter une entité d'état.

Les exemples courants d'entités d'état incluent :

  • Agrégations sur une fenêtre temporelle.
  • Clés uniques dans une jointure entre deux Streams.

Pour déclarer un filigrane sur un DataFrame en streaming, spécifiez un champ de Timestamp et un threshold de latence. À mesure que de nouvelles données arrivent, le gestionnaire d'état suit le timestamp le plus récent dans le champ spécifié et traite uniquement les enregistrements dans le threshold de latence.

Les queries traitent toujours les enregistrements qui arrivent dans le threshold. Les queries peuvent encore traiter les enregistrements qui arrivent en dehors du threshold, mais cela n'est pas garanti.

L'exemple suivant applique un watermark threshold de 10 minutes à un nombre fenêtré :

Python
from pyspark.sql.functions import window

(df
.withWatermark("event_time", "10 minutes")
.groupBy(
window("event_time", "5 minutes"),
"id")
.count()
)

Dans cet exemple :

  • La colonne event_time est utilisée pour définir un filigrane de 10 minutes et une fenêtre glissante de 5 minutes.
  • Un décompte est collecté pour chaque id observé pour chaque fenêtre de 5 minutes non chevauchante.
  • L'information d'état est maintenue pour chaque décompte jusqu'à ce que la fin de la fenêtre soit 10 minutes plus ancienne que le event_time le plus récent observé.
important

Dans une Opération groupBy() et window(), référencez les colonnes par nom, "<colName>" ou col("<colName>"), pour vous assurer que le marqueur de temps d’événement est conservé. En Scala, vous pouvez également utiliser $colName.

Comment les watermarks affectent-ils le temps de traitement et le throughput ?

Les modes de sortie contrôlent quand une query avec des filigranes écrit des données dans le récepteur. Les filigranes sont essentiels pour le contrôle du throughput dans le streaming avec état, car ils réduisent la quantité totale d’informations d’état en mémoire. Tous les modes de sortie ne sont pas pris en charge pour toutes les Opérations avec état. Voir Filigranes et mode de sortie pour les agrégations fenêtrées.

La sélection d'une durée de filigrane présente des compromis :

  • Des filigranes plus courts réduisent la latence des queries, car les queries stockent moins d'informations d'état et écrivent les résultats une fois que chaque durée de query est terminée. Cependant, les seuils courts ont une faible tolérance pour les données tardives.
  • Des filigranes plus longs ont une tolérance élevée pour les données tardives. Cependant, les longs filigranes augmentent la latence des query, car les query doivent stocker davantage d’information d’état et attendre d’écrire les résultats après une durée de filigrane plus longue.

Filigranes et mode de sortie pour les agrégations par fenêtre

Le tableau suivant présente le comportement de traitement pour les query avec agrégation sur un Timestamp et un watermark :

Mode de résultat

Comportement

Ajouter

La query écrit des lignes dans la table cible une fois le watermark threshold dépassé. Toutes les écritures sont retardées en fonction du threshold de latence. L'ancien état d'agrégation est supprimé après le dépassement du threshold.

Mettre à jour

La query écrit des lignes dans la table cible à mesure que les résultats sont calculés, et la query peut mettre à jour et écraser des lignes à mesure que de nouvelles données arrivent. L'ancien état d'agrégation est supprimé une fois le threshold dépassé.

Terminé

L'état d'agrégation n'est pas abandonné. La query réécrit la table cible pour chaque Trigger.

Mode de résultat

Comportement

Ajouter

La query écrit des lignes dans la table cible une fois le watermark threshold dépassé. Toutes les écritures sont retardées en fonction du threshold de latence. L'ancien état d'agrégation est supprimé après le dépassement du threshold.

Mettre à jour

La query écrit des lignes dans la table cible à mesure que les résultats sont calculés, et la query peut mettre à jour et écraser des lignes à mesure que de nouvelles données arrivent. L'ancien état d'agrégation est supprimé une fois le threshold dépassé.

Terminé

L'état d'agrégation n'est pas abandonné. La query réécrit la table cible pour chaque Trigger.

Filigranes et modes de sortie pour les jointures Stream-Stream

Les jointures entre plusieurs streams ne prennent en charge que le mode d'ajout. Les queries écrivent les enregistrements correspondants pour chaque batch.

Pour les jointures internes, Databricks vous recommande de définir un seuil de threshold sur chaque source de données de streaming pour permettre à la query de rejeter les informations d'état pour les anciens enregistrements. Sans filigranes, Structured Streaming tente de joindre chaque clé des deux côtés de la jointure sur chaque Trigger, ce qui pourrait affecter les performances.

Pour les jointures externes, le watermarking est obligatoire. Lorsqu’un enregistrement ne correspond pas, la query écrit une valeur nulle pour cette clé. Étant donné que les jointures ne prennent en charge que le mode d’ajout, les enregistrements non correspondants ne sont pas écrits tant que le threshold de retard n’est pas dépassé.

Contrôler le threshold des données tardives avec une politique de watermarks multiples

Pour plusieurs entrées Structured Streaming, vous pouvez définir plusieurs filigranes pour contrôler les seuils de tolérance pour les données tardives. Les filigranes vous permettent de contrôler les informations d'état et la latence.

Une query de streaming peut avoir plusieurs flux d'entrée qui sont unis ou joints ensemble. Pour les Opérations avec état, chaque Stream d'entrée pourrait nécessiter un threshold différent pour la tolérance aux données tardives. Spécifiez ces threshold en utilisant withWatermark("eventTime", delay) sur chaque flux d'entrée. Voici un exemple de query avec des jointures Stream-Stream.

Python
input_stream1 = ...      # delays up to 1 hour
input_stream2 = ... # delays up to 2 hours

(input_stream1.withWatermark("eventTime1", "1 hour")
.join(
input_stream2.withWatermark("eventTime2", "2 hours"),
joinCondition)
)

Lors de l'exécution de la query avec des Opérations à état, Structured Streaming suit individuellement le temps maximal d'événement pour chaque Stream d'entrée, calcule les filigranes en fonction du délai correspondant et détermine un filigrane global unique. Par default, Structured Streaming utilise le minimum comme filigrane global. Si un Stream prend du retard par rapport aux autres, un filigrane global minimum empêche la query de marquer accidentellement des données comme étant en retard. Par exemple, cela peut se produire lorsqu'un des flux cesse de recevoir des données en raison de défaillances en amont. Le filigrane global se déplace en toute sécurité au rythme du Stream le plus lent et retarde la sortie de la requête si nécessaire.

Pour réduire la latence, définissez spark.sql.streaming.multipleWatermarkPolicy sur max (la valeur par default est min) afin d'utiliser le watermark du Stream's le plus rapide comme watermark global. Cependant, cette configuration supprime les données des flux les plus lents. Databricks vous recommande d’appliquer cette configuration avec prudence.

Appliquer des filigranes aux opérations distinctes

L'distinct Opération suit chaque enregistrement unique dans l'état. Sans watermark, l'état croît indéfiniment et peut entraîner des problèmes de mémoire. Spécifiez un watermark sur un champ Timestamp pour délimiter l'état et supprimer les anciens enregistrements une fois que le threshold est passé.

L’exemple suivant applique un filigrane à une opération distinct :

Python
streamingDf = spark.readStream. ...  # columns: eventTime, id, value, ...

# Apply watermark before distinct operation
(streamingDf
.withWatermark("eventTime", "1 hour")
.distinct()
)

Dans cet exemple, la query de streaming supprime les enregistrements en double qui arrivent dans l'heure suivant le eventTime observé le plus récent. La query supprime les informations d'état pour la déduplication après le dépassement du threshold.

important

Pour dédupliquer des colonnes spécifiques au lieu de toutes les colonnes, utilisez dropDuplicates() ou dropDuplicatesWithinWatermark() au lieu de distinct. Consultez Supprimer les doublons dans le filigrane.

Supprimer les doublons dans le watermark

Dans Databricks Runtime 13.3 LTS ou supérieur, vous pouvez utiliser un identifiant unique pour dédupliquer les enregistrements au sein d'un threshold de filigrane.

Structured Streaming garantit un traitement exactement une fois, mais ne déduplique pas les enregistrements des sources de données. Utilisez dropDuplicatesWithinWatermark pour supprimer les doublons sur n'importe quel champ, même lorsque les champs diffèrent entre les enregistrements en double, comme l'heure de l'événement ou l'heure d'arrivée.

Avec dropDuplicatesWithinWatermark, les queries dédupliquent toujours les enregistrements qui arrivent dans le watermark threshold. Les queries peuvent également dédupliquer les enregistrements qui arrivent en dehors du threshold, mais ce n’est pas garanti. Pour garantir que les queries suppriment tous les doublons, définissez le watermark threshold de manière à ce qu'il soit supérieur à la différence de Timestamp maximale entre les événements en double.

Vous devez spécifier un watermark pour utiliser la méthode dropDuplicatesWithinWatermark :

Python
streamingDf = spark.readStream. ...

# deduplicate using guid column with watermark based on eventTime column
(streamingDf
.withWatermark("eventTime", "10 hours")
.dropDuplicatesWithinWatermark(["guid"])
)

Exemples de cas d'utilisation

Les exemples suivants présentent des cas d'utilisation avancés du fenêtrage :

Utiliser les fenêtres glissantes pour calculer les totaux de ventes horaires

Les fenêtres basculantes sont de taille fixe avec des intervalles non chevauchants. Chaque ligne d'entrée appartient à exactement une fenêtre. Utilisez des fenêtres basculantes pour compute des agrégations discrètes sur des périodes, telles que les totaux des ventes horaires :

Python
from pyspark.sql.functions import window, sum

hourly_sales = (orders
.withWatermark("timestamp", "1 hour")
.groupBy(window("timestamp", "1 hour"))
.agg(sum("amount").alias("total_sales"))
)

Dans cet exemple :

  • window("timestamp", "1 hour") regroupe les commandes en intervalles d'une heure sans chevauchement, tels que de 5 h à 6 h et de 6 h à 7 h.
  • withWatermark("timestamp", "1 hour") garde l'agrégat de chaque fenêtre en état jusqu'à ce que le timestamp de fin de fenêtre soit supérieur d'une heure au timestamp de commande maximum.

Utiliser les fenêtres glissantes pour calculer des agrégats glissants

Les fenêtres glissantes sont de taille fixe avec des intervalles qui peuvent se chevaucher. Une seule ligne peut appartenir à plusieurs fenêtres. Utilisez des fenêtres glissantes pour calculer des agrégats mobiles, tels que les ventes sur une période mobile de 6 heures :

Python
from pyspark.sql.functions import window, sum

rolling_sales = (orders
.withWatermark("timestamp", "1 hour")
.groupBy(window("timestamp", "6 hours", slideDuration="1 hour"))
.agg(sum("amount").alias("total_sales"))
)

Dans cet exemple :

  • window("timestamp", "6 hours", slideDuration="1 hour") regroupe les commandes en intervalles de 6 heures qui avancent d'une heure, par exemple, de 5 h à 11 h et de 6 h à 12 h.
  • withWatermark("timestamp", "1 hour") garde l'agrégat de chaque fenêtre en état jusqu'à ce que le timestamp de fin de fenêtre soit supérieur d'une heure au timestamp de commande maximum.
  • slideDuration doit être inférieur ou égal à windowDuration.

Utiliser les fenêtres de session pour vérifier l'activité de l'utilisateur

Les fenêtres de session n'ont pas de taille fixe. Une fenêtre s'ouvre lorsqu'une ligne arrive et se ferme après une durée d'intervalle sans nouvelles lignes. Utilisez les fenêtres de session pour agréger les pics d'activité entre de longues périodes d'inactivité, telles que les vues de page d'un utilisateur au cours d'une période de 30 minutes :

Python
from pyspark.sql.functions import session_window, sum

sessionized_page_views = (activity
.withWatermark("timestamp", "1 hour")
.groupBy("user_id", session_window("timestamp", gapDuration="30 minutes"))
.agg(sum("page_views").alias("total_page_views"))
)

Dans cet exemple :

  • session_window("timestamp", gapDuration="30 minutes") Ouvre une fenêtre lorsque la première vue de page arrive. Chaque vue de page ultérieure qui arrive dans les 30 minutes étend la fenêtre. Si aucune vue de page n'arrive dans les 30 minutes, la fenêtre se ferme et la vue de page suivante start une nouvelle fenêtre.
  • withWatermark("timestamp", "1 hour") maintient l'agrégat de chaque session en état jusqu'à ce que le Timestamp de fin de fenêtre soit supérieur d'1 heure au Timestamp maximal de la vue de page.
  • L'argument timeColumn pour window() et session_window() doit être de TimestampType ou TimestampNTZType.
  • Utilisez current_timestamp() pour définir des fenêtres basées sur le temps de traitement plutôt que sur le temps de l’événement.
  • Vous pouvez définir des durées de fenêtre allant des microsecondes jusqu’aux jours. Les durées d’un mois et plus ne sont pas prises en charge.
  • Utilisez le mode de sortie complete avec des agrégations par fenêtre pour conserver indéfiniment l'état de toutes les fenêtres. Utilisez le mode de sortie append avec un filigrane approprié pour limiter la croissance de l'état et éviter les problèmes de mémoire pour les grands ensembles de données. Pour plus de détails sur le comportement du mode de sortie, consultez Filigranes et mode de sortie pour les agrégations par fenêtre.