Scala och Java användardefinierade funktioner (UDF: er) i Unity Catalog

Den här sidan beskriver hur du skapar Scala och Java användardefinierade funktioner (UDF: er), registrerar dem i Unity Catalog och delar dem i beräkningsmiljöer. Med UDF:er för Unity Catalog kan du återanvända befintlig JVM-logik med styrnings- och åtkomstkontroller för Unity Catalog.

Till skillnad från Scala-UDF:er med sessionsomfattning, som är begränsade till en enda anteckningsbok eller ett kluster, är registrerade UDF:er i Unity Catalog:

  • Styrd: Hanteras med behörigheter och åtkomstkontroller för Unity-katalogen.
  • Återanvändbar: Delas mellan team, notebook-filer, jobb och SQL-lager.
  • Identifieringsbar: Synlig i Katalogutforskaren och systemtabeller.
  • Isolerad: Kör i sandbox-miljö med en engångskostnad för kallstart per session. Efterföljande anrop är snabba.

Requirements

Arbetsytan måste vara aktiverad för Unity Catalog. Följande ytterligare krav gäller.

Beräkning: Alla beräkningstyper stöds, inklusive serverlösa anteckningsböcker och jobb, SQL warehouses och Spark Declarative Pipelines i Lakeflow. Klassisk beräkning kräver Databricks Runtime 18.2 eller senare. På serverlösa beräknings- och SQL-lager måste UDF-definitionen ange miljöversion 4 eller senare i fältet environment_version . Det här kravet gäller för UDF-definitionen, inte för den anropande notebook-filen eller jobbet. Se serverlösa miljöversioner.

Utveckling:

  • Scala: 2.13.16. Scala 2.12 stöds inte.
  • JDK: 17.
  • Paketering: En fet JAR som innehåller alla beroenden från tredje part som används av UDF.

Behörigheter:

  • Skapa en UDF: USAGE och CREATE FUNCTION i schemat och USAGE i katalogen.
  • Kör en UDF: EXECUTE på funktionen och USAGE i schemat och katalogen.
  • Få åtkomst till JAR-filen: READ VOLUME på volymen där JAR-filen lagras.

Mer information om Behörigheter för Unity-katalogen finns i Hantera privilegier i Unity Catalog .

Skapa din UDF-JAR

Paketera din kompilerade kod som en JAR och ladda upp den till en Unity Catalog-volym innan du registrerar UDF. Välj en byggmetod:

Skapa lokalt

Följ de här stegen för att skapa en fet JAR med hjälp av en lokal utvecklingsmiljö.

Konfigurera din miljö

Installera de verktyg som krävs på den lokala datorn. Följande kommandon gäller för macOS. För andra plattformar installerar du JDK 17 och sbt (Scala) eller Maven (Java) med hjälp av plattformens pakethanterare.

Scala

Installera JDK 17 och sbt:

brew install openjdk@17
brew install sbt

Kontrollera installationen:

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

Java

Installera JDK 17 och Maven:

brew install openjdk@17
brew install maven

Kontrollera installationen:

java -version   # Should show Java 17
mvn --version   # Should show Maven version

Skapa projektet

Konfigurera ett projekt i Scala eller Java.

Scala

Skapa ett nytt Scala-projekt med :sbt

sbt new scala/scala-seed.g8

När du uppmanas till det anger du ett projektnamn (till exempel my-udf-project).

Konfigurera build.sbt

Ersätt innehållet i build.sbt filen med följande konfiguration:

scalaVersion := "2.13.16"

ThisBuild / organization := "com.example"

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

Aktivera plugin-programmet sbt-assembly

Skapa eller redigera project/assembly.sbt och lägg till:

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

Det här plugin-programmet skapar en fet JAR som innehåller alla dina beroenden.

Java

Skapa ett nytt Maven-projekt med hjälp av snabbstartsarketypen:

mvn archetype:generate \
  -DgroupId=com.example \
  -DartifactId=my-udf \
  -DarchetypeArtifactId=maven-archetype-quickstart \
  -DinteractiveMode=false

Det här kommandot skapar maven-standardprojektstrukturen med src/main/java och src/test/java kataloger.

Konfigurera pom.xml

I den genererade pom.xml filen i taggarna <project></project> lägger du till ett <properties> block med följande konfiguration:

<properties>
  <maven.compiler.source>17</maven.compiler.source>
  <maven.compiler.target>17</maven.compiler.target>
  <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

Lägg även till ett <build>-block med följande konfiguration i <project></project>-taggarna:

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

maven-shade-plugin Skapar en fet JAR som innehåller alla dina beroenden.

Skriv din UDF

