fonctions définies par l'utilisateur pandas
Une fonction définie par l'utilisateur (UDF) pandas — également appelée UDF vectorisée — est une fonction définie par l'utilisateur qui utilise Apache Arrow pour transférer des données et pandas pour travailler avec les données. Les UDF pandas permettent des opérations vectorisées qui peuvent augmenter les performances jusqu'à 100 fois par rapport aux UDF Python ligne par ligne.
Pour plus d'information, consultez le billet de blog Nouvelles annotations de type Python et UDF Pandas dans la prochaine version d'Apache Spark 3.0.
Vous définissez une UDF pandas en utilisant le mot-clé pandas_udf comme décorateur et enveloppez la fonction avec une annotation de type Python.
Cet article décrit les différents types d'UDF pandas et montre comment utiliser les UDF pandas avec des annotations de type.
UDF Série à Série
Vous utilisez une UDF pandas Série à Série pour vectoriser les opérations scalaires.
Vous pouvez les utiliser avec des APIs telles que select et withColumn.
La fonction Python doit accepter une série pandas comme entrée et renvoyer une série pandas de même longueur. Spécifiez ces types à l'aide d'annotations de type Python. Spark exécute une UDF pandas en divisant les données en **batches** de lignes, en appelant la fonction pour chaque **batch**, puis en concaténant les résultats.
L'exemple suivant montre comment créer une UDF pandas qui calcule le produit de 2 colonnes.
import pandas as pd
from pyspark.sql.functions import col, pandas_udf
from pyspark.sql.types import LongType
# Declare the function and create the UDF
def multiply_func(a: pd.Series, b: pd.Series) -> pd.Series:
return a * b
multiply = pandas_udf(multiply_func, returnType=LongType())
# The function for a pandas_udf should be able to execute with local pandas data
x = pd.Series([1, 2, 3])
print(multiply_func(x, x))
# 0 1
# 1 4
# 2 9
# dtype: int64
# Create a Spark DataFrame, 'spark' is an existing SparkSession
df = spark.createDataFrame(pd.DataFrame(x, columns=["x"]))
# Execute function as a Spark vectorized UDF
df.select(multiply(col("x"), col("x"))).show()
# +-------------------+
# |multiply_func(x, x)|
# +-------------------+
# | 1|
# | 4|
# | 9|
# +-------------------+
Iterator of Series to Iterator of Series UDF
Une UDF d'itérateur est identique à une UDF pandas scalaire, à l'exception de :
-
La fonction Python
- Prend un itérateur de batchs au lieu d'un seul batch d'entrée en tant qu'entrée.
- Renvoie un itérateur de batchs de sortie au lieu d'un seul batch de sortie.
-
La longueur de la sortie entière dans l'itérateur doit être la même que la longueur de l'entrée entière.
-
L'UDF pandas encapsulée prend une seule colonne Spark comme entrée.
Vous devez spécifier l'annotation de type Python comme
Iterator[pandas.Series] -> Iterator[pandas.Series].
Cette UDF pandas est utile lorsque l'exécution de l'UDF nécessite l'initialisation d'un état, par exemple, le chargement d'un fichier de modèle de machine learning pour appliquer l'inférence à chaque batch d'entrée.
L'exemple suivant montre comment créer une UDF pandas avec prise en charge des itérateurs.
import pandas as pd
from typing import Iterator
from pyspark.sql.functions import col, pandas_udf, struct
pdf = pd.DataFrame([1, 2, 3], columns=["x"])
df = spark.createDataFrame(pdf)
# When the UDF is called with the column,
# the input to the underlying function is an iterator of pd.Series.
@pandas_udf("long")
def plus_one(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
for x in batch_iter:
yield x + 1
df.select(plus_one(col("x"))).show()
# +-----------+
# |plus_one(x)|
# +-----------+
# | 2|
# | 3|
# | 4|
# +-----------+
# In the UDF, you can initialize some state before processing batches.
# Wrap your code with try/finally or use context managers to ensure
# the release of resources at the end.
y = 1 # value captured by the UDF closure
@pandas_udf("long")
def plus_y(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
try:
for x in batch_iter:
yield x + y
finally:
pass # release resources here, if any
df.select(plus_y(col("x"))).show()
# +---------+
# |plus_y(x)|
# +---------+
# | 2|
# | 3|
# | 4|
# +---------+
Itérateur de plusieurs Series vers un itérateur de Series UDF
Un Iterator of multiple Series to Iterator of Series UDF a des caractéristiques et des restrictions similaires à un Iterator of Series to Iterator of Series UDF. La fonction spécifiée prend un itérateur de batches et produit un itérateur de batches. C’est également utile lorsque l’exécution de l’UDF nécessite l’initialisation d’un certain état.
Les différences sont :
- La fonction Python sous-jacente prend un itérateur d'un tuple de pandas Series.
- L'UDF pandas encapsulée prend plusieurs colonnes Spark en entrée.
Vous spécifiez les indications de type comme Iterator[Tuple[pandas.Series, ...]] -> Iterator[pandas.Series].
from typing import Iterator, Tuple
import pandas as pd
from pyspark.sql.functions import col, pandas_udf, struct
pdf = pd.DataFrame([1, 2, 3], columns=["x"])
df = spark.createDataFrame(pdf)
@pandas_udf("long")
def multiply_two_cols(
iterator: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
for a, b in iterator:
yield a * b
df.select(multiply_two_cols("x", "x")).show()
# +-----------------------+
# |multiply_two_cols(x, x)|
# +-----------------------+
# | 1|
# | 4|
# | 9|
# +-----------------------+
UDF de série à scalaire
Les UDF pandas de séries à scalaires sont similaires aux fonctions d'agrégation Spark.
Une UDF pandas de série à scalaire définit une agrégation d'une ou plusieurs séries pandas vers une valeur scalaire, où chaque série pandas représente une colonne Spark.
Vous utilisez une UDF pandas scalaire de type Series avec des APIs telles que select, withColumn, groupBy.agg et
pyspark.sql.Window.
Vous exprimez l’indication de type comme pandas.Series, ... -> Any. Le type de retour doit être un type de données primitif, et le scalaire renvoyé peut être soit un type primitif Python, par exemple, int ou float, soit un type de données NumPy tel que numpy.int64 ou numpy.float64. Any devrait idéalement
être un type scalaire spécifique.
Ce type d'UDF ne prend pas en charge l'agrégation partielle et toutes les données de chaque groupe sont chargées en mémoire.
L'exemple suivant montre comment utiliser ce type d'UDF pour compute la moyenne avec les opérations select, groupBy et window :
import pandas as pd
from pyspark.sql.functions import pandas_udf
from pyspark.sql import Window
df = spark.createDataFrame(
[(1, 1.0), (1, 2.0), (2, 3.0), (2, 5.0), (2, 10.0)],
("id", "v"))
# Declare the function and create the UDF
@pandas_udf("double")
def mean_udf(v: pd.Series) -> float:
return v.mean()
df.select(mean_udf(df['v'])).show()
# +-----------+
# |mean_udf(v)|
# +-----------+
# | 4.2|
# +-----------+
df.groupby("id").agg(mean_udf(df['v'])).show()
# +---+-----------+
# | id|mean_udf(v)|
# +---+-----------+
# | 1| 1.5|
# | 2| 6.0|
# +---+-----------+
w = Window \
.partitionBy('id') \
.rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
df.withColumn('mean_v', mean_udf(df['v']).over(w)).show()
# +---+----+------+
# | id| v|mean_v|
# +---+----+------+
# | 1| 1.0| 1.5|
# | 1| 2.0| 1.5|
# | 2| 3.0| 6.0|
# | 2| 5.0| 6.0|
# | 2|10.0| 6.0|
# +---+----+------+
Pour une utilisation détaillée, consultez pyspark.sql.functions.pandas_udf.
Utilisation
Définition de la taille de batch Arrow
Cette configuration n'a aucun impact sur le compute serverless ou sur le compute configuré avec le mode d'accès standard et Databricks Runtime de 13.3 LTS à 14.2. Sur le compute serverless, la plateforme gère en interne le dimensionnement des batchs Arrow.
Les partitions de données dans Spark sont converties en lots d'enregistrements Arrow, ce qui peut temporairement entraîner une utilisation élevée de la mémoire dans la JVM. Pour éviter d'éventuelles exceptions de mémoire insuffisante, vous pouvez ajuster la taille des lots d'enregistrements Arrow en définissant la configuration spark.sql.execution.arrow.maxRecordsPerBatch sur un entier qui détermine le nombre maximal de lignes pour chaque batch. La valeur par default est de 10 000 enregistrements par lot. Si le nombre de colonnes est élevé, la valeur doit être ajustée en conséquence. En utilisant cette limite, chaque partition de données est divisée en 1 ou plusieurs lots d'enregistrements pour le traitement.
Timestamp avec sémantique de fuseau horaire
Spark stocke en interne les Timestamp en tant que valeurs UTC, et les données de Timestamp importées sans fuseau horaire spécifié sont converties de l'heure locale en UTC avec une résolution en microsecondes.
Lorsque les données de timestamp sont exportées ou affichées dans Spark, le fuseau horaire de session est utilisé pour localiser les valeurs de timestamp. Le fuseau horaire de session est défini avec la configuration spark.sql.session.timeZone et default est le fuseau horaire local système de la JVM. pandas utilise un type datetime64 avec une résolution en nanosecondes, datetime64[ns], avec un fuseau horaire facultatif par colonne.
Lorsque les données d'horodatage sont transférées de Spark à pandas, elles sont converties en nanosecondes et chaque colonne est convertie dans le fuseau horaire de la session Spark, puis localisée à ce fuseau horaire, ce qui supprime le fuseau horaire et affiche les valeurs comme heure locale. Cela se produit lors de l’appel de toPandas() ou pandas_udf avec des colonnes de timestamp.
Lorsque les données Timestamp sont transférées de pandas vers Spark, elles sont converties en microsecondes UTC. Cela se produit lors de l'appel
createDataFrame avec un DataFrame pandas ou lors du retour d'un
timestamp à partir d'une UDF pandas. Ces conversions sont effectuées
automatiquement pour garantir que Spark dispose des données dans le format attendu, il
n'est donc pas nécessaire d'effectuer ces conversions vous-même. Toutes les valeurs en nanosecondes sont tronquées.
Une UDF standard charge les données Timestamp en tant qu'objets datetime Python, ce qui est différent d'un Timestamp pandas. Pour obtenir les meilleures performances, nous vous recommandons d'utiliser les fonctionnalités de séries chronologiques de pandas lorsque vous travaillez avec des Timestamp dans une UDF pandas. Pour plus de détails, consultez la fonctionnalité Série temporelle / Date.
Exemple de Notebook
Le Notebook suivant illustre les améliorations de performances que vous pouvez obtenir avec les UDF pandas :