Aller au contenu principal

Fonctions définies par l'utilisateur (UDF) Scala et Java dans Unity Catalog

Cette page décrit comment créer des fonctions définies par l'utilisateur (UDF) Scala et Java, les enregistrer dans Unity Catalog et les partager entre les environnements de compute. Les UDF Unity Catalog vous permettent de réutiliser la logique JVM existante avec la gouvernance et les contrôles d'accès Unity Catalog.

Contrairement aux UDF Scala à étendue de session qui sont limitées à un seul Notebook ou cluster, les UDF enregistrées dans Unity Catalog sont :

  • Géré : Géré avec les autorisations et les contrôles d'accès d'Unity Catalog.
  • **Réutilisable** : Partagé entre les équipes, les Notebooks, les Jobs et les SQL Warehouse.
  • Détectable : Visible dans l'Explorateur de catalogues et les tables système.
  • **Isolé** : Exécuter dans des sandboxes avec un coût de start à froid unique par session. Les appels ultérieurs sont rapides.

Exigences

Votre Workspace doit être activé pour Unity Catalog. Les exigences supplémentaires suivantes s’appliquent.

Compute : Tous les types de compute sont pris en charge, y compris les Notebooks et les Jobs serverless, les SQL Warehouses et les Spark Declarative Pipelines sur Lakeflow. Le compute classique nécessite Databricks Runtime 18.2 ou version ultérieure. Sur le Serverless compute et les SQL Warehouse, la définition d’UDF doit spécifier la version d’environnement 4 ou supérieure dans le champ environment_version. Cette exigence s’applique à la définition d’UDF, et non au Notebook ou au Job appelant. Consultez les versions d’environnement Serverless.

Développement :

  • Scala : 2.13.16. Scala 2,12 n’est pas pris en charge.
  • JDK : 17.
  • **Packaging** : Un JAR fat contenant toutes les dépendances tierces utilisées par l'UDF.

Autorisations :

  • Créez une UDF : USAGE et CREATE FUNCTION sur le schéma, et USAGE sur le catalogue.
  • Exécuter une UDF : EXECUTE sur la fonction, et USAGE sur le schéma et le catalogue.
  • Accédez au fichier JAR : READ VOLUME sur le volume où le JAR est stocké.

Voir Gérer les privilèges dans Unity Catalog pour plus d'informations sur les autorisations Unity Catalog.

Construisez votre JAR UDF

Packagez votre code compilé en tant que JAR et uploadez-le vers un volume Unity Catalog avant d'enregistrer l'UDF. Choisissez une méthode de build :

Construire localement

Suivez ces étapes pour créer un fichier JAR volumineux à l'aide d'un environnement de développement local.

Configurez votre environnement

Installez les outils requis sur votre machine locale. Les commandes suivantes sont pour macOS. Pour les autres plateformes, installez JDK 17 et sbt (Scala) ou Maven (Java) à l'aide du gestionnaire de package de votre plateforme.

Installez JDK 17 et sbt :

Bash
brew install openjdk@17
brew install sbt

Vérifiez votre installation :

Bash
java -version   # Should show Java 17
sbt --version # Should show sbt version

Créez 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

Lorsque vous y êtes invité, saisissez un nom de projet (par exemple, my-udf-project).

Configurer build.sbt

Remplacez le contenu de votre fichier build.sbt par la configuration suivante :

Scala
scalaVersion := "2.13.16"

ThisBuild / organization := "com.example"

lazy val myUDF = (project in file("."))
.settings(
name := "my-udf"
)

Activer le plug-in sbt-assembly

Créez ou modifiez project/assembly.sbt et ajoutez :

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

Ce plug-in crée un JAR volumineux contenant toutes vos dépendances.

Écrivez votre UDF

Lorsque vous écrivez votre UDF, référez-vous aux types de données pris en charge et aux mappages de langage pour voir comment les types Scala et Java correspondent aux types SQL.

Votre gestionnaire d'UDF doit répondre aux exigences suivantes :

  • Scala : définissez le gestionnaire comme une méthode sur un object (pas un class). La valeur HANDLER se résout en une méthode sur un object Scala.
  • Java : Définissez le gestionnaire comme une méthode public static.
  • Signature : Les types, l'ordre et le type de retour des paramètres de la méthode doivent correspondre à la liste des arguments et au RETURNS type dans votre CREATE FUNCTION instruction.
  • Scalaire uniquement : Le gestionnaire doit renvoyer une seule valeur scalaire. Les types de retour de table ne sont pas pris en charge.
  • Autonome : Le gestionnaire doit opérer uniquement sur ses arguments d'entrée. Il ne peut pas utiliser les Spark API ou dépendre des packages Spark principaux. Consultez les Limitations.
