UDF Scala et Java à l'échelle de la session

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 à l'échelle de la session sont créées dans un notebook ou un travail 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
UDF Scala inline Définissez une fonction UDF dans un notebook à l’aide d’une fonction Scala ou d’une fonction lambda. Étendue de session. Non pris en charge avec le calcul serverless.
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. Étendue de session. Non pris en charge avec le calcul 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é. Non pris en charge avec le calcul 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 UDF 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.registerJavaFunction né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.

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"))

Fichiers avec UDF

Important

Cette fonctionnalité est en version bêta. Les administrateurs d’espace de travail peuvent contrôler l’accès à cette fonctionnalité à partir de la page Aperçus . Consultez Gérer les préversions d’Azure Databricks.

Le type Scala pour un fichier est FileRef. Utilisez-le comme paramètre ou type de retour dans une UDF, soit comme type de premier niveau, soit imbriqué. Pour le type et ses règles d’imbriquement, voir FILE type.

Pour lire le contenu d’un fichier dans une UDF, appelez l’un des éléments suivants sur un FileRef:

  • asLocalFile(): Retourne un java.io.File que vous pouvez transmettre à n’importe quelle bibliothèque qui accepte un chemin.
  • open(): Renvoie un java.io.InputStream permettant de lire les octets du fichier. L’appelant ferme la porte.

Pour produire un nouveau FileRef dans une UDF, appelez l’une des méthodes statiques suivantes :

  • FileRef.create(uri): Crée une référence au fichier en uri.
  • FileRef.fromBytes(bytes, destinationPath, contentType): Charge bytes dans le chemin de volume destinationPath en tant que fichier externe et renvoie une référence.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Téléverse un fichier local vers le chemin de volume destinationPath en tant que fichier externe et renvoie une référence.

Retourner un FileRef à partir d’une UDF qui écrit dans une colonne FILE MANAGED n’est pas pris en charge.

Pour des exemples en Python, Scala et SQL, y compris le traitement d’images, la détection de types de fichiers et l’extraction de trames vidéo, voir Fichiers de processus avec UDF.

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 serverless à partir de Databricks Runtime 18 LTS ou version ultérieure. 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 pom.xml généré, dans les balises <project></project>, ajoutez un bloc <properties> et configurez le maven-shade-plugin pour générer un fichier JAR fat :

<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 fat

Empaqueter votre UDF compilé dans un fichier JAR fat.

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 :

  1. Dans votre espace de travail Azure Databricks, cliquez sur l’icône 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 Charger sur ce volume et sélectionnez votre fichier JAR.
  5. Cliquez sur Télécharger.
  6. 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 de « 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 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 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 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 :

  • Rendez l'UDF sensible aux valeurs null, et effectuez une vérification de valeur null dans l'UDF
  • Utilisez les expressions IF ou CASE WHEN pour 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 UDF scalaires 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