Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Ważna
Funkcje UDF w językach Scala i Java można rejestrować w Unity Catalog w celu zarządzania, ponownego wykorzystania i łatwiejszego odnajdywania. Zobacz funkcje zdefiniowane przez użytkownika (UDF) w językach Scala i Java w Unity Catalog.
Na tej stronie opisano, jak tworzyć funkcje UDF języków Scala i Java o zakresie sesji w Azure Databricks. UDF o zakresie sesji są tworzone w notatniku lub zadaniu i obowiązują tylko w bieżącej sesji SparkSession. Aby zapoznać się z dokumentacją języka SQL, zobacz Zewnętrzne funkcje skalarne zdefiniowane przez użytkownika (UDF).
Wybierz swoje podejście
Funkcję UDF w języku Scala lub Java można zdefiniować na następujące sposoby. Aby porównać wszystkie typy funkcji UDF w różnych językach, sposobach zarządzania i środowiskach obliczeniowych, zobacz Funkcje UDF zarządzane przez Unity Catalog i funkcje UDF o zakresie sesji.
| Approach | Opis |
|---|---|
| Funkcja UDF inline w języku Scala | Zdefiniuj funkcję UDF w notatniku przy użyciu funkcji w języku Scala lub wyrażenia lambda. Zakres sesji. Nieobsługiwane w obliczeniach bezserwerowych. |
| Funkcja zdefiniowana przez użytkownika Java z pliku JAR | Zarejestruj wstępnie skompilowaną klasę UDF z pliku JAR przy użyciu polecenia spark.udf.registerJavaFunction. Zakres sesji. Obsługiwane w obliczeniach bezserwerowych. |
| Funkcja UDF Scala lub Java zarządzana przez Unity Catalog | Zarejestruj funkcję UDF w Unity Catalog, aby zapewnić nadzór, możliwość ponownego użycia i wykrywalność. Obsługiwane w obliczeniach bezserwerowych. |
Requirements
- Funkcje UDF języka Scala w obliczeniach z włączonym Unity Catalog, działających w standardowym trybie dostępu, wymagają środowiska Databricks Runtime w wersji 14.2 lub nowszej.
- Obsługa instancji ARM dla funkcji UDF w języku Scala w klastrach z włączonym katalogiem Unity wymaga Runtime Databricks 15.2 lub nowszego.
- Rejestrowanie funkcji UDF języka Java z pliku JAR za pomocą
spark.udf.registerJavaFunctionwymaga środowiska Databricks Runtime w wersji 18 LTS lub nowszej. Zobacz Jak zarejestrować funkcję UDF języka Java z pliku JAR.
Ważna
Skompiluj plik JAR dla tych samych wersji języka Scala i platformy Apache Spark, co obliczenia, które go uruchamiają. Niezgodność może spowodować niepowodzenie funkcji UDF podczas rejestracji lub wywołania.
- Klasyczne obliczenia: dopasuj wersje Scala i Spark do wersji Databricks Runtime. Zobacz sekcję Środowisko systemowe w dokumencie Informacje o wydaniu Databricks Runtime: wersje i zgodność dla swojej wersji. Na przykład środowisko Databricks Runtime 18 LTS używa języka Scala 2.13.16 i apache Spark 4.0.
- Obliczenia bezserwerowe: Dopasuj wersję języka Scala do wersji środowiska. Zobacz wersje środowiska bezserwerowego.
Oznacz zależność Apache Spark jako provided, aby nie została dołączona do pliku JAR. Uwzględnij tylko zależności zewnętrzne używane przez funkcję UDF.
Rejestrowanie funkcji jako UDF
Zarejestruj funkcję Scala jako UDF za pomocą spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Wywołaj UDF w Spark SQL
Utwórz widok tymczasowy, a następnie wywołaj funkcję UDF w zapytaniu SQL:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Używanie funkcji UDF z ramkami danych
Można również wywołać funkcję zdefiniowaną przez użytkownika (UDF) za pomocą interfejsu DataFrame API:
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"))
Pliki z UDF
Ważna
Ta funkcja jest dostępna w wersji beta. Administratorzy obszaru roboczego mogą kontrolować dostęp do tej funkcji ze strony Podglądy . Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.
Typ pliku w Scali to FileRef. Używaj go jako parametru lub typu zwrotnego w UDF, albo jako typ najwyższego poziomu, albo zagnieżdżony. Dla typu i jego reguł zagnieżdżania zobacz FILE typ.
Aby odczytać zawartość pliku w UDF, wywołaj jeden z następujących na :FileRef
-
asLocalFile(): Zwraca ajava.io.File, które możesz przekazać do dowolnej biblioteki akceptującej ścieżkę. -
open(): Zwraca ajava.io.InputStreamdo odczytu bajtów pliku. Dzwoniący zamyka telefon.
Aby wygenerować nowy FileRef w UDF, wywołaj jedną z następujących metod statycznych:
-
FileRef.create(uri): Tworzy odwołanie do pliku w .uri -
FileRef.fromBytes(bytes, destinationPath, contentType): Przesyłabytesdo ścieżkidestinationPathVolume jako plik zewnętrzny i zwraca referencję. -
FileRef.fromLocalFile(localFile, destinationPath, contentType): Przesyła lokalny plik do ścieżkidestinationPathVolume jako plik zewnętrzny i zwraca referencję.
Nie obsługuje się zwrotu z FileRef UDF, który zapisuje do kolumny FILE MANAGED .
Przykłady w Python, Scala i SQL, w tym przetwarzanie obrazów, wykrywanie typów plików i ekstrakcję klatek wideo, można znaleźć w artykule Pliki procesowe z UDF.
Zarejestruj funkcję Java UDF z pliku JAR
Spakuj funkcję zdefiniowaną przez użytkownika do pliku JAR, dodaj go do sesji za pomocą polecenia spark.addArtifact i zarejestruj klasę UDF za pomocą polecenia spark.udf.registerJavaFunction.
Uwaga
Obsługiwane w standardowym trybie dostępu i bezserwerowych obliczeniach w środowisku Databricks Runtime 18 LTS lub nowszym. Zarejestrowana funkcja ma zakres sesji i nie jest rejestrowana w Unity Catalog.
Poniższe kroki przeprowadzają przez tworzenie projektu, napisanie klasy UDF, zbudowanie grubego pliku JAR oraz jego zarejestrowanie.
Krok 1. Tworzenie projektu
Skonfiguruj projekt w języku Scala lub Java.
Scala
Utwórz nowy projekt Scala przy użyciu polecenia sbt:
sbt new scala/scala-seed.g8
Zastąp zawartość build.sbt pliku następującym kodem. Ustaw scalaVersion oraz wersję spark-sql tak, aby odpowiadały Twoim zasobom obliczeniowym:
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"
)
Włącz wtyczkę sbt-assembly, aby utworzyć fat JAR. Utwórz lub edytuj project/assembly.sbt i dodaj:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Utwórz nowy projekt Maven przy użyciu archetypu szybkiego startu:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
To polecenie tworzy standardową strukturę projektu Maven z katalogami src/main/java i .src/test/java
W wygenerowanym pliku pom.xml, wewnątrz tagów <project></project> dodaj blok <properties> i skonfiguruj element maven-shade-plugin, aby zbudować fat JAR:
<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>
Krok 2: Napisz swoją klasę UDF
Klasa UDF musi zaimplementować jeden z org.apache.spark.sql.api.java.UDF interfejsów (UDF1 za pośrednictwem UDF22), gdzie liczba wskazuje, ile argumentów wejściowych przyjmuje funkcja UDF. Zaimplementuj metodę call(), używając własnej logiki.
Program obsługi musi być klasą Java.
spark.udf.registerJavaFunction ładuje klasę za pomocą mechanizmu refleksji, więc musi to być publiczna klasa najwyższego poziomu (lub publiczna klasa zagnieżdżona static) z publicznym konstruktorem bezargumentowym. Scala class lub object nie spełnia tego wymagania i kończy się niepowodzeniem w czasie wywołania. Plik JAR można skompilować za pomocą biblioteki sbt, ale sama klasa UDF musi być zapisana w Java.
Utwórz 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;
}
}
Krok 3. Tworzenie tłuszczu JAR
Spakuj skompilowany UDF w gruby plik JAR.
Scala
W katalogu głównym projektu uruchom polecenie:
sbt clean assembly
Plik fat JAR jest tworzony w target/scala-2.13/ pod nazwą taką jak my-udf-assembly-0.1.0-SNAPSHOT.jar.
Java
W katalogu głównym projektu uruchom polecenie:
mvn clean package
Plik fat JAR jest tworzony w target/ pod nazwą taką jak my-udf-1.0-SNAPSHOT.jar.
Krok 4: Przekaż plik JAR do woluminu Unity Catalog
Prześlij plik JAR do woluminu Unity Catalog, aby środowisko obliczeniowe mogło uzyskać do niego dostęp. Jeśli jeszcze nie masz woluminu, utwórz go:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Przekaż plik JAR do woluminu przy użyciu Eksploratora wykazu:
- W obszarze roboczym Azure Databricks kliknij pozycję
Katalog aby otworzyć Eksplorator Katalogu.
- Wybierz katalog, a następnie wybierz schemat zawierający wolumin.
- Kliknij nazwę woluminu.
- Kliknij Prześlij do tego woluminu i wybierz plik JAR.
- Kliknij Przekaż.
- Po zakończeniu przesyłania kliknij nazwę pliku JAR, a następnie kliknij Kopiuj ścieżkę, aby skopiować ścieżkę woluminu. Na przykład
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Potrzebujesz tej ścieżki w następnym kroku.
Krok 5: Zarejestrować i wywołać UDF
Dodaj plik JAR do sesji przy użyciu ścieżki woluminu, zarejestruj klasę UDF i wywołaj go z usługi 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()
W środowisku obliczeniowym w trybie bezserwerowym i w trybie dostępu standardowego musisz podać jawnie określony typ zwracany. Pominięcie typu zwracanego kończy się niepowodzeniem w przypadku UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Funkcje agregujące zdefiniowane przez użytkownika (UDAFs) nie są obsługiwane w systemie registerJavaFunction.
Zapytanie zwraca wynik działania UDF, potwierdzając, że funkcja jest zarejestrowana i można ją wywołać:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Kolejność oceny i sprawdzanie wartości null
Spark SQL (w tym SQL oraz interfejsy API DataFrame i Dataset) nie gwarantuje kolejności obliczania podwyrażeń. Platforma Spark nie ocenia danych wejściowych operatora ani funkcji od lewej do prawej. Wyrażenia logiczne AND i OR nie mają semantyki zwarć od lewej do prawej.
Nie polegaj na skutkach ubocznych ani na kolejności oceniania wyrażeń boolowskich, ani na kolejności klauzul WHERE i HAVING. Optymalizator zapytań może zmienić kolejność tych wyrażeń i klauzul. Jeśli funkcja UDF opiera się na semantyce zwarciowej na potrzeby sprawdzania wartości null, platforma Spark nie gwarantuje, że sprawdzanie wartości null jest uruchamiane przed funkcją UDF. Przykład:
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
Ta WHERE klauzula nie gwarantuje, że Spark wywoła funkcję UDF strlen po odfiltrowaniu wartości null.
Aby obsługiwać sprawdzanie wartości null, Databricks zaleca jedno z poniższych rozwiązań:
- Spraw, aby sama funkcja zdefiniowana przez użytkownika obsługiwała wartości null, i wykonuj sprawdzanie wartości null wewnątrz tej funkcji
- Użyj wyrażeń
IFlubCASE WHENdo sprawdzenia wartości null i wywołania funkcji zdefiniowanej przez użytkownika w gałęzi warunkowej.
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
Typizowane interfejsy API zestawu danych
Uwaga
Ta funkcja jest obsługiwana w klastrach z obsługą Unity Catalog w standardowym trybie dostępu w środowisku Databricks Runtime 15.4 lub nowszym.
Użyj typowych interfejsów API zestawu danych, aby uruchamiać przekształcenia, takie jak mapowanie, filtrowanie i agregacje w zestawach danych z funkcją zdefiniowaną przez użytkownika.
W poniższym przykładzie użyto interfejsu map() API, aby zmodyfikować liczbę w kolumnie wynikowej na prefiksowany ciąg:
spark.range(3).map(f => s"row-$f").show()
W tym przykładzie użyto map(), ale ten sam wzorzec dotyczy także innych typowanych interfejsów API Dataset, takich jak filter(), mapPartitions(), foreach(), foreachPartition(), reduce() i flatMap().
Funkcje UDF w języku Scala i zgodność środowiska Databricks Runtime
Poniższe funkcje wymagają co najmniej określonych wersji Databricks Runtime w klastrach z włączonym Unity Catalog, działających w standardowym (współdzielonym) trybie dostępu.
| Funkcja | Minimalna wersja środowiska Databricks Runtime |
|---|---|
| Skalarne funkcje zdefiniowane przez użytkownika | Databricks Runtime 14.2 |
Dataset.map, , Dataset.mapPartitions, Dataset.filter, , Dataset.reduceDataset.flatMap |
Databricks Runtime 15.4 |
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups |
Databricks Runtime 15.4 |
(Przesyłanie strumieniowe) foreachWriter Sink |
Databricks Runtime 15.4 |
(Przesyłanie strumieniowe) foreachBatch |
Databricks Runtime 16.1 |
(Przesyłanie strumieniowe) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks Runtime 16.2 |
spark.udf.registerJavaFunction (funkcja Java zdefiniowana przez użytkownika z pliku JAR) |
Databricks Runtime 18 LTS |