remarque

Pour Scala, un gestionnaire avec un type de paramètre primitif (tel que Int) est ignoré et renvoie NULL lorsque n’importe quel argument d’entrée est SQL NULL. Pour recevoir et gérer les valeurs NULL, enveloppez le paramètre dans Option, tel que Option[Int].

Créez un objet Scala dans src/main/scala/com/example/MyUDF.scala et définissez votre fonction UDF.

Exemple basique

Scala
package com.example

object MyUDF {
def addOne(x: Int): Int = x + 1
}

Exemple avec dépendance externe

Pour utiliser des bibliothèques externes, ajoutez-les à votre fichier build.sbt :

Scala
scalaVersion := "2.13.16"

ThisBuild / organization := "com.example"

lazy val myUDF = (project in file("."))
.settings(
name := "currency-udf",
libraryDependencies ++= Seq(
"org.apache.commons" % "commons-lang3" % "3.12.0"
)
)

Utilisez ensuite la dépendance dans votre UDF :

Scala
package com.example

import org.apache.commons.lang3.StringUtils

object CurrencyUDF {
private val rates: Map[String, Double] = Map(
"USD" -> 1.0,
"EUR" -> 1.1,
"GBP" -> 1.3,
"JPY" -> 0.007
)

def convertToUSD(price: Double, currency: String): Double = {
require(currency != null, "Currency must not be null")

val normalizedCurrency = StringUtils.upperCase(currency)

rates.get(normalizedCurrency) match {
case Some(rate) => price * rate
case None => throw new IllegalArgumentException(s"Unsupported currency: $currency")
}
}
}

Testez votre UDF avec des tests unitaires avant de la déployer. Voir Tester les UDF localement.

remarque

Votre UDF s'exécute dans un sandbox isolé sans session Spark active, de sorte qu'il ne peut pas utiliser les Spark API depuis le corps de la fonction. Par exemple, vous ne pouvez pas créer ou opérer sur des DataFrames ou des Datasets, exécuter spark.sql(...), ou accéder à SparkSession ou SparkContext. L'UDF doit être une logique autonome basée sur ses arguments d'entrée. Il ne peut pas non plus dépendre des packages Spark core.

Créez votre JAR fat

Construisez votre projet pour créer un JAR autonome contenant toutes les dépendances.

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.

upload votre JAR vers un volume Unity Catalog

Si vous n'avez pas encore de volume Unity Catalog, créez-en un :

SQL
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';

Si d'autres utilisateurs doivent exécuter l'UDF, accordez-leur READ VOLUME sur le volume :

SQL
GRANT READ VOLUME ON VOLUME my_catalog.my_schema.udf_jars TO `user@example.com`;

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.
  7. Cliquez sur Copier le chemin pour copier le chemin du volume dans votre presse-papiers. Par exemple, /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar (Scala) ou /Volumes/my_catalog/my_schema/udf_jars/my-udf-1.0-SNAPSHOT.jar (Java). Vous avez besoin de ce chemin lorsque vous enregistrez l’UDF.

Créer dans un Notebook

Vous pouvez compiler une UDF, l'empaqueter en tant que JAR et l'upload vers un volume Unity Catalog directement depuis un Notebook Databricks. Cette approche fonctionne pour les petites UDF sans dépendances. Pour les UDF avec des bibliothèques tierces, utilisez Build locally.

La cellule Python suivante écrit une UDF Java qui nettoie une chaîne (supprime les espaces blancs, réduit les espaces répétés et met en minuscules), la compile avec JDK 17, l'empaquette en JAR et la copie dans un volume Unity Catalog. Mettez à jour le volume_path pour qu’il pointe vers un volume existant sur lequel vous avez la permission WRITE VOLUME.

Python
import os
import subprocess
import shutil

build_dir = "/tmp/udf_build"
package_dir = f"{build_dir}/src/com/databricks/udf"
classes_dir = f"{build_dir}/classes"
os.makedirs(package_dir, exist_ok=True)
os.makedirs(classes_dir, exist_ok=True)

# The UDF handler: a public static method on a plain Java class.
# The doubled backslashes produce a single backslash in the Java source (\\s+).
udf_code = """package com.databricks.udf;
public class StringCleanUDF {
public static String clean(String input) {
if (input == null) return null;
return input.trim().replaceAll("\\\\s+", " ").toLowerCase();
}
}
"""
with open(f"{package_dir}/StringCleanUDF.java", "w") as f:
f.write(udf_code)

