UDF-y Scala i Java w zakresie bieżącej sesji

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.registerJavaFunction wymaga ś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.

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 a java.io.File , które możesz przekazać do dowolnej biblioteki akceptującej ścieżkę.
  • open(): Zwraca a java.io.InputStream do 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ła bytes do ścieżki destinationPath Volume jako plik zewnętrzny i zwraca referencję.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Przesyła lokalny plik do ścieżki destinationPath Volume 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:

  1. W obszarze roboczym Azure Databricks kliknij pozycję Ikona Danych.Katalog aby otworzyć Eksplorator Katalogu.
  2. Wybierz katalog, a następnie wybierz schemat zawierający wolumin.
  3. Kliknij nazwę woluminu.
  4. Kliknij Prześlij do tego woluminu i wybierz plik JAR.
  5. Kliknij Przekaż.
  6. 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ń IF lub CASE WHEN do 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