Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Importante
Le funzioni definite dall'utente Scala e Java possono essere registrate in Unity Catalog ai fini della governance, del riutilizzo e della facilità di individuazione. Vedi le funzioni definite dall'utente (UDF) in Scala e Java in Unity Catalog.
Questa pagina descrive come creare UDF in Scala e Java limitate alla sessione in Azure Databricks. Le UDF con ambito di sessione vengono definite in un notebook o in un job e si applicano solo alla SparkSession corrente. Per informazioni di riferimento sul linguaggio SQL, vedere Funzioni scalari definite dall'utente esterne.
Scegliere l'approccio
È possibile definire una UDF Scala o Java nei seguenti modi. Per confrontare tutti i tipi di UDF per linguaggio, governance e calcolo, consulta UDF governate da Unity Catalog e UDF con ambito di sessione.
| Avvicinarsi | Descrizione |
|---|---|
| Funzione inline definita dall'utente in Scala | Definisci una UDF in un notebook tramite una funzione Scala o una lambda. Limitato alla sessione. Non supportato nell'ambiente di elaborazione serverless. |
| UDF Java da un file JAR | Registrare una classe UDF precompilata da un JAR usando spark.udf.registerJavaFunction. Limitato alla sessione. Supportato nel calcolo serverless. |
| UDF Scala o Java governata da Unity Catalog | Registra una UDF in Unity Catalog per la governance, il riutilizzo e la rintracciabilità. Supportato in ambiente serverless. |
Requisiti
- Le UDF Scala nelle risorse di calcolo con Unity Catalog abilitato e modalità di accesso standard richiedono Databricks Runtime 14.2 o versioni successive.
- Il supporto dell'istanza ARM per le funzioni definite dall'utente Scala nei cluster abilitati per Unity Catalog richiede Databricks Runtime 15.2 o versione successiva.
- La registrazione di una funzione Java definita dall'utente da un file JAR con
spark.udf.registerJavaFunctionrichiede Databricks Runtime 18.3 o versioni successive. Vedi Registrare una UDF Java a partire da un file JAR.
Importante
Compilare il file JAR con le stesse versioni di Scala e Apache Spark dell'ambiente di calcolo che lo esegue. Un'incompatibilità può causare il fallimento della UDF in fase di registrazione o di chiamata.
- Elaborazione classica: abbina le versioni di Scala e Spark a quella di Databricks Runtime. Consultare la sezione Ambiente di sistema di Versioni e compatibilità delle note sulla versione di Databricks Runtime per la propria versione. Ad esempio, Databricks Runtime 18.3 usa Scala 2.13.16 e Apache Spark 4.0.
- Calcolo serverless: Fai corrispondere la versione di Scala alla versione dell'ambiente. Consultare le versioni dell'ambiente serverless .
Contrassegna la dipendenza Apache Spark come provided in modo che non venga inclusa nel tuo JAR. Includi solo le dipendenze di terze parti utilizzate dalla funzione definita dall'utente.
Registrare una funzione come UDF
Registrare una funzione Scala come funzione definita dall'utente usando spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Invocare l’UDF in Spark SQL
Crea una vista temporanea, quindi richiama l'UDF in una query SQL:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Usare la UDF con DataFrames
È anche possibile chiamare una UDF tramite l'API DataFrame:
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"))
Registrare una UDF Java da un JAR
Crea un pacchetto di una UDF come JAR, aggiungilo alla sessione con spark.addArtifact e registra la classe UDF con spark.udf.registerJavaFunction.
Nota
Supportato in modalità di accesso standard e nel calcolo serverless in Databricks Runtime 18.3 o versione successiva. La funzione registrata è con ambito sessione e non è registrata in Unity Catalog.
I passaggi seguenti illustrano la creazione di un progetto, la scrittura di una classe UDF, la compilazione di un file JAR fat e la registrazione.
Passaggio 1: Creare il progetto
Configurare un progetto in Scala o Java.
Scala
Creare un nuovo progetto Scala usando sbt:
sbt new scala/scala-seed.g8
Sostituire il contenuto del build.sbt file con quanto segue. Impostare scalaVersion e la spark-sql versione in modo che corrispondano al calcolo:
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"
)
Abilitare il plug-in sbt-assembly per creare un file JAR fat. Creare o modificare project/assembly.sbt e aggiungere:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Creare un nuovo progetto Maven usando l'archetipo di avvio rapido:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Questo comando crea la struttura standard di un progetto Maven con le directory src/main/java e src/test/java.
Nel file generato pom.xml, all'interno dei tag <project></project>, aggiungi un blocco <properties> e configura maven-shade-plugin per creare un JAR fat:
<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>
Passaggio 2: Scrivere la classe UDF
La classe UDF deve implementare una delle interfacce org.apache.spark.sql.api.java.UDF (da UDF1 a UDF22), dove il numero indica quanti argomenti di input accetta la UDF. Implementa il metodo call() con la tua logica.
Il gestore deve essere una classe Java.
spark.udf.registerJavaFunction carica la classe tramite reflection, pertanto deve essere una classe pubblica di primo livello (o static annidata) con un costruttore no-arg pubblico. Uno Scala class o object non soddisfa questo requisito e genera un errore al momento della chiamata. È possibile compilare il file JAR con sbt, ma la classe UDF stessa deve essere scritta in Java.
Creare 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;
}
}
Passaggio 3: Crea il tuo JAR fat
Crea un fat JAR contenente la UDF compilata.
Scala
Dalla directory radice del progetto eseguire:
sbt clean assembly
Il file JAR fat viene creato in target/scala-2.13/ con un nome come my-udf-assembly-0.1.0-SNAPSHOT.jar.
Java
Dalla directory radice del progetto eseguire:
mvn clean package
Il file JAR fat viene creato in target/ con un nome come my-udf-1.0-SNAPSHOT.jar.
Passaggio 4: Caricare il file JAR in un volume del catalogo Unity
Carica il file JAR in un volume di Unity Catalog in modo che le risorse di calcolo possano accedervi. Se non hai già un volume, creane uno:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Caricare il file JAR nel volume usando Esplora cataloghi:
- Nell'area di lavoro Azure Databricks, fare clic su
Catalogo per aprire Esplora Catalogo.
- Selezionare il catalogo, quindi selezionare lo schema che contiene il volume.
- Fare clic sul nome del volume.
- Fare clic su Carica in questo volume e selezionare il file JAR.
- Fare clic su Carica.
- Al termine del caricamento, fare clic sul nome del file JAR e quindi su Copia percorso per copiare il percorso del volume. Ad esempio:
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Questo percorso è necessario nel passaggio successivo.
Passaggio 5: Registrare e chiamare l'UDF
Aggiungi il file JAR alla sessione usando il percorso del volume, registra la classe UDF e richiamala in 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()
Nell’ambiente di calcolo con modalità di accesso serverless e standard, è necessario specificare un tipo di ritorno esplicito. Omettere il tipo restituito genera un errore con UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Le funzioni di aggregazione definite dall'utente non sono supportate con registerJavaFunction.
La query restituisce l'output della UDF, confermando che la funzione è registrata e può essere richiamata:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
"Ordine di valutazione e controllo dei null"
Spark SQL (inclusi SQL e le API dataframe e set di dati) non garantisce l'ordine di valutazione della sottoespressione. Spark non valuta gli input di un operatore o di una funzione da sinistra a destra. Le espressioni e AND logiche OR non hanno semantica di corto circuito da sinistra a destra.
Non basarsi sugli effetti collaterali o sull'ordine di valutazione delle espressioni booleane o sull'ordine delle WHERE clausole e HAVING . Query Optimizer può riordinare queste espressioni e clausole. Se una funzione definita dall'utente si basa sulla semantica di corto circuito per il controllo null, Spark non garantisce che il controllo Null venga eseguito prima della funzione definita dall'utente. Per esempio:
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
Questa WHERE condizione non garantisce che Spark esegua la strlen UDF dopo aver escluso i valori null.
Per gestire il controllo null, Databricks consiglia una delle opzioni seguenti:
- Fare in modo che l'UDF sia in grado di gestire i valori null ed eseguire il controllo dei valori null all'interno dell'UDF
- Usare le espressioni
IFoCASE WHENper eseguire il controllo di null e richiamare l'UDF in un ramo condizionale
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 del set di dati tipizzato
Nota
Questa funzionalità è supportata nei cluster abilitati per Unity Catalog con modalità di accesso standard in Databricks Runtime 15.4 e versioni successive.
Usare le API del set di dati tipizzato per eseguire trasformazioni, ad esempio mapping, filtro e aggregazioni nei set di dati con una funzione definita dall'utente.
L'esempio seguente usa l'API map() per modificare un numero in una colonna di risultato in una stringa con prefisso:
spark.range(3).map(f => s"row-$f").show()
Questo esempio utilizza map(), ma lo stesso schema si applica ad altre API di Dataset tipizzati come filter(), mapPartitions(), foreach(), foreachPartition(), reduce() e flatMap().
Funzionalità di UDF in Scala e compatibilità con Databricks Runtime
Le funzionalità seguenti richiedono versioni minime di Databricks Runtime nei cluster abilitati per Unity Catalog in modalità di accesso standard (condiviso).
| Caratteristica / Funzionalità | Versione minima di Databricks Runtime |
|---|---|
| Funzioni definite dall'utente scalari | 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 |
(Streaming) foreachWriter Sink |
Databricks Runtime 15.4 |
(Streaming) foreachBatch |
Databricks Runtime 16.1 |
(Streaming) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks Runtime 16.2 |
spark.udf.registerJavaFunction (UDF Java da un file JAR) |
Databricks Runtime 18.3 |