Fonctions utilisateur Scala et Java à portée de 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 à 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.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"))

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 :

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