Aller au contenu principal

UDF Scala et Java définies au niveau de la session

info

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

UDF Scala en ligne

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.

UDF Java à partir d'un JAR

Enregistrer une classe UDF précompilée depuis un JAR à l'aide de spark.udf.registerJavaFunction. Limité à la session. Pris en charge sur le compute Serverless.

UDF Scala ou Java gérée par Unity Catalog

Enregistrez une UDF dans Unity Catalog pour la gouvernance, la réutilisation et la visibilité. Pris en charge sur le compute serverless.

Approche

Description

UDF Scala en ligne

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.

UDF Java à partir d'un JAR

Enregistrer une classe UDF précompilée depuis un JAR à l'aide de spark.udf.registerJavaFunction. Limité à la session. Pris en charge sur le compute Serverless.

UDF Scala ou Java gérée par Unity Catalog

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.registerJavaFunction nécessite Databricks Runtime 18.3 ou une version ultérieure. Consultez Enregistrer une UDF Java à partir d'un JAR.
important

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.

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:

Scala
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 :

Scala
spark.range(1, 20).createOrReplaceTempView("test")
SQL
%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 :

Scala
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.

remarque

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.

Créer un nouveau projet Scala à l'aide de sbt:

Bash
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 :

Scala
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 :

Scala
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")

É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:

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.

Depuis le répertoire racine de votre projet, exécutez :

Bash
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.

É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 :

SQL
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:

  1. Dans votre Workspace Databricks, cliquez sur Icône de données. **Catalogue** pour ouvrir l’Explorateur de catalogues.
  2. Sélectionnez le catalogue, puis sélectionnez le schéma qui contient votre volume.
  3. Cliquez sur le nom du volume.
  4. Cliquez sur Upload vers ce volume et sélectionnez votre fichier JAR.
  5. Click upload .
  6. 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 :

Python
# 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 :

Output
+----------+
| 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 :

Scala
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 IF ou CASE WHEN pour effectuer la vérification de nullité et appeler la UDF dans une Branch conditionnelle.
Scala
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

remarque

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 :

Scala
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

Dataset.map, Dataset.mapPartitions, Dataset.filter, Dataset.reduce, Dataset.flatMap

Databricks Runtime 15.4

KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups

Databricks Runtime 15.4

(streaming) foreachWriter Sink

Databricks Runtime 15.4

(streaming) foreachBatch

Databricks Runtime 16.1

(streaming) KeyValueGroupedDataset.flatMapGroupsWithState

Databricks Runtime 16.2

spark.udf.registerJavaFunction (UDF Java à partir d'un JAR)

Databricks Runtime 18.3

Fonctionnalité

Version minimale de Databricks Runtime

UDF scalaires

Databricks Runtime 14,2

Dataset.map, Dataset.mapPartitions, Dataset.filter, Dataset.reduce, Dataset.flatMap

Databricks Runtime 15.4

KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups

Databricks Runtime 15.4

(streaming) foreachWriter Sink

Databricks Runtime 15.4

(streaming) foreachBatch

Databricks Runtime 16.1

(streaming) KeyValueGroupedDataset.flatMapGroupsWithState

Databricks Runtime 16.2

spark.udf.registerJavaFunction (UDF Java à partir d'un JAR)

Databricks Runtime 18.3