Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Viktigt!
Scala och Java UDF:er kan registreras i Unity Catalog för styrning, återanvändning och identifiering. Se Scala och Java användardefinierade funktioner (UDF: er) i Unity Catalog.
Den här sidan beskriver hur du skapar sessionsomfångsbegränsade Scala- och Java UDF:er i Azure Databricks. Sessionsomfattande UDF:er definieras i en notebook-fil eller ett jobb och gäller endast för den aktuella SparkSession. Information om SQL-språkreferensen finns i Externa användardefinierade skalära funktioner (UDF:er).
Välj din metod
Du kan definiera en Scala eller Java UDF på följande sätt. Information om hur du jämför alla UDF-typer mellan språk, styrning och beräkning finns i Unity Catalog governed vs. session scoped UDF (Enhetskatalog styrd jämfört med sessionsomfattande UDF:er).
| Tillvägagångssätt | Description |
|---|---|
| Inline Scala UDF | Definiera en UDF i en notebook-fil med hjälp av en Scala-funktion eller lambda. Sessionsbaserad. Stöds inte vid serverlös beräkning. |
| Java UDF från en JAR | Registrera en fördefinierad UDF-klass från en JAR med .spark.udf.registerJavaFunction Sessionsbaserad. Stöds vid serverlös beräkning. |
| Unity Catalog-styrd Scala eller Java UDF | Registrera en UDF i Unity Catalog för styrning, återanvändning och identifiering. Stöds vid serverlös beräkning. |
Kravspecifikation
- Scala-UDF:er på beräkningsresurser med Unity Catalog aktiverat och med standardåtkomstläge kräver Databricks Runtime 14.2 eller högre.
- ARM-instansstöd för Scala-UDF:er i Unity Catalog-aktiverade kluster kräver Databricks Runtime 15.2 eller senare.
- Registrering av en Java UDF från en JAR med
spark.udf.registerJavaFunctionkräver Databricks Runtime 18 LTS eller senare. Se Registrera en Java UDF från en JAR.
Viktigt!
Skapa din JAR mot samma Scala- och Apache Spark-versioner som den beräkning som kör den. Ett matchningsfel kan leda till att UDF misslyckas vid registreringen eller anropstiden.
- Klassisk beräkning: Matcha Scala- och Spark-versionerna av din Databricks Runtime-version. Se avsnittet Systemmiljö i versionsanteckningarna för Databricks Runtime och kompatibilitet för din version. Databricks Runtime 18 LTS använder till exempel Scala 2.13.16 och Apache Spark 4.0.
- Serverlös beräkning: Matcha Scala-versionen av din miljöversion. Se Miljöversioner.
Markera Apache Spark-beroendet som provided, så att det inte buntas med i din JAR. Inkludera endast beroenden från tredje part som din UDF använder.
Registrera en funktion som en UDF
Registrera en Scala-funktion som en UDF med hjälp av spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Anropa UDF i Spark SQL
Skapa en temporär vy och anropa sedan UDF i en SQL-fråga:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Använda UDF med DataFrames
Du kan också anropa en UDF med hjälp av DataFrame-API:et:
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"))
Filer med UDF
Viktigt!
Den här funktionen finns i Beta. Arbetsyteadministratörer kan styra åtkomsten till den här funktionen från sidan Förhandsversioner . Se Hantera förhandsversioner av Azure Databricks.
Scala-typen för en fil är FileRef. Använd det som en parameter eller returtyp i en UDF, antingen som en typ på toppnivå eller som kapslad typ. För typen och dess nästlingsregler, se FILE typ.
För att läsa innehållet i en fil i en UDF, anropa en av följande på en FileRef:
-
asLocalFile(): Returnerar enjava.io.Filesom du kan skicka till vilket bibliotek som helst som accepterar en väg. -
open(): Returnerar enjava.io.InputStreamför att läsa filens byte. Uppringaren stänger den.
För att producera en ny FileRef i en UDF, anropa en av följande statiska metoder:
-
FileRef.create(uri): Skapar en referens till filen viduri. -
FileRef.fromBytes(bytes, destinationPath, contentType): Överförbytestill sökvägen för volymendestinationPathsom en extern fil och returnerar en referens. -
FileRef.fromLocalFile(localFile, destinationPath, contentType): Laddar upp en lokal fil till volymsökvägendestinationPathsom en extern fil och returnerar en referens.
Att returnera ett FileRef från en UDF som skriver till en FILE MANAGED-kolumn stöds inte.
För exempel i Python, Scala och SQL, inklusive bildbehandling, filtypsdetektion och videoramsextraktion, se Processfiler med UDF:er.
Registrera en Java UDF från en JAR
Paketera en UDF som en JAR, lägg till den i sessionen med spark.addArtifactoch registrera UDF-klassen med spark.udf.registerJavaFunction.
Kommentar
Stöds i standardåtkomstläge och serverlös beräkning i Databricks Runtime 18 LTS eller senare. Den registrerade funktionen är begränsad till sessionen och är inte registrerad i Unity Catalog.
Följande steg går igenom hur du skapar ett projekt, skriver en UDF-klass, skapar en fet JAR och registrerar den.
Steg 1: Skapa projektet
Konfigurera ett projekt i Scala eller Java.
Scala
Skapa ett nytt Scala-projekt med :sbt
sbt new scala/scala-seed.g8
Ersätt innehållet i build.sbt filen med följande. Ange scalaVersion och versionen spark-sql så att de matchar din beräkning:
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"
)
Aktivera plugin-programmet sbt-assembly för att skapa en fet JAR. Skapa eller redigera project/assembly.sbt och lägg till:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
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.
I den genererade pom.xml lägger du till ett <project></project>-block inom taggarna <properties> och konfigurerar maven-shade-plugin för att bygga en fat JAR-fil:
<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>
Steg 2: Skriv din UDF-klass
Din UDF-klass måste implementera ett av gränssnitten org.apache.spark.sql.api.java.UDF (UDF1 via UDF22), där talet anger hur många indataargument som UDF tar.
call() Implementera metoden med din logik.
Hanteraren måste vara en Java-klass.
spark.udf.registerJavaFunction läser in klassen via reflektion, så den måste vara en offentlig klass på översta nivån (eller static-kapslad) med en offentlig konstruktor utan argument. En Scala class eller object uppfyller inte det här kravet och misslyckas vid anrop. Du kan skapa JAR med sbt, men själva UDF-klassen måste skrivas i Java.
Skapa 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;
}
}
Steg 3: Skapa din feta JAR
Paketera din kompilerade UDF till en fet JAR.
Scala
Kör följande från projektets rotkatalog:
sbt clean assembly
Den feta 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.
Steg 4: Ladda upp jar-filen till en Unity Catalog-volym
Ladda upp JAR-filen till en Unity Catalog-volym så att din beräkning kan komma åt den. Om du inte redan har en volym skapar du en:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Ladda upp din JAR-fil till volymen med Catalog Explorer:
- På din Azure Databricks-arbetsyta klickar du på
Katalog för att öppna Katalogutforskaren.
- Välj katalogen och välj sedan det schema som innehåller volymen.
- Klicka på volymnamnet.
- Klicka på Ladda upp till den här volymen och välj din JAR-fil.
- Klicka på Överför.
- När uppladdningen är klar klickar du på namnet på JAR-filen och klickar sedan på Kopiera sökväg för att kopiera volymsökvägen. Till exempel
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Du behöver den här sökvägen i nästa steg.
Steg 5: Registrera och anropa UDF
Lägg till JAR-filen i sessionen med dess volymsökväg, registrera UDF-klassen och anropa den från 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()
Vid beräkning av serverlöst och standardåtkomstläge måste du skicka en explicit returtyp. Det går inte att utelämna returtypen med UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Användardefinierade aggregeringsfunktioner (UDAF:er) stöds inte med registerJavaFunction.
Frågan returnerar UDF-utdata och bekräftar att funktionen är registrerad och anropsbar:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Utvärderingsordning och nullkontroll
Spark SQL (inklusive SQL och API:er för DataFrame och datauppsättning) garanterar inte ordningen för utvärdering av underuttryck. Spark utvärderar inte indata för en operator eller funktion från vänster till höger. Logiska AND uttryck och OR uttryck har inte vänster-till-höger-kortslutningssemantik.
Förlita dig inte på bieffekterna eller utvärderingsordningen för booleska uttryck, eller på ordningen på WHERE- och HAVING-satser. Frågeoptimeraren kan ändra ordning på dessa uttryck och satser. Om en UDF förlitar sig på kortslutningssemantik för null-kontroll garanterar Spark inte att null-kontrollen körs före UDF. Ett exempel:
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
Den här WHERE satsen garanterar inte att Spark anropar UDF när den filtrerar strlen bort null-värden.
För att hantera null-kontroll rekommenderar Databricks något av följande:
- Gör själva UDF:en medveten om null-värden och hantera null-kontrollen inuti UDF:en
- Använd
IFellerCASE WHENuttryck för att göra null-kontrollen och anropa UDF i en villkorsstyrd gren
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:er för typade datauppsättningar
Kommentar
Den här funktionen stöds i Unity Catalog-aktiverade kluster med standardåtkomstläge i Databricks Runtime 15.4 och senare.
Använd inskrivna API:er för datauppsättningar för att köra transformeringar som mappning, filter och aggregeringar på datauppsättningar med en användardefinierad funktion.
I följande exempel används API:et map() för att ändra ett tal i en resultatkolumn till en prefixsträng:
spark.range(3).map(f => s"row-$f").show()
I det här exemplet används map(), men samma mönster gäller för andra typerade API:er för datauppsättningar som filter(), mapPartitions(), foreach(), foreachPartition(), reduce()och flatMap().
Scala UDF-funktioner och Databricks Runtime-kompatibilitet
Följande funktioner kräver minsta Databricks Runtime-versioner på Unity Catalog-aktiverade kluster i standardåtkomstläge (delad).
| Egenskap | Minsta Databricks Runtime-version |
|---|---|
| Skalära UDF:er | 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 |
(Direktuppspelning) foreachWriter Sink |
Databricks Runtime 15.4 |
(Direktuppspelning) foreachBatch |
Databricks Runtime 16.1 |
(Direktuppspelning) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks Runtime 16.2 |
spark.udf.registerJavaFunction(Java UDF från en JAR) |
Databricks Runtime 18 LTS |