# Compile with JDK 17 to match Environment Version 4.
subprocess.run(
["javac", "--release", "17", "-d", classes_dir, f"{package_dir}/StringCleanUDF.java"],
check=True,
)

# Package the compiled class into a JAR.
jar_path = f"{build_dir}/string_clean_udf.jar"
subprocess.run(["jar", "cf", jar_path, "-C", classes_dir, "."], check=True)

# Copy the JAR to a Unity Catalog volume.
volume_path = "/Volumes/my_catalog/my_schema/udf_jars/string_clean_udf.jar"
os.makedirs(os.path.dirname(volume_path), exist_ok=True)
shutil.copy2(jar_path, volume_path)

print(f"JAR uploaded to: {volume_path}")

Une fois que le JAR est dans le volume, enregistrez l'UDF. Utilisez LANGUAGE JAVA et définissez le HANDLER sur la méthode entièrement qualifiée, telle que com.databricks.udf.StringCleanUDF.clean.

Enregistrez votre UDF dans Unity Catalog

Après avoir créé et upload votre JAR, utilisez l'instruction CREATE FUNCTION pour enregistrer votre UDF dans Unity Catalog.

SQL
CREATE OR REPLACE FUNCTION my_catalog.my_schema.add_one(x INT)
RETURNS INT
LANGUAGE SCALA
DETERMINISTIC
ENVIRONMENT (
java_dependencies = '["/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar"]',
environment_version = '4'
)
HANDLER 'com.example.MyUDF.addOne';

L'instruction CREATE FUNCTION utilise les paramètres suivants :

  • LANGUAGE: Le langage de l'UDF.

  • HANDLER: chemin d’accès entièrement qualifié à la méthode, au format 'package.Object.method' (Scala) ou 'package.ClassName.method' (Java).

  • DETERMINISTIC: Déclare que la fonction retourne toujours la même sortie pour la même entrée, permettant l'optimisation des requêtes.

remarque

Supprimez DETERMINISTIC si votre fonction appelle des APIs externes ou présente tout autre comportement non déterministe.

  • ENVIRONMENT: Définit l’environnement d’exécution pour l’UDF.

    • java_dependencies: Un tableau JSON de chemins de fichiers JAR dans vos volumes Unity Catalog. Il s'agit du chemin de fichier que vous avez copié à l'étape précédente. Utilisez des guillemets simples autour du tableau et des guillemets doubles autour des chemins.
    • environment_version: Doit être '4' ou supérieur pour les UDF Scala et Java. La version 4 de l'environnement spécifie Scala 2.13.16 et JDK 17. Consultez les versions de l'environnement Serverless.

Appelez votre UDF dans SQL et les notebooks

Après l'enregistrement, vous pouvez appeler l'UDF dans les query SQL, les notebooks et les vues :

SQL
-- Simple select
SELECT my_catalog.my_schema.add_one(5) AS result;

-- With table data
SELECT
id,
price,
currency,
my_catalog.my_schema.convert_to_usd(price, currency) AS price_usd
FROM my_catalog.my_schema.transactions;

-- Filtering
SELECT *
FROM my_catalog.my_schema.products
WHERE my_catalog.my_schema.convert_to_usd(price, currency) > 100;

-- Aggregation
SELECT
category,
SUM(my_catalog.my_schema.convert_to_usd(price, currency)) AS total_usd
FROM my_catalog.my_schema.sales
GROUP BY category;

Gouvernance et partage

Utilisez les autorisations de Unity Catalog pour contrôler qui peut exécuter votre UDF et pour la rendre visible au sein de votre organisation.

Accorder des autorisations

Utilisez Catalog Explorer ou SQL pour accorder les autorisations nécessaires aux autres utilisateurs afin qu’ils puissent exécuter vos UDF.

  1. Dans la barre latérale, cliquez sur Icône de données. Catalogue .
  2. Sélectionnez le catalogue, puis sélectionnez le schéma qui contient votre fonction.
  3. Cliquez sur le nom de la fonction.
  4. Dans l'onglet Autorisations , cliquez sur Accorder .
  5. Sélectionnez les principaux auxquels vous souhaitez accorder l'accès, puis sélectionnez l'autorisation EXECUTE.
  6. Cliquez sur Confirmer .

Révoquer les autorisations

