UDF Scala et Java définies au niveau de la session
Les UDF Scala et Java peuvent être enregistrées dans Unity Catalog pour la gouvernance, la réutilisation et la découvrabilité. Consultez les fonctions définies par l’utilisateur (UDF) Scala et Java dans Unity Catalog.
Cette page décrit comment créer des UDF Scala et Java limitées à la session dans Databricks. Les UDF limitées à la session sont définies dans un Notebook ou un Job et s'appliquent uniquement à la SparkSession actuelle. Pour la référence du langage SQL, consultez Fonctions scalaires externes définies par l'utilisateur (UDF).
Choisissez votre approche
Vous pouvez définir une UDF Scala ou Java des manières suivantes. Pour comparer tous les types d'UDF entre les langages, la gouvernance et le compute, consultez UDF régies par Unity Catalog et UDF à portée de session.
Approche | Description |
|---|---|
Définir une UDF dans un notebook à l'aide d'une fonction Scala ou lambda. Limité à la session. Non pris en charge sur le compute serverless. | |
Enregistrer une classe UDF précompilée depuis un JAR à l'aide de | |
Enregistrez une UDF dans Unity Catalog pour la gouvernance, la réutilisation et la visibilité. Pris en charge sur le compute serverless. |
Exigences
- Les UDF Scala sur les computes compatibles Unity Catalog avec le mode d'accès standard nécessitent Databricks Runtime 14,2 ou une version ultérieure.
- La prise en charge des instances ARM pour les UDF Scala sur les clusters compatibles avec Unity Catalog nécessite Databricks Runtime 15.2 ou version ultérieure.
- L'enregistrement d'une UDF Java à partir d'un JAR avec
spark.udf.registerJavaFunctionnécessite Databricks Runtime 18.3 ou une version ultérieure. Consultez Enregistrer une UDF Java à partir d'un JAR.
Créez votre JAR avec les mêmes versions de Scala et d'Apache Spark que le compute qui l'exécute. Une incompatibilité peut entraîner l'échec de l'UDF lors de l'enregistrement ou de l'appel.
- Compute classique : veuillez faire correspondre les versions de Scala et Spark de votre Databricks Runtime. Veuillez consulter la section Environnement système des versions et de la compatibilité des notes de version de Databricks Runtime pour votre version. Par exemple, Databricks Runtime 18,3 utilise Scala 2.13.16 et Apache Spark 4,0.
- Compute serverless : faites correspondre la version Scala à la version de votre environnement. Consultez Versions d’environnement Serverless.
Marquez la dépendance Apache Spark comme provided afin qu'elle ne soit pas intégrée à votre JAR. Incluez uniquement les dépendances tierces utilisées par votre UDF.
Enregistrer une fonction en tant qu'UDF
Enregistrez une fonction Scala en tant qu'UDF à l'aide de spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Appelez l'UDF dans Spark SQL
Créez une vue temporaire, puis appelez l'UDF dans une requête SQL :
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Utiliser les UDF avec les DataFrames
Vous pouvez également appeler une 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"))
Enregistrer une UDF Java à partir d'un JAR
Packagez une UDF en tant que JAR, ajoutez-le à votre session avec spark.addArtifact, et enregistrez la classe UDF avec spark.udf.registerJavaFunction.
Pris en charge en mode d'accès standard et compute serverless dans Databricks Runtime 18,3 ou supérieure. La fonction enregistrée est limitée à la session et n'est pas enregistrée dans Unity Catalog.
Les étapes suivantes décrivent la création d'un projet, l'écriture d'une classe UDF, la création d'un fichier JAR volumineux et son enregistrement.
Étape 1 : Créer votre projet
Configurez un projet en Scala ou Java.
- Scala
- Java
Créer un nouveau projet Scala à l'aide de sbt:
sbt new scala/scala-seed.g8
Remplacez le contenu de votre fichier build.sbt par ce qui suit. Définissez scalaVersion et la version spark-sql pour qu'elles correspondent à votre compute :
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 plugin sbt-assembly pour créer un fat JAR. Créez ou modifiez project/assembly.sbt et ajoutez :
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Créez un nouveau 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 pom.xml généré, dans les balises <project></project>, ajoutez un bloc <properties> et configurez le maven-shade-plugin pour créer un JAR volumineux :
<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 : Écrivez votre classe UDF
Votre classe UDF doit implémenter l'une des org.apache.spark.sql.api.java.UDF interfaces (UDF1 à UDF22), où le nombre indique le nombre d'arguments d'entrée que l'UDF prend. Implémentez la méthode call() 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 premier niveau (ou imbriquée static) avec un constructeur public sans argument. Un Scala class ou object ne satisfait pas à cette exigence et échoue au moment de l'appel. Vous pouvez créer le JAR avec sbt, mais la classe UDF elle-même doit être écrite en Java.
Créer 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 : créer votre JAR fat
Empaquetez votre UDF compilée dans un fat JAR.
- Scala
- Java
Depuis le répertoire racine de votre projet, exécutez :
sbt clean assembly
Le fat JAR est créé dans target/scala-2.13/ avec un nom tel que my-udf-assembly-0.1.0-SNAPSHOT.jar.
Depuis le répertoire racine de votre projet, exécutez :
mvn clean package
Le fat JAR est créé dans target/ avec un nom tel que my-udf-1.0-SNAPSHOT.jar.
Étape 4 : Upload de votre JAR vers un volume Unity Catalog
upload le JAR vers un volume Unity Catalog afin que votre compute puisse y accéder. Si vous n'en avez pas déjà un, créez-en un :
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Upload votre fichier JAR dans le volume à l’aide de Catalog Explorer:
- Dans votre Workspace 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 Upload vers ce volume et sélectionnez votre fichier JAR.
- Click upload .
- Une fois l'upload terminé, cliquez sur le nom de votre fichier JAR, puis cliquez sur **Copier le chemin** 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 : Enregistrez et appelez l'UDF
Ajoutez le JAR à votre session en utilisant son chemin de volume, enregistrez la classe UDF et appelez-la depuis 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()
Sur le compute serverless et en mode d’accès standard, vous devez passer un type de retour explicite. 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 (UDAFs) ne sont pas prises en charge avec registerJavaFunction.
La query renvoie la sortie de l'UDF, confirmant que la fonction est enregistrée et appelable :
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Ordre d'évaluation et vérification des valeurs nulles
Spark SQL (y compris SQL et les APIs 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 de court-circuitage de gauche à droite.
Ne vous fiez pas aux effets secondaires ou à l'ordre d'évaluation des expressions booléennes, ni à l'ordre des clauses WHERE et HAVING. L'optimiseur de query peut réorganiser ces expressions et clauses. Si une UDF repose sur la sémantique d'évaluation en court-circuit pour la vérification de nullité, Spark ne garantit pas que la vérification de nullité s'exécute avant l'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 clause WHERE ne garantit pas que Spark invoque l'UDF strlen après avoir filtré les valeurs nulles.
Pour gérer la vérification des valeurs nulles, Databricks recommande l'une des méthodes suivantes :
- Rendez l'UDF elle-même sensible aux valeurs nulles et effectuez une vérification des valeurs nulles au sein de l'UDF
- Utilisez les expressions
IFouCASE WHENpour effectuer la vérification de nullité et appeler la UDF dans une Branch 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
APIs de Dataset typées
Cette fonctionnalité est prise en charge sur les clusters compatibles Unity Catalog avec le mode d'accès standard dans Databricks Runtime 15,4 et versions ultérieures.
Utilisez les APIs Dataset typées pour exécuter des transformations telles que map, filter et des agrégations sur des Datasets 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 une 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 APIs de dataset typées telles que filter(), mapPartitions(), foreach(), foreachPartition(), reduce() et flatMap().
Fonctionnalités des UDF Scala et compatibilité avec Databricks Runtime
Les fonctionnalités suivantes requièrent des versions minimales de Databricks Runtime sur des clusters avec Unity Catalog activé, en mode d'accès standard (partagé).
Fonctionnalité | Version minimale de Databricks Runtime |
|---|---|
UDF scalaires | Databricks Runtime 14,2 |
| Databricks Runtime 15.4 |
| Databricks Runtime 15.4 |
(streaming) | Databricks Runtime 15.4 |
(streaming) | Databricks Runtime 16.1 |
(streaming) | Databricks Runtime 16.2 |
| Databricks Runtime 18.3 |