Aller au contenu principal

withWatermark

Définit un filigrane temporel d'événement pour ce DataFrame. Un filigrane suit un point dans le temps avant lequel nous supposons qu'aucune donnée tardive supplémentaire n'arrivera.

Syntaxe

withWatermark(eventTime: str, delayThreshold: str)

parameter

parameter

Type

Description

eventTime

str

le nom de la colonne qui contient l'heure de l'événement de la ligne.

delayThreshold

str

le délai minimal à attendre pour que les données arrivent en retard, par rapport au dernier enregistrement qui a été traité sous forme d'intervalle (par exemple... « 1 minute » ou « 5 heures »).

parameter

Type

Description

eventTime

str

le nom de la colonne qui contient l'heure de l'événement de la ligne.

delayThreshold

str

le délai minimal à attendre pour que les données arrivent en retard, par rapport au dernier enregistrement qui a été traité sous forme d'intervalle (par exemple... « 1 minute » ou « 5 heures »).

Renvoie

DataFrame: DataFrame filigrané.

Notes

Cette fonctionnalité est uniquement destinée à Structured Streaming.

Spark utilisera ce filigrane à plusieurs fins :

  • Savoir quand une agrégation de fenêtre temporelle donnée peut être finalisée et donc émise lors de l'utilisation de modes de sortie qui n'autorisent pas les mises à jour.
  • Pour minimiser la quantité d'état que nous devons conserver pour les agrégations en cours.

Le filigrane actuel est calculé en examinant le MAX(eventTime) observé sur l'ensemble des partitions dans la query, moins un delayThreshold spécifié par l'utilisateur. En raison du coût de la coordination de cette valeur entre les partitions, le filigrane réel utilisé n'est garanti d'être en retard d'au moins delayThreshold par rapport à l'heure réelle de l'événement.

Exemples

Python
from pyspark.sql import Row
from pyspark.sql.functions import timestamp_seconds
df = spark.readStream.format("rate").load().selectExpr(
"value % 5 AS value", "timestamp")
df.select("value", df.timestamp.alias("time")).withWatermark("time", '10 minutes')
# DataFrame[value: bigint, time: timestamp]