När du skriver UDF läser du Datatyper för datatyper som stöds och Språkmappningar för att se hur Scala och Java typer mappas till SQL-typer.

UDF-hanteraren måste uppfylla följande krav:

  • Scala: Definiera hanteraren som en metod på en object (inte en class). Värdet HANDLER motsvarar en metod på en Scala object.
  • Java: Definiera hanteraren som en public static metod.
  • Signatur: Metodens parametertyper, ordning och returtyp måste motsvara argumentlistan och RETURNS-typen i din CREATE FUNCTION-sats.
  • Endast skalär: Hanteraren måste returnera ett enda skalärt värde. Tabellreturtyper stöds inte.
  • Fristående: Hanteraren får endast fungera på indataargumenten. Den kan inte använda Spark-API:er eller vara beroende av Spark-kärnpaket. Se Begränsningar.

Note

I Scala utelämnas en hanterare med en primitiv parametertyp (till exempel Int) och returnerar NULL när något indataargument är SQL NULL. Om du vill ta emot och hantera NULL värden omsluter du parametern i Option, till exempel Option[Int].

Scala

Skapa ett Scala-objekt i src/main/scala/com/example/MyUDF.scala och definiera din UDF-funktion.

Grundläggande exempel

package com.example

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

Exempel med externt beroende

Om du vill använda externa bibliotek lägger du till dem i build.sbt filen:

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

Använd sedan beroendet i din UDF:

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

Testa din UDF med enhetstester innan du distribuerar. Se Testa UDF:er lokalt.

Java

Skapa en Java-klass i src/main/java/com/example/MyUDF.java och definiera din UDF som en offentlig statisk metod.

Grundläggande exempel

package com.example;

public class MyUDF {
    public static int addOne(int x) {
        return x + 1;
    }
}

Exempel med externt beroende

Om du vill använda externa bibliotek lägger du till dem i <dependencies> avsnittet i pom.xml filen:

<dependencies>
    <dependency>
        <groupId>org.apache.commons</groupId>
        <artifactId>commons-lang3</artifactId>
        <version>3.12.0</version>
    </dependency>
</dependencies>

Använd sedan beroendet i din UDF:

package com.example;

import org.apache.commons.lang3.StringUtils;
import java.util.Map;
import java.util.HashMap;

public class CurrencyUDF {
    private static final Map<String, Double> rates = new HashMap<>();

    static {
        rates.put("USD", 1.0);
        rates.put("EUR", 1.1);
        rates.put("GBP", 1.3);
        rates.put("JPY", 0.007);
    }

    public static double convertToUSD(double price, String currency) {
        if (currency == null) {
            throw new IllegalArgumentException("Currency must not be null");
        }

        String normalizedCurrency = StringUtils.upperCase(currency);

        if (!rates.containsKey(normalizedCurrency)) {
            throw new IllegalArgumentException("Unsupported currency: " + currency);
        }

        return price * rates.get(normalizedCurrency);
    }
}

Testa din UDF med enhetstester innan du distribuerar. Se Testa UDF:er lokalt.

Note

Din UDF körs i en isolerad sandbox-miljö utan någon aktiv Spark-session, så den kan inte använda Spark-API:er inifrån funktionstexten. Du kan till exempel inte skapa eller arbeta med DataFrames eller Datauppsättningar, köra spark.sql(...)eller komma åt SparkSession eller SparkContext. UDF måste vara fristående logik över sina indataargument. Det kan inte heller vara beroende av Spark-kärnpaket.

Skapa din fat JAR

Skapa projektet för att skapa en fet JAR som innehåller alla beroenden.

Scala

Kör från projektets rotkatalog:

sbt clean assembly

Den fat JAR-filen skapas i target/scala-2.13/ med ett namn som my-udf-assembly-0.1.0-SNAPSHOT.jar.

Java

Kör följande från projektets rotkatalog:

mvn clean package

Den feta JAR-filen skapas i target/ med ett namn som my-udf-1.0-SNAPSHOT.jar.

Ladda upp din JAR-fil till en Unity Catalog-volym

Om du inte redan har en Unity Catalog-volym skapar du en:

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

Om andra användare behöver köra UDF beviljar du dem READ VOLUME på volymen:

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

Ladda upp din JAR-fil till volymen med Catalog Explorer:

  1. På din Azure Databricks-arbetsyta klickar du på dataikonen.Katalog för att öppna Katalogutforskaren.
  2. Välj katalogen och välj sedan det schema som innehåller volymen.
  3. Klicka på volymnamnet.
  4. Klicka på Ladda upp till den här volymen och välj din JAR-fil.
  5. Klicka på Överför.
  6. När uppladdningen är klar klickar du på namnet på JAR-filen.
  7. Klicka Kopiera sökväg för att kopiera volymsökvägen till urklipp. Till exempel /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar (Scala) eller /Volumes/my_catalog/my_schema/udf_jars/my-udf-1.0-SNAPSHOT.jar (Java). Du behöver den här sökvägen när du registrerar UDF.

