Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
Important
Scala en Java UDF's kunnen worden geregistreerd in Unity Catalog voor governance, hergebruik en detectie. Zie Scala en Java door de gebruiker gedefinieerde functies (UDF's) in Unity Catalog.
Op deze pagina wordt beschreven hoe u scala met sessiebereik en Java UDF's maakt in Azure Databricks. UDF's met sessiebereik worden gedefinieerd in een notebook of taak en zijn alleen van toepassing op de huidige SparkSession. Zie Externe door de gebruiker gedefinieerde scalaire functies (UDF's) voor de sql-taalverwijzing.
Kies uw benadering
U kunt een Scala of Java UDF op de volgende manieren definiëren. Als u alle UDF-typen wilt vergelijken voor talen, governance en berekeningen, raadpleegt u Unity Catalog beheerd versus UDF's met sessiebereik.
| Approach | Beschrijving |
|---|---|
| Inline Scala UDF | Definieer een UDF in een notebook met behulp van een Scala-functie of lambda. Sessiegebonden. Niet ondersteund op serverloze berekeningen. |
| Java UDF vanuit een JAR | Registreer een vooraf gecompileerde UDF-klasse uit een JAR met behulp van spark.udf.registerJavaFunction. Sessiegebonden. Ondersteund op serverloze berekeningen. |
| Scala- of Java-UDF onder Unity Catalog-beheer | Registreer een UDF in Unity Catalog voor governance, hergebruik en vindbaarheid. Ondersteund op serverloze berekeningen. |
Requirements
- Scala UDF's op compute met Unity Catalog met standaardtoegangsmodus vereisen Databricks Runtime 14.2 of hoger.
- Arm-exemplaarondersteuning voor Scala UDF's op clusters met Unity Catalog vereist Databricks Runtime 15.2 of hoger.
- Voor het registreren van een Java UDF vanuit een JAR met
spark.udf.registerJavaFunctionis Databricks Runtime 18 LTS of hoger vereist. Zie Een Java UDF registreren vanuit een JAR.
Important
Bouw uw JAR op basis van dezelfde Scala- en Apache Spark-versies als de berekening waarmee het wordt uitgevoerd. Een afwijking kan ertoe leiden dat de UDF niet werkt bij de registratie of tijdens het aanroepen.
- Klassieke rekenkracht: overeenkomen met de Scala- en Spark-versies van uw Databricks Runtime-versie. Zie de sectie Systeemomgeving van de releaseopmerkingen over versies en compatibiliteit van Databricks Runtime voor uw versie. Databricks Runtime 18 LTS maakt bijvoorbeeld gebruik van Scala 2.13.16 en Apache Spark 4.0.
- Serverloze berekening: komt overeen met de Scala-versie van uw omgevingsversie. Zie Serverloze Omgevingsversies.
Markeer de Apache Spark-afhankelijkheid als provided zodat deze niet in uw JAR wordt gebundeld. Neem alleen afhankelijkheden van derden op die door uw UDF worden gebruikt.
Een functie registreren als een UDF
Registreer een Scala-functie als UDF met behulp van spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
De UDF aanroepen in Spark SQL
Maak een tijdelijke weergave en roep vervolgens de UDF aan in een SQL-query:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
UDF gebruiken met DataFrames
U kunt ook een UDF aanroepen met behulp van de 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"))
Een Java UDF registreren vanuit een JAR
Pak een UDF als JAR in, voeg deze toe aan uw sessie met spark.addArtifacten registreer de UDF-klasse bij spark.udf.registerJavaFunction.
Notitie
Ondersteund in de standaardtoegangsmodus en serverloze berekeningen in Databricks Runtime 18 LTS of hoger. De geregistreerde functie is sessiegebonden en is niet geregistreerd in Unity Catalog.
De volgende stappen helpen u bij het maken van een project, het schrijven van een UDF-klasse, het bouwen van een fat JAR en het registreren ervan.
Stap 1: Uw project maken
Stel een project in Scala of Java in.
Scala
Maak een nieuw Scala-project met behulp van sbt:
sbt new scala/scala-seed.g8
Vervang de inhoud van het build.sbt bestand door het volgende. Stel scalaVersion en de versie van spark-sql zo in dat deze overeenkomen met uw rekenomgeving:
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"
)
Schakel de invoegtoepassing sbt-assembly in om een vet JAR-bestand te bouwen. Maak of bewerk project/assembly.sbt en voeg toe:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Maak een nieuw Maven-project met behulp van het quickstart-archetype:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Met deze opdracht maakt u de standaard Maven-projectstructuur met src/main/java en src/test/java directory's.
Voeg in de gegenereerde pom.xml, binnen de <project></project>-tags, een <properties>-blok toe en configureer de maven-shade-plugin om een fat JAR te bouwen:
<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>
Stap 2: Uw UDF-klasse schrijven
Uw UDF-klasse moet een van de org.apache.spark.sql.api.java.UDF interfaces (UDF1 via UDF22) implementeren, waarbij het getal aangeeft hoeveel invoerargumenten de UDF gebruikt. Implementeer de call() methode met uw logica.
De handler moet een Java klasse zijn.
spark.udf.registerJavaFunction laadt de klasse via reflectie, dus de klasse moet een public-klasse op het hoogste niveau zijn (of een static geneste klasse) met een openbare constructor zonder argumenten. Een Scala class of object voldoet niet aan deze vereiste en mislukt tijdens het aanroepen. U kunt de JAR bouwen met sbt, maar de UDF-klasse zelf moet worden geschreven in Java.
Maak 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;
}
}
Stap 3: uw vet-JAR bouwen
Pak uw gecompileerde UDF in een vet JAR-bestand.
Scala
Voer vanuit de hoofdmap van uw project het volgende uit:
sbt clean assembly
De fat JAR wordt gemaakt in target/scala-2.13/ met een naam als my-udf-assembly-0.1.0-SNAPSHOT.jar.
Java
Voer vanuit de hoofdmap van uw project het volgende uit:
mvn clean package
De fat JAR wordt gemaakt in target/ met een naam als my-udf-1.0-SNAPSHOT.jar.
Stap 4: Uw JAR uploaden naar een Unity Catalog-volume
Upload het JAR-bestand naar een Unity Catalog-volume , zodat uw rekenproces er toegang toe heeft. Als u nog geen volume hebt, maakt u er een:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Upload uw JAR-bestand naar het volume met behulp van Catalog Explorer:
- Klik in uw Azure Databricks werkruimte op
Catalog om Catalog Explorer te openen.
- Selecteer de catalogus en selecteer vervolgens het schema dat uw volume bevat.
- Klik op de volumenaam.
- Klik op Uploaden naar dit volume en selecteer uw JAR-bestand.
- Klik op Uploaden.
- Nadat het uploaden is voltooid, klikt u op de naam van het JAR-bestand en klikt u vervolgens op Pad kopiëren om het volumepad te kopiëren. Bijvoorbeeld:
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. U hebt dit pad nodig in de volgende stap.
Stap 5: De UDF registreren en aanroepen
Voeg het JAR-bestand toe aan uw sessie met behulp van het volumepad, registreer de UDF-klasse en roep deze aan vanuit 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()
Op serverloze en standaardtoegangsmodus berekenen moet u een expliciet retourtype doorgeven. Het weglaten van het retourtype mislukt met UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Door de gebruiker gedefinieerde statistische functies (UDAF's) worden niet ondersteund met registerJavaFunction.
De query retourneert de UDF-uitvoer, waarbij wordt bevestigd dat de functie is geregistreerd en aanroepbaar is:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Evaluatievolgorde en nullcontrole
Spark SQL (inclusief SQL en dataframe- en gegevensset-API's) garandeert niet de volgorde van de subexpressie-evaluatie. Spark evalueert de invoer van een operator of functie niet van links naar rechts. Logische expressies AND en OR hebben geen kortsluitsemantiek van links naar rechts.
Vertrouw niet op de bijwerkingen of evaluatievolgorde van Boole-expressies of de volgorde van WHERE en HAVING componenten. De queryoptimalisatie kan deze expressies en componenten opnieuw ordenen. Als een UDF afhankelijk is van kortsluitingssemantiek voor null-controle, garandeert Spark niet dat de null-controle wordt uitgevoerd vóór de UDF. Voorbeeld:
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
Deze WHERE component garandeert niet dat Spark de strlen UDF aanroept nadat null-waarden zijn gefilterd.
Databricks raadt een van de volgende opties aan om null-controles af te handelen:
- De UDF zelf null-aware maken en null-controle uitvoeren in de UDF
- Gebruik
IFofCASE WHENexpressies om de null-controle uit te voeren en de UDF aan te roepen in een voorwaardelijke vertakking
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's voor getypte datasets
Notitie
Deze functie wordt ondersteund op clusters met Unity Catalog-functionaliteit met standaardtoegangsmodus in Databricks Runtime 15.4 en hoger.
Gebruik getypte gegevensset-API's om transformaties uit te voeren, zoals toewijzing, filter en aggregaties op gegevenssets met een door de gebruiker gedefinieerde functie.
In het volgende voorbeeld wordt de map() API gebruikt om een getal in een resultaatkolom te wijzigen in een voorvoegseltekenreeks:
spark.range(3).map(f => s"row-$f").show()
In dit voorbeeld wordt gebruikgemaakt van map(), maar hetzelfde patroon geldt ook voor andere getypeerde Dataset-API's, zoals filter(), mapPartitions(), foreach(), foreachPartition(), reduce() en flatMap().
Scala UDF-functies en compatibiliteit met Databricks Runtime
Voor de volgende functies zijn minimale Databricks Runtime-versies op clusters met Unity Catalog-functionaliteit vereist in de standaardtoegangsmodus (gedeeld).
| Eigenschap | Minimale Databricks Runtime-versie |
|---|---|
| Scalaire UDFs | Databricks Runtime 14.2 is een geavanceerd platform voor dataverwerking en analyse. |
Dataset.map
Dataset.mapPartitions, Dataset.filter, Dataset.reduceDataset.flatMap |
Databricks Runtime 15.4 |
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups |
Databricks Runtime 15.4 |
(Streamen) foreachWriter Sink |
Databricks Runtime 15.4 |
(Streamen) foreachBatch |
Databricks Runtime 16.1 |
(Streamen) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks Runtime 16.2 |
spark.udf.registerJavaFunction(Java UDF vanuit een JAR) |
Databricks Runtime 18 LTS |