Aller au contenu principal

Observer

Définissez des métriques (nommées) à observer sur le DataFrame. Cette méthode renvoie un DataFrame « observé » qui renvoie le même résultat que l'entrée, avec les garanties suivantes : il va calculer les agrégats définis (métriques) sur toutes les données qui traversent le dataset à ce moment-là. Il affichera la valeur des colonnes agrégées définies dès que nous atteindrons un point d'achèvement.

Syntaxe​

observe(observation: Union["Observation", str], *exprs: Column)

parameter​

parameter

Type

Description

observation

Observation ou chaîne

str pour spécifier le nom, ou une instance Observation pour obtenir la métrique.

exprs

Colonne

expressions de colonne (Column).

parameter

Type

Description

observation

Observation ou chaîne

str pour spécifier le nom, ou une instance Observation pour obtenir la métrique.

exprs

Colonne

expressions de colonne (Column).

Renvoie​

DataFrame: le DataFrame observé.

Notes​

Lorsque observation est Observation, cette méthode ne prend en charge que les batch query. Lorsque observation est une chaîne de caractères, cette méthode fonctionne pour les requêtes par batch et en streaming. L'exécution continue n'est actuellement pas encore prise en charge.

Exemples​

Python
from pyspark.sql import Observation, functions as sf
df = spark.createDataFrame([(2, "Alice"), (5, "Bob")], schema=["age", "name"])
observation = Observation("my metrics")
observed_df = df.observe(observation,
sf.count(sf.lit(1)).alias("count"), sf.max("age"))
observed_df.count()
# 2
observation.get
# {'count': 2, 'max(age)': 5}