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}