Utilisez Catalog Explorer ou SQL pour révoquer les autorisations des autres utilisateurs.

  1. Dans la barre latérale, cliquez sur Icône de données. Catalogue .
  2. Sélectionnez le catalogue, puis sélectionnez le schéma qui contient votre fonction.
  3. Cliquez sur le nom de la fonction.
  4. Dans l'onglet **tab**, cochez la case en regard du principal auquel vous souhaitez révoquer l'accès. Cliquez sur Révoquer .
  5. Dans la notification, cliquez sur Révoquer .

Découvrir les UDF

Pour trouver les UDF gérées dans Unity Catalog, query la table information_schema.routines, en remplaçant les valeurs my_catalog et my_schema :

SQL
SELECT
routine_catalog,
routine_schema,
routine_name,
routine_definition,
created
FROM system.information_schema.routines
WHERE routine_catalog = 'my_catalog'
AND routine_schema = 'my_schema';

Mettez à jour votre UDF

Pour mettre à jour une UDF Unity Catalog existante avec un nouveau code :

  1. Apportez des modifications à votre code localement.

  2. Reconstruisez le JAR avec un nouveau numéro de version.

    • Scala : sbt clean assembly (par exemple, my-udf-assembly-0.2.0-SNAPSHOT.jar)
    • Java : mvn clean package (par exemple, my-udf-2.0-SNAPSHOT.jar)
  3. Upload le nouveau JAR dans le volume Unity Catalog.

  4. Utilisez CREATE OR REPLACE FUNCTION avec le même nom de fonction pour mettre à jour l'UDF. Vérifiez que vous référencez le dernier JAR dans votre java_dependencies.

Databricks utilise le nouveau code lors de la prochaine invocation. Vous n'avez pas besoin de redémarrer votre cluster.

Optimisation des performances

Latence de start à froid

Le premier appel UDF dans une session initialise le sandbox isolé, ce qui ajoute de la latence. Les appels ultérieurs dans la même session sont plus rapides. Tenez-en compte lors de l'évaluation ou de la conception de charges de travail sensibles à la latence.

Mise en cache des calculs coûteux

Si votre UDF effectue une initialisation ou un compute coûteux, mettez le résultat en cache pour ne le calculer qu'une seule fois.

Utilisez un champ val dans l'objet Scala pour mettre en cache le résultat :

Scala
package example

object CachedUDF {
// Computed once and cached
val expensiveData: Map[String, Double] = {
// Load data from somewhere expensive
Map("key1" -> 1.0, "key2" -> 2.0)
}

def lookup(key: String): Double = {
expensiveData.getOrElse(key, 0.0)
}
}

Utilisez DÉTERMINISTE lorsque cela est approprié

Marquez votre UDF comme DETERMINISTIC si elle produit toujours le même résultat pour la même entrée. Ceci permet à l'optimiseur de query de mettre en cache les résultats et d'améliorer les performances.

Limitations

  • Seules les UDF scalaires sont prises en charge. Les fonctions d'agrégation définies par l'utilisateur (UDAFs) et les fonctions de table définies par l'utilisateur (UDTFs) ne sont pas prises en charge.
  • Les UDF s'exécutent dans un sandbox isolé sans session Spark active. Les Spark APIs (SparkSession, SparkContext, spark.sql(...), DataFrame et opérations sur les datasets) ne sont pas disponibles.
  • Les UDF ne peuvent pas dépendre des packages Spark principaux.
  • Les UDF n'ont pas accès aux fichiers de workspace ou aux volumes Unity Catalog lors de l'exécution.

Bonnes pratiques

Databricks recommande les pratiques suivantes :

  • Gérez la version de vos fichiers JAR. Par exemple, my-udf-0.1.0.jar, my-udf-0.2.0.jar.
  • Valider les mappages de types SQL avant le déploiement. Voir les mappages de langues.
  • Accordez les autorisations READ VOLUME et EXECUTE uniquement aux utilisateurs qui doivent exécuter l'UDF. Utilisez la propriété de groupe pour les UDF partagées entre les équipes.

Test des UDF localement

Testez votre UDF avec des tests unitaires avant de déployer en production.

Pour tester src/main/scala/example/MyUDF.scala, créez un fichier de test dans src/test/scala/example/MyUDFTest.scala:

Scala
package example

import org.scalatest.funsuite.AnyFunSuite

class MyUDFTest extends AnyFunSuite {
test("addOne should add 1 to input") {
assert(MyUDF.addOne(5) == 6)
}

test("addOne should handle negative numbers") {
assert(MyUDF.addOne(-1) == 0)
}
}

Ajoutez la dépendance de test à build.sbt:

Scala
libraryDependencies += "org.scalatest" %% "scalatest" % "3.2.15" % Test

Pour exécuter les tests :

Bash
sbt test

Ressources supplémentaires