Aller au contenu principal

dropDuplicatesWithinWatermark

Renvoie un nouveau DataFrame avec les lignes dupliquées supprimées, en tenant éventuellement compte uniquement de certaines colonnes, dans les limites du filigrane.

Syntaxe

dropDuplicatesWithinWatermark(subset: Optional[List[str]] = None)

parameter

parameter

Type

Description

subset

Liste des noms de colonne, facultatif

Liste des colonnes à utiliser pour la comparaison des doublons (default toutes les colonnes).

parameter

Type

Description

subset

Liste des noms de colonne, facultatif

Liste des colonnes à utiliser pour la comparaison des doublons (default toutes les colonnes).

Renvoie

DataFrame: DataFrame sans doublons.

Notes

Ceci ne fonctionne qu'avec les streaming DataFrame, et le filigrane du DataFrame d'entrée doit être défini via withWatermark.

Pour un DataFrame de streaming, cela conservera toutes les données sur l'ensemble des triggers comme état intermédiaire afin de supprimer les lignes dupliquées. L'état sera conservé pour garantir la sémantique : « Les événements sont dédupliqués tant que la distance temporelle entre les événements les plus anciens et les plus récents est inférieure au threshold de délai du watermark. » Il est recommandé aux utilisateurs de définir le threshold de délai du watermark plus long que les différences maximales de timestamp entre les événements dupliqués.

Note : les données arrivées trop tard et plus anciennes que le filigrane seront abandonnées.

Prend en charge Spark Connect.

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]

df.dropDuplicatesWithinWatermark()

df.dropDuplicatesWithinWatermark(['value'])