Que sont les fonctions définies par l'utilisateur (UDF) ?
Les fonctions définies par l'utilisateur (UDF) vous permettent de réutiliser et de partager du code qui étend les capacités intégrées sur Databricks. Utilisez les UDF pour effectuer des tâches spécifiques telles que des calculs complexes, des transformations ou des manipulations de données personnalisées.
Quand utiliser une UDF ou une fonction Apache Spark ?
Utilisez des UDF pour la logique difficile à exprimer avec les fonctions Apache Spark intégrées. Les fonctions Apache Spark intégrées sont optimisées pour le traitement distribué et offrent de meilleures performances à grande échelle. Pour plus d'informations, consultez Fonctions.
Databricks recommande les UDF pour les query ad hoc, le nettoyage manuel des données, l'analyse exploratoire des données et les Opérations sur des datasets de petite à moyenne taille. Les cas d'utilisation courants des UDF incluent le chiffrement des données, le déchiffrement, le hachage, l'analyse JSON et la validation.
Utilisez les méthodes Apache Spark pour les opérations sur de grands datasets et toute charge de travail exécutée régulièrement ou en continu, y compris les Jobs ETL et les opérations de streaming.
Comprendre les types d'UDF
Sélectionnez un type d'UDF dans les tabs suivantes pour afficher une description, un exemple et un link pour en savoir plus.
- Scalar UDF
- Batch Scalar UDFs
- Non-Scalar UDFs
- UDAF
- UDTFs
Les UDF scalaires opèrent sur une seule ligne et renvoient une seule valeur de résultat pour chaque ligne. Ils peuvent être régis par Unity Catalog ou limités à la session.
L'exemple suivant utilise une UDF scalaire pour calculer la longueur de chaque nom dans une colonne name et ajouter la valeur dans une nouvelle colonne name_length.
+-------+-------+
| name | score |
+-------+-------+
| alice | 10.0 |
| bob | 20.0 |
| carol | 30.0 |
| dave | 40.0 |
| eve | 50.0 |
+-------+-------+
-- Create a SQL UDF for name length
CREATE OR REPLACE FUNCTION main.test.get_name_length(name STRING)
RETURNS INT
RETURN LENGTH(name);
-- Use the UDF in a SQL query
SELECT name, main.test.get_name_length(name) AS name_length
FROM your_table;
+-------+-------+-------------+
| name | score | name_length |
+-------+-------+-------------+
| alice | 10.0 | 5 |
| bob | 20.0 | 3 |
| carol | 30.0 | 5 |
| dave | 40.0 | 4 |
| eve | 50.0 | 3 |
+-------+-------+-------------+
Pour implémenter ceci dans un Notebook Databricks en utilisant PySpark :
from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType
@udf(returnType=IntegerType())
def get_name_length(name):
return len(name)
df = df.withColumn("name_length", get_name_length(df.name))
# Show the result
display(df)
Consultez les fonctions définies par l'utilisateur (UDF) SQL et Python dans Unity Catalog et les fonctions scalaires définies par l'utilisateur (UDF) Python.
Traitez les données par batchs tout en conservant une parité 1:1 des lignes d'entrée/sortie. Cela réduit la surcharge des opérations ligne par ligne pour le traitement des données à grande échelle. Les UDF de batch conservent également l'état entre les lots pour s'exécuter plus efficacement, réutiliser les ressources et gérer les calculs complexes qui nécessitent un contexte à travers les blocs de données.
Ils peuvent être régis par Unity Catalog ou être limités à la session.
La fonction UDF Python Batch Unity Catalog suivante calcule l'IMC lors du traitement par batch des lignes :
+-------------+-------------+
| weight_kg | height_m |
+-------------+-------------+
| 90 | 1.8 |
| 77 | 1.6 |
| 50 | 1.5 |
+-------------+-------------+
%sql
CREATE OR REPLACE FUNCTION main.test.calculate_bmi_pandas(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
AS $$
import pandas as pd
from typing import Iterator, Tuple
def handler_function(batch_iter: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
for weight_series, height_series in batch_iter:
yield weight_series / (height_series ** 2)
$$;
select main.test.calculate_bmi_pandas(cast(70 as double), cast(1.8 as double));
+--------+
| BMI |
+--------+
| 27.8 |
| 30.1 |
| 22.2 |
+--------+
Consultez les fonctions définies par l'utilisateur (UDF) SQL et Python dans Unity Catalog et les fonctions définies par l'utilisateur (UDF) Python par batch dans Unity Catalog.
Les UDF non scalaires fonctionnent sur des datasets/colonnes entiers avec des ratios d’entrée/sortie flexibles (1 ou de nombreux).
Les UDF pandas par batch à l'échelle de la session peuvent être des types suivants :
- Série à série
- Itérateur de série à itérateur de série
- Itérateur de plusieurs Series à Itérateur de Series
- Série vers scalaire
Voici un exemple d'UDF pandas de Series à Series.
from pyspark.sql.functions import pandas_udf
import pandas as pd
df = spark.createDataFrame([(70, 1.75), (80, 1.80), (60, 1.65)], ["Weight", "Height"])
@pandas_udf("double")
def calculate_bmi_pandas(weight: pd.Series, height: pd.Series) -> pd.Series:
return weight / (height ** 2)
df.withColumn("BMI", calculate_bmi_pandas(df["Weight"], df["Height"])).display()
Les UDAFs opèrent sur plusieurs lignes et renvoient un résultat agrégé unique. Les UDAFs sont uniquement à portée de session.
L'exemple UDAF suivant agrège les scores par longueur du nom.
from pyspark.sql.functions import pandas_udf
from pyspark.sql import SparkSession
import pandas as pd
# Define a pandas UDF for aggregating scores
@pandas_udf("int")
def total_score_udf(scores: pd.Series) -> int:
return scores.sum()
# Group by name length and aggregate
result_df = (df.groupBy("name_length")
.agg(total_score_udf(df["score"]).alias("total_score")))
display(result_df)
+-------------+-------------+
| name_length | total_score |
+-------------+-------------+
| 3 | 70.0 |
| 4 | 40.0 |
| 5 | 40.0 |
+-------------+-------------+
Consultez les fonctions définies par l'utilisateur pandas pour Python et les fonctions d'agrégation définies par l'utilisateur Scala (UDAFs).
Une UDTF prend un ou plusieurs arguments d'entrée et renvoie plusieurs lignes (et éventuellement plusieurs colonnes) pour chaque ligne d'entrée. Ils peuvent être régis par Unity Catalog ou limités à la session.
La fonction UDTF suivante crée une table en utilisant une liste fixe de deux arguments entiers :
CREATE OR REPLACE FUNCTION get_sum_diff(x INT, y INT)
RETURNS TABLE (sum INT, diff INT)
LANGUAGE PYTHON
HANDLER 'GetSumDiff'
AS $$
class GetSumDiff:
def eval(self, x: int, y: int):
yield x + y, x - y
$$;
SELECT * FROM get_sum_diff(10, 3);
+-----+------+
| sum | diff |
+-----+------+
| 13 | 7 |
+-----+------+
Pour implémenter ceci dans un Notebook Databricks en utilisant PySpark :
from pyspark.sql.functions import lit, udtf
@udtf(returnType="sum: int, diff: int")
class GetSumDiff:
def eval(self, x: int, y: int):
yield x + y, x - y
GetSumDiff(lit(1), lit(2)).show()
Voir les UDTF de Unity Catalog et les UDTF à étendue de session.
UDFs régies par Unity Catalog par rapport aux UDFs délimitées par session
Unity Catalog persiste les UDF régies par Unity Catalog pour une gouvernance, une réutilisation et une découvrabilité améliorées. Vous définissez les UDFs limitées à la session dans un notebook ou un job, limitées à la SparkSession actuelle. Vous pouvez définir et accéder aux UDFs limitées à la session à l'aide de SQL, Python ou Scala.
Utilisez le tableau suivant pour décider entre les deux catégories, puis consultez les aide-mémoires qui suivent pour les types spécifiques d'UDF dans chacun d'eux.
Considération | UDF régies par Unity Catalog | UDFs délimités à la session |
|---|---|---|
Idéal pour | Partage sécurisé de fonctions entre les équipes, les Notebooks, les Jobs et les SQL Warehouse. | Développement rapide et itératif dans un seul Notebook ou Job. |
Langages | SQL, Python, Scala et Java. | SQL, Python et Scala. |
Gouvernance et partage | Régies par les autorisations d'Unity Catalog et découvrables dans l'Explorateur de catalogues. | Limité à la SparkSession actuelle. Non régi ni partagé. |
Persistance | Persisté dans Unity Catalog et réutilisable entre les sessions. | Existe uniquement pour la session actuelle. |
Aide-mémoire des UDF régies par Unity Catalog
Les fonctions définies par l'utilisateur (UDF) régies par Unity Catalog permettent de définir, d'utiliser, de partager en toute sécurité et de régir des fonctions personnalisées dans tous les environnements de calcul. Consultez les fonctions définies par l'utilisateur (UDF) SQL et Python dans Unity Catalog.
Type d'UDF | Compute pris en charge | Description |
|---|---|---|
UDF Python Unity Catalog |
| Définissez une UDF en Python et enregistrez-la dans Unity Catalog pour la gouvernance. Les fonctions UDF scalaires opèrent sur une seule ligne et renvoient une seule valeur de résultat pour chaque ligne. |
UDF Python Unity Catalog en batch |
| Définissez une UDF en Python et enregistrez-la dans Unity Catalog pour la gouvernance. Opérations batch sur plusieurs valeurs et retour de plusieurs valeurs. Réduit la surcharge des opérations ligne par ligne pour le traitement de données à grande échelle. |
Unity Catalog Python UDTF |
| Définir une UDTF en Python et l'enregistrer dans Unity Catalog pour la gouvernance. Une UDTF prend un ou plusieurs arguments d'entrée et renvoie plusieurs lignes (et éventuellement plusieurs colonnes) pour chaque ligne d'entrée. |
UDF Unity Catalog Scala ou Java |
| Définissez une UDF en Scala ou en Java et enregistrez-la dans Unity Catalog pour la gouvernance. Les UDF scalaires opèrent sur une seule ligne et renvoient une seule valeur de résultat pour chaque ligne. Nécessite Scala 2.13.16, JDK 17 et version 4 de l'environnement. |
Aide-mémoire des UDFs à portée de session pour le compute isolé de l'utilisateur
Vous définissez des UDF délimitées à la session dans un Notebook ou un Job, délimitées à la SparkSession actuelle. Vous pouvez définir et accéder à des UDF délimitées à la session en utilisant SQL, Python ou Scala.
Type d'UDF | Compute pris en charge | Description |
|---|---|---|
Python scalaire |
| Les fonctions UDF scalaires opèrent sur une seule ligne et renvoient une seule valeur de résultat pour chaque ligne. |
Python non scalaire |
| Les UDF non scalaires incluent |
UDTF Python |
| Une UDTF prend un ou plusieurs arguments d'entrée et renvoie plusieurs lignes (et éventuellement plusieurs colonnes) pour chaque ligne d'entrée. |
UDF scalaires Scala |
| Les fonctions UDF scalaires opèrent sur une seule ligne et renvoient une seule valeur de résultat pour chaque ligne. |
UDF Scala ou Java à partir d'un JAR |
| Enregistrer une classe UDF précompilée depuis un JAR à l'aide de |
UDAFs Scala |
| Les UDAFs opèrent sur plusieurs lignes et renvoient un seul résultat agrégé. |
Considérations relatives aux performances
-
Fonctions intégrées et UDF SQL sont les options les plus efficaces.
-
Les UDF Scala sont généralement plus rapides que les UDF Python.
- Les fonctions UDF Scala non isolées s’exécutent dans la Machine virtuelle Java (JVM), ce qui leur évite la surcharge liée au déplacement des données à l’intérieur et à l’extérieur de la JVM.
- Les UDF Scala isolées doivent faire transiter les données à l'intérieur et à l'extérieur de la JVM, mais elles peuvent toujours être plus rapides que les UDF Python parce qu'elles gèrent la mémoire plus efficacement.
-
UDF Python et UDF pandas ont tendance à être plus lents que les UDF Scala car ils doivent sérialiser les données et les déplacer hors de la JVM vers l'interpréteur Python.
- Les UDF Pandas sont jusqu'à 100 fois plus rapides que les UDF Python car ils utilisent Apache Arrow pour réduire les coûts de sérialisation.