Skapa i ett anteckningsblock

Du kan kompilera en UDF, paketera den som en JAR och ladda upp den till en Unity Catalog-volym direkt från en Azure Databricks notebook-fil. Den här metoden fungerar för små, beroendefria UDF:er. För UDF:er med bibliotek från tredje part använder du Skapa lokalt.

Följande Python cell skriver en Java UDF som rensar en sträng (trimmar blanksteg, minimerar upprepade blanksteg och gemener), kompilerar den med JDK 17, paketerar den som en JAR och kopierar den till en Unity Catalog-volym. Uppdatera så att den volume_path pekar på en befintlig volym som du har WRITE VOLUME behörighet till.

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

Efter att JAR-filen finns i volymen registrerar du UDF:n. Använd LANGUAGE JAVA och ange HANDLER till den fullständigt kvalificerade metoden, till exempel com.databricks.udf.StringCleanUDF.clean.

Registrera din UDF i Unity Catalog

När du har skapat och laddat upp din JAR använder du -instruktionen CREATE FUNCTION för att registrera din UDF i Unity Catalog.

Scala

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

Java

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

Instruktionen CREATE FUNCTION använder följande parametrar:

  • LANGUAGE: Språket i UDF.

  • HANDLER: Fullständigt kvalificerad sökväg till metoden, i formatet 'package.Object.method' (Scala) eller 'package.ClassName.method' (Java).

  • DETERMINISTIC: Deklarerar att funktionen alltid returnerar samma utdata för samma indata, vilket möjliggör frågeoptimering.

    Note

    Ta bort DETERMINISTIC om funktionen anropar externa API:er eller har något annat icke-deterministiskt beteende.

  • ENVIRONMENT: Definierar körningsmiljön för UDF.

    • java_dependencies: En JSON-matris med JAR-filsökvägar i dina Unity Catalog-volymer. Det här är den filsökväg som du kopierade i föregående steg. Använd enkla citattecken runt matrisen och dubbla citattecken runt sökvägar.
    • environment_version: Måste vara '4' eller högre för Scala och Java UDF:er. Miljöversion 4 anger Scala 2.13.16 och JDK 17. Se serverlösa miljöversioner.

Anropa din UDF i SQL och notebook-filer

Efter registreringen kan du anropa UDF i SQL-frågor, notebook-filer och vyer:

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

Styrning och delning

Använd Behörigheter för Unity Catalog för att styra vem som kan köra din UDF och för att göra den identifierbar i hela organisationen.

Bevilja behörigheter

Använd Catalog Explorer eller SQL för att bevilja de behörigheter som krävs för att andra användare ska kunna köra dina UDF:er.

Katalogutforskaren

  1. I sidofältet klickar du på dataikonen.Katalog.
  2. Välj katalogen och välj sedan det schema som innehåller din funktion.
  3. Klicka på funktionsnamnet.
  4. På fliken Behörigheter klickar du på Bevilja.
  5. Välj de huvudkonton som du vill bevilja åtkomst till och välj behörigheten EXECUTE .
  6. Klicka på Bekräfta.

SQL

Kör följande kommando i en notebook-fil eller Databricks SQL-redigeraren för att bevilja EXECUTE behörigheter till en användare eller grupp.

-- Grant to a specific user
GRANT EXECUTE ON FUNCTION my_catalog.my_schema.add_one TO `user@example.com`;

-- Grant to a group
GRANT EXECUTE ON FUNCTION my_catalog.my_schema.add_one TO `data-engineers`;

Återkalla behörigheter

Använd Catalog Explorer eller SQL för att återkalla behörigheter från andra användare.

Katalogutforskaren

  1. I sidofältet klickar du på dataikonen.Katalog.
  2. Välj katalogen och välj sedan det schema som innehåller din funktion.
  3. Klicka på funktionsnamnet.
  4. På fliken Behörigheter markerar du kryssrutan bredvid det huvudnamn som du vill återkalla åtkomsten från. Klicka på Återkalla.
  5. I meddelandet klickar du på Återkalla.

SQL

Kör följande kommando i en notebook-fil eller Databricks SQL-redigeraren för att återkalla EXECUTE behörigheter från en användare eller grupp.

-- Revoke from specific user
REVOKE EXECUTE ON FUNCTION my_catalog.my_schema.add_one FROM `user@example.com`;

