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 |
|---|---|---|
| str | le nom de la colonne qui contient l'heure de l'événement de la ligne. |
| 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
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]