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 ou chaîne |
|
| 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
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}