-- Revoke from a group
REVOKE EXECUTE ON FUNCTION my_catalog.my_schema.add_one FROM `data-engineers`;

Identifiera UDF:er

Om du vill hitta UDF:er som hanteras i Unity Catalog, fråga tabellen information_schema.routines och ersätt värdena my_catalog och my_schema:

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

Uppdatera din UDF

Så här uppdaterar du en befintlig UDF för Unity Catalog med ny kod:

  1. Gör ändringar i koden lokalt.
  2. Återskapa JAR-filen med ett nytt versionsnummer.
    • Scala: sbt clean assembly (till exempel my-udf-assembly-0.2.0-SNAPSHOT.jar)
    • Java: mvn clean package (till exempel my-udf-2.0-SNAPSHOT.jar)
  3. Ladda upp den nya JAR-filen till Unity Catalog-volymen.
  4. Använd CREATE OR REPLACE FUNCTION med samma funktionsnamn för att uppdatera UDF. Kontrollera att du refererar till den senaste JAR-filen i java_dependencies.

Azure Databricks använder den nya koden vid nästa anrop. Du behöver inte starta om klustret.

Prestandaoptimering

Svarstid för kallstart

Det första UDF-anropet i en session initierar den isolerade sandboxen, vilket ökar svarstiden. Efterföljande anrop i samma session går snabbare. Ta hänsyn till detta vid benchmarking eller utformning av svarstidskänsliga arbetsbelastningar.

Cachelagring av dyra beräkningar

Om din UDF utför dyr initiering eller beräkning cachelagrar du resultatet för att bara beräkna det en gång.

Scala

Använd ett val fält i Scala-objektet för att cachelagera resultatet:

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

Java

Använd ett static-fält med ett statiskt initieringsblock för att cacha resultatet:

package example;

import java.util.Map;
import java.util.HashMap;

public class CachedUDF {
    // Computed once and cached
    private static Map<String, Double> expensiveData;

    static {
        // Load data from somewhere expensive
        expensiveData = new HashMap<>();
        expensiveData.put("key1", 1.0);
        expensiveData.put("key2", 2.0);
    }

    public static double lookup(String key) {
        return expensiveData.getOrDefault(key, 0.0);
    }
}

Använd DETERMINISTIC när det är lämpligt

Markera din UDF som DETERMINISTIC om den alltid genererar samma utdata för samma indata. På så sätt kan frågeoptimeraren cachelagra resultat och förbättra prestandan.

Limitations

  • Endast skalära UDF:er stöds. Användardefinierade aggregeringsfunktioner (UDAF:er) och användardefinierade tabellfunktioner (UDF: er) stöds inte.
  • UDF:er körs i en isolerad sandbox-miljö utan någon aktiv Spark-session. Spark-API:er (SparkSession, SparkContext, spark.sql(...), DataFrame- och Datauppsättningsåtgärder) är inte tillgängliga.
  • UDF:er kan inte vara beroende av Spark-kärnpaket.
  • UDF:er har inte åtkomst till arbetsytefiler eller Unity Catalog-volymer vid körning.

Metodtips

Databricks rekommenderar följande metoder:

  • Version dina JAR-filer. Till exempel my-udf-0.1.0.jar, my-udf-0.2.0.jar.
  • Verifiera SQL-typmappningar före distributionen. Se Språkmappningar.
  • Bevilja READ VOLUME och EXECUTE behörigheter endast till användare som behöver köra UDF. Använd gruppägarskap för UDF:er som delas mellan team.

Testa UDF:er lokalt

Testa din UDF med enhetstester innan du driftsätter i produktion.

Scala

Om du vill testa src/main/scala/example/MyUDF.scalaskapar du en testfil i src/test/scala/example/MyUDFTest.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)
  }
}

Lägg till testberoendet i build.sbt:

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

Så här kör du testerna:

sbt test

Java

Om du vill testa src/main/java/com/example/MyUDF.javaskapar du en testfil i src/test/java/com/example/MyUDFTest.java:

package com.example;

import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.*;

public class MyUDFTest {
    @Test
    public void testAddOne() {
        assertEquals(6, MyUDF.addOne(5));
    }

    @Test
    public void testAddOneWithNegativeNumbers() {
        assertEquals(0, MyUDF.addOne(-1));
    }
}

Lägg till JUnit-beroendet i avsnittet <dependencies> i din pom.xml:

<dependency>
    <groupId>org.junit.jupiter</groupId>
    <artifactId>junit-jupiter</artifactId>
    <version>5.10.0</version>
    <scope>test</scope>
</dependency>

Så här kör du testerna:

mvn test

Ytterligare resurser