Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
Important
Scala et Java UDF peuvent être inscrits dans le catalogue Unity pour la gouvernance, la réutilisation et la détectabilité. Consultez les fonctions définies par l’utilisateur en Scala et Java (UDF) dans Unity Catalog.
Cette page explique comment créer des fonctions UDF scala et Java délimitées à la session dans Azure Databricks. Les UDF à portée de session sont créées dans un notebook ou un job et s’appliquent uniquement à la SparkSession en cours. Pour obtenir la référence du langage SQL, consultez fonctions scalaires définies par l’utilisateur externes (UDF) .
Choisir votre approche
Vous pouvez définir une fonction UDF Scala ou Java de la manière suivante. Pour comparer tous les types de fonctions UDF selon les langages, la gouvernance et le calcul, consultez Fonctions UDF régies par Unity Catalog et fonctions UDF limitées à la session.
| Approche | Description |
|---|---|
| Inline Scala UDF | Définissez une fonction UDF dans un notebook à l’aide d’une fonction Scala ou d’une fonction lambda. Limitée à la session. Non pris en charge avec le calcul sans serveur. |
| Java UDF à partir d’un fichier JAR | Enregistrer une classe UDF précompilée à partir d’une archive JAR à l’aide de spark.udf.registerJavaFunction. Limitée à la session. Pris en charge sur l’environnement serverless. |
| UDF Scala ou Java régie par Unity Catalog | Inscrivez une fonction UDF dans le catalogue Unity pour la gouvernance, la réutilisation et la détectabilité. Pris en charge sur l’environnement serverless. |
Spécifications
- Les UDF Scala sur un calcul compatible avec Unity Catalog en mode d’accès standard nécessitent Databricks Runtime 14.2 ou version ultérieure.
- La prise en charge des instances ARM pour les UDFs Scala sur Unity Catalog nécessite Databricks Runtime 15.2 ou une version ultérieure.
- L’inscription d’un Java UDF à partir d’un fichier JAR
spark.udf.registerJavaFunctionnécessite Databricks Runtime 18 LTS ou version ultérieure. Consultez Inscrire un Java UDF à partir d’un fichier JAR.
Important
Générez votre fichier JAR sur les mêmes versions Scala et Apache Spark que le calcul qui l’exécute. Une incompatibilité peut faire échouer l’UDF lors de l’enregistrement ou de l’appel.
- Calcul classique : correspond aux versions Scala et Spark de votre version databricks Runtime. Consultez la section Environnement système des notes de version et compatibilité de Databricks Runtime pour votre version. Par exemple, Databricks Runtime 18 LTS utilise Scala 2.13.16 et Apache Spark 4.0.
- Calcul sans serveur : faites correspondre la version Scala à la Version de l’environnement. Consultez les Versions de l’environnement serverless.
Marquez la dépendance Apache Spark de provided sorte qu’elle n’est pas regroupée dans votre fichier JAR. Incluez uniquement les dépendances tierces utilisées par votre UDF.
Inscrire une fonction en tant que fonction définie par l’utilisateur
Inscrivez une fonction Scala en tant que fonction UDF à l’aide spark.udf.registerde :
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Appeler la UDF dans Spark SQL
Créez une vue temporaire, puis appelez la fonction UDF dans une requête SQL :
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Utiliser une UDF avec des DataFrames
Vous pouvez également appeler une fonction UDF à l’aide de l’API DataFrame :
import org.apache.spark.sql.functions.{col, udf}
val squared = udf((s: Long) => s * s)
display(spark.range(1, 20).select(squared(col("id")) as "id_squared"))
Inscrire un Java UDF à partir d’un fichier JAR
Regroupez une UDF dans un fichier JAR, ajoutez-la à votre session avec spark.addArtifact et enregistrez la classe UDF avec spark.udf.registerJavaFunction.
Remarque
Pris en charge en mode d’accès standard et avec le calcul sans serveur à partir de Databricks Runtime 18 LTS. La fonction inscrite est limitée à la session et n’est pas inscrite dans le catalogue Unity.
Les étapes suivantes vous guident dans la création d’un projet, l’écriture d’une classe UDF, la création d’un fichier JAR gras et son inscription.
Étape 1 : Créer votre projet
Configurez un projet en Scala ou Java.
Scala
Créez un projet Scala à l’aide de sbt:
sbt new scala/scala-seed.g8
Remplacez le contenu de votre build.sbt fichier par ce qui suit. Définissez scalaVersion et la version spark-sql pour qu’ils correspondent à votre environnement de calcul :
scalaVersion := "2.13.16"
ThisBuild / organization := "com.example"
lazy val myUDF = (project in file("."))
.settings(
name := "my-udf",
libraryDependencies += "org.apache.spark" %% "spark-sql" % "4.0.0" % "provided"
)
Activez le plug-in sbt-assembly pour générer un fichier JAR fat. Créez ou modifiez project/assembly.sbt et ajoutez :
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Créez un projet Maven à l’aide de l’archétype de démarrage rapide :
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Cette commande crée la structure de projet Maven standard avec les répertoires src/main/java et src/test/java.
Dans le fichier généré pom.xml, dans les <project></project> balises, ajoutez un <properties> bloc et configurez-le maven-shade-plugin pour générer un fichier JAR gras :
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.5.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
Étape 2 : Écrire votre classe UDF
Votre classe UDF doit implémenter l’une des org.apache.spark.sql.api.java.UDF interfaces (UDF1 via UDF22), où le nombre indique le nombre d’arguments d’entrée que prend la fonction UDF. Implémentez la call() méthode avec votre logique.
Le gestionnaire doit être une classe Java.
spark.udf.registerJavaFunction charge la classe par réflexion. Il doit donc s’agir d’une classe publique de niveau supérieur (ou static imbriquée) avec un constructeur public no-arg public. Une Scala class ou object ne répond pas à cette exigence et échoue au moment de l’appel. Vous pouvez générer le fichier JAR avec sbt, mais la classe UDF elle-même doit être écrite dans Java.
Créez src/main/java/com/example/MyIntegerUDF.java :
package com.example;
import org.apache.spark.sql.api.java.UDF1;
public class MyIntegerUDF implements UDF1<Integer, Integer> {
@Override
public Integer call(Integer x) {
return x + 1;
}
}
Étape 3 : Générez votre JAR autonome
Empaqueter votre UDF compilé dans un fichier JAR gras.
Scala
À partir du répertoire racine de votre projet, exécutez :
sbt clean assembly
Le fichier JAR fat est créé dans target/scala-2.13/ sous un nom tel que my-udf-assembly-0.1.0-SNAPSHOT.jar.
Java
À partir du répertoire racine de votre projet, exécutez :
mvn clean package
Le fichier JAR fat est créé dans target/ sous un nom tel que my-udf-1.0-SNAPSHOT.jar.
Étape 4 : Charger votre fichier JAR dans un volume de catalogue Unity
Chargez le fichier JAR dans un volume de catalogue Unity afin que votre calcul puisse y accéder. Si vous n’avez pas encore de volume, créez-en un :
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Chargez votre fichier JAR sur le volume à l’aide de l’Explorateur de catalogues :
- Dans votre espace de travail Azure Databricks, cliquez sur
Catalogue pour ouvrir l’Explorateur de catalogues.
- Sélectionnez le catalogue, puis sélectionnez le schéma qui contient votre volume.
- Cliquez sur le nom du volume.
- Cliquez sur Charger sur ce volume et sélectionnez votre fichier JAR.
- Cliquez sur Télécharger.
- Une fois le chargement terminé, cliquez sur le nom de votre fichier JAR, puis cliquez sur Copier le chemin d’accès pour copier le chemin du volume. Par exemple :
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Vous avez besoin de ce chemin à l’étape suivante.
Étape 5 : Enregistrer et appeler la fonction UDF
Ajoutez le fichier JAR à votre session à l’aide de son chemin d’accès au volume, inscrivez la classe UDF et appelez-la à partir de Spark SQL :
# Add the JAR containing your UDF class to the session
spark.addArtifact("/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar")
# Register the UDF class, providing the SQL function name,
# the fully qualified class name, and the return type
from pyspark.sql.types import IntegerType
spark.udf.registerJavaFunction(
"my_udf",
"com.example.MyIntegerUDF",
IntegerType(),
)
# Call the UDF from Spark SQL
spark.sql("SELECT my_udf(21)").show()
En mode de calcul serverless et en mode d’accès standard, vous devez spécifier explicitement un type de retour. L’omission du type de retour échoue avec UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Les fonctions d’agrégation définies par l’utilisateur (UDAF) ne sont pas prises en charge avec registerJavaFunction.
La requête retourne la sortie UDF, confirmant que la fonction est inscrite et pouvant être appelée :
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Ordre d’évaluation et vérification de valeurs nulles
Spark SQL (y compris SQL et les API DataFrame et Dataset) ne garantit pas l’ordre d’évaluation des sous-expressions. Spark n’évalue pas les entrées d’un opérateur ou d’une fonction de gauche à droite. Les expressions logiques AND et OR n’ont pas de sémantique d’évaluation en court-circuit de gauche à droite.
Ne vous fiez pas aux effets secondaires ni à l’ordre d’évaluation des expressions booléennes, ni à l’ordre des clauses WHERE et HAVING. L’optimiseur de requête peut réorganiser ces expressions et clauses. Si une fonction UDF s’appuie sur la sémantique de court-circuitage pour la vérification null, Spark ne garantit pas que la vérification null s’exécute avant la fonction UDF. Par exemple:
spark.udf.register("strlen", (s: String) => s.length)
spark.sql("select s from test1 where s is not null and strlen(s) > 1") // no guarantee
Cette WHERE clause ne garantit pas que Spark appelle la strlen fonction UDF après qu’elle filtre les valeurs null.
Pour gérer la vérification null, Databricks recommande l’une des opérations suivantes :
- Rendre l’UDF elle-même capable de gérer les valeurs nulles et vérifier les valeurs nulles dans l’UDF
- Utilisez les expressions
IFouCASE WHENpour effectuer la vérification de valeurs nulles et appelez la fonction définie par l’utilisateur dans une branche conditionnelle
spark.udf.register("strlen_nullsafe", (s: String) => if (s != null) s.length else -1)
spark.sql("select s from test1 where s is not null and strlen_nullsafe(s) > 1") // ok
spark.sql("select s from test1 where if(s is not null, strlen(s), null) > 1") // ok
API des jeux de données typés
Remarque
Cette fonctionnalité est prise en charge sur les clusters avec catalogue Unity avec le mode d’accès standard dans Databricks Runtime 15.4 et versions ultérieures.
Utilisez des API de jeu de données typées pour exécuter des transformations telles que la carte, le filtre et les agrégations sur les jeux de données avec une fonction définie par l’utilisateur.
L’exemple suivant utilise l’API map() pour modifier un nombre dans une colonne de résultats en chaîne préfixée :
spark.range(3).map(f => s"row-$f").show()
Cet exemple utilise map(), mais le même modèle s’applique à d’autres API de jeu de données typées telles que filter(), , mapPartitions(), foreach(), foreachPartition(), reduce(), et flatMap().
Fonctionnalités UDF Scala et compatibilité Databricks Runtime
Les fonctionnalités suivantes nécessitent des versions minimales de Databricks Runtime sur des clusters compatibles avec le catalogue Unity en mode d’accès standard (partagé).
| Caractéristique | Version minimale de Databricks Runtime |
|---|---|
| Fonctions scalaires définies par l'utilisateur | Databricks Runtime 14.2 |
Dataset.map, , Dataset.mapPartitionsDataset.filter, , Dataset.reduceDataset.flatMap |
Databricks Runtime 15.4 |
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups |
Databricks Runtime 15.4 |
(Diffusion) foreachWriter Sink |
Databricks Runtime 15.4 |
(Diffusion) foreachBatch |
Databricks Runtime 16.1 |
(Diffusion) KeyValueGroupedDataset.flatMapGroupsWithState |
Version 16.2 de Databricks Runtime |
spark.udf.registerJavaFunction(Java UDF à partir d’un fichier JAR) |
Databricks Runtime 18 LTS |