Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Important
UDF ve Scale a Javě lze zaregistrovat v Unity Catalogu pro účely správy, opakovaného použití a snadnějšího vyhledávání. Viz Scala a Java uživatelem definované funkce (UDF) v katalogu Unity.
Tato stránka popisuje, jak v Azure Databricks vytvořit uživatelsky definované funkce (UDF) v jazycích Scala a Java s rozsahem platnosti relace. UDF s rozsahem relace jsou definovány v poznámkovém bloku nebo úloze a platí pouze pro aktuální SparkSession. Referenční informace k jazyku SQL najdete v tématu Externí uživatelem definované skalární funkce (UDF).
Volba přístupu
Scala nebo Java UDF můžete definovat následujícími způsoby. Chcete-li porovnat všechny typy UDF napříč jazyky, správou a výpočetními prostředími, přečtěte si UDF řízené pomocí Unity Catalogu vs. UDF s rozsahem relace.
| Approach | Description |
|---|---|
| Vložená uživatelsky definovaná funkce Scala | Definujte UDF v poznámkovém bloku pomocí funkce nebo lambda výrazu v jazyce Scala. Vymezená relace. Nepodporuje se na bezserverových výpočetních prostředcích. |
| Java UDF ze souboru JAR | Zaregistrujte předkompilovanou třídu UDF ze souboru JAR pomocí spark.udf.registerJavaFunction. Vymezená relace. Podporováno na výpočetních prostředcích bez serveru. |
| Scala nebo Java UDF spravovaná službou Unity Catalog | Zaregistrujte uživatelsky definovanou funkci (UDF) v Unity Catalog pro správu, opakované použití a snadnější vyhledatelnost. Podporováno na výpočetních prostředcích bez serveru. |
Požadavky
- Funkce UDF v jazyce Scala ve výpočetních prostředcích s povoleným Unity Catalog a standardním režimem přístupu vyžadují Databricks Runtime verze 14.2 nebo novější.
- Podpora instancí ARM pro uživatelsky definované funkce Scala v clusterech s podporou Unity Catalog vyžaduje Databricks Runtime 15.2 nebo vyšší.
- Registrace Java UDF ze souboru JAR pomocí
spark.udf.registerJavaFunctionvyžaduje Databricks Runtime 18 LTS nebo novější. Viz Zaregistrujte Java UDF ze souboru JAR.
Important
Sestavte soubor JAR pro stejné verze jazyka Scala a Apache Spark jako výpočetní prostředí, ve kterém běží. Nesoulad může způsobit selhání UDF při registraci nebo volání.
- Klasické výpočetní prostředí: Porovná verze Scala a Sparku vaší verze Databricks Runtime. Viz část Systémové prostředí v dokumentu Poznámky k verzi a kompatibilitě prostředí Databricks Runtime pro vaši verzi. Například Databricks Runtime 18 LTS používá Scala 2.13.16 a Apache Spark 4.0.
- Bezserverové výpočty: Verze jazyka Scala musí odpovídat verzi vašeho prostředí. Vizte Verze prostředí.
Označte závislost Apache Sparku tak provided , aby nebyla součástí souboru JAR. Zahrňte pouze závislosti třetích stran, které vaše uživatelsky definovaná funkce používá.
Registrace funkce jako UDF
Zaregistrujte funkci Scala jako UDF pomocí spark.udf.register:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Zavolejte UDF v Spark SQL
Vytvořte dočasný pohled a potom v dotazu SQL zavolejte funkci definovanou uživatelem (UDF):
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
Použití UDF s DataFrames
Funkci definovanou uživatelem můžete také volat pomocí rozhraní API datového rámce:
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"))
Soubory s UDF
Important
Tato funkce je v beta verzi. Správci pracovního prostoru můžou řídit přístup k této funkci ze stránky Previews . Viz Manage Azure Databricks preview.
Typ souboru ve Scala je .FileRef Použijte jej jako parametr nebo návratový typ v UDF, buď jako typ nejvyšší úrovně, nebo jako vnořený typ. Informace o typu a jeho pravidlech vnoření naleznete v části FILE typ.
Chcete-li přečíst obsah souboru v UDF, zavolejte na FileRef jednu z následujících možností:
-
asLocalFile(): Vrátí objektjava.io.File, který můžete předat libovolné knihovně, která akceptuje cestu. -
open(): Vrací objektjava.io.InputStreampro čtení bajtů ze souboru. Volající to zavře.
Pro vytvoření nového FileRef v UDF zavolejte jednu z následujících statických metod:
-
FileRef.create(uri): Vytvoří odkaz na soubor v umístěníuri. -
FileRef.fromBytes(bytes, destinationPath, contentType): Nahrajebytesdo cesty svazkudestinationPathjako externí soubor a vrátí referenci. -
FileRef.fromLocalFile(localFile, destinationPath, contentType): Nahraje lokální soubor do cestydestinationPathVolume jako externí soubor a vrátí referenci.
Vrácení hodnoty FileRef z funkce UDF, která zapisuje do sloupce FILE MANAGED, není podporováno.
Pro příklady v Python, Scala a SQL, včetně zpracování obrazu, detekce typů souborů a extrakce video snímků, viz Process files with UDF.
Zaregistrujte uživatelsky definovanou funkci jazyka Java ze souboru JAR
Zabalte UDF jako soubor JAR, přidejte ho do své relace pomocí spark.addArtifact a zaregistrujte třídu UDF pomocí spark.udf.registerJavaFunction.
Poznámka:
Podporováno ve standardním režimu přístupu a v bezserverových výpočetních prostředcích v Databricks Runtime 18 LTS a vyšší. Registrovaná funkce je vymezená relací a není zaregistrovaná v katalogu Unity.
Následující kroky vás provedou vytvořením projektu, napsáním třídy UDF definované uživatelem, sestavením tlustého souboru JAR a jeho registrací.
Krok 1: Vytvoření projektu
Nastavte projekt v jazyce Scala nebo Java.
Scala
Vytvořte nový projekt Scala pomocí sbt:
sbt new scala/scala-seed.g8
Obsah souboru nahraďte build.sbt následujícím kódem. Nastavte scalaVersion a verzi spark-sql tak, aby odpovídaly vašemu výpočetnímu prostředí:
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"
)
Povolte plugin sbt-assembly, abyste mohli sestavit fat JAR. Vytvořte nebo upravte project/assembly.sbt a přidejte:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Vytvořte nový projekt Maven pomocí archetypu rychlého startu:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Tento příkaz vytvoří standardní strukturu projektu Maven s src/main/java adresáři a src/test/java adresáři.
Do vygenerovaného pom.xml uvnitř značek <project></project> přidejte blok <properties> a nakonfigurujte maven-shade-plugin tak, aby sestavilo 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: Napište svou třídu UDF
Vaše třída UDF musí implementovat jedno z rozhraní org.apache.spark.sql.api.java.UDF (UDF1 až UDF22), kde číslo označuje, kolik vstupních argumentů UDF přijímá. Implementujte metodu call() pomocí vaší logiky.
Obslužná třída musí být třídou jazyka Java.
spark.udf.registerJavaFunction načítá třídu pomocí reflexe, takže musí jít o veřejnou třídu nejvyšší úrovně (nebo veřejnou vnořenou třídu static) s veřejným bezparametrovým konstruktorem. Scala class nebo object nesplňuje tento požadavek a v době volání selže. Soubor JAR můžete sestavit pomocí sbt, ale samotná třída UDF musí být napsána v Javě.
Vytvořit 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: Sestavit fat JAR soubor
Zabalte zkompilovanou uživatelsky definovanou funkci (UDF) do fat JARu.
Scala
V kořenovém adresáři projektu spusťte:
sbt clean assembly
Soubor fat JAR je vytvořen v target/scala-2.13/ s názvem, například my-udf-assembly-0.1.0-SNAPSHOT.jar.
Java
V kořenovém adresáři projektu spusťte:
mvn clean package
Soubor fat JAR je vytvořen v target/ s názvem, například my-udf-1.0-SNAPSHOT.jar.
Krok 4: Nahrání souboru JAR do svazku katalogu Unity
Nahrajte soubor JAR do svazku v Unity Catalogu, aby k němu mělo vaše výpočetní prostředí přístup. Pokud ještě žádný diskový svazek nemáte, vytvořte ho:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Nahrajte soubor JAR do svazku pomocí Průzkumníka katalogu:
- V pracovním prostoru Azure Databricks klikněte na
Klikněte na Katalog pro otevření Průzkumníka katalogu.
- Vyberte katalog a pak vyberte schéma, které obsahuje váš svazek.
- Klikněte na název svazku.
- Klikněte na Nahrát do tohoto svazku a vyberte svůj soubor JAR.
- Klikněte na tlačítko Odeslat.
- Po dokončení nahrávání klikněte na název souboru JAR a potom klikněte na Kopírovat cestu, čímž zkopírujete cestu ke svazku. Například:
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Tuto cestu potřebujete v dalším kroku.
Krok 5: Zaregistrujte a zavolejte funkci definovanou uživatelem (UDF)
Přidejte soubor JAR do relace pomocí cesty ke svazku, zaregistrujte třídu UDF a volejte ji v jazyce 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()
V bezserverovém a standardním režimu přístupu musíte předat explicitní návratový typ. Vynechání návratového typu selže s chybou UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Uživatelsky definované agregační funkce (UDAFy) nejsou v registerJavaFunction podporovány.
Dotaz vrátí výstup UDF, který potvrzuje, že je funkce zaregistrovaná a volatelná:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Kontrola pořadí vyhodnocení a hodnoty null
Spark SQL (včetně rozhraní SQL a rozhraní API datových rámců a datových sad) nezaručuje pořadí vyhodnocení dílčího výrazu. Spark nevyhodnocuje vstupy operátoru nebo funkce zleva doprava. Logické AND a OR výrazy nemají sémantiku zkratování zleva doprava.
Nespoléhejte na vedlejší efekty ani na pořadí vyhodnocování logických výrazů ani na pořadí klauzulí WHERE a HAVING. Optimalizátor dotazů může změnit pořadí těchto výrazů a klauzulí. Pokud UDF spoléhá při kontrole na hodnotu null na semantiku zkráceného vyhodnocování, Spark nezaručuje, že se kontrola hodnoty null provede před spuštěním UDF. Příklady:
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
Tato WHERE klauzule nezaručuje, že Spark spustí UDF strlen poté, co odfiltruje hodnoty null.
Pro zpracování kontroly null doporučuje Databricks některou z následujících možností:
- Upravte samotnou UDF tak, aby pracovala s hodnotou null, a kontrolu na hodnotu null provádějte uvnitř UDF.
- Použijte
IFneboCASE WHENvýrazy k ověření hodnoty null a pro vyvolání funkce definované uživatelem v podmíněné větvi.
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
Nativně typované API pro datové sady
Poznámka:
Tato funkce je podporována v clusterech s podporou katalogu Unity se standardním režimem přístupu v Databricks Runtime 15.4 a novějším.
Pomocí typed Dataset API můžete spouštět transformace, jako je mapování, filtrování a agregace u datových sad s uživatelem definovanou funkcí.
Následující příklad používá rozhraní API map() ke změně čísla ve výsledném sloupci na řetězec s předponou:
spark.range(3).map(f => s"row-$f").show()
Tento příklad používá map(), ale stejný vzor lze použít i pro další rozhraní API pro typované datové sady, například filter(), mapPartitions(), foreach(), foreachPartition(), reduce() a flatMap().
Funkce UDF Scala a kompatibilita Databricks Runtime
Následující funkce vyžadují minimální verze prostředí Databricks Runtime v clusterech s povoleným katalogem Unity ve standardním (sdíleném) režimu přístupu.
| Vlastnost | Minimální verze Databricks Runtime |
|---|---|
| Skalární funkce definované uživatelem | 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 |
(Streamování) foreachWriter Sink |
Databricks Runtime 15.4 |
(Streamování) foreachBatch |
Databricks Runtime 16.1 |
(Streamování) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks Runtime 16.2 |
spark.udf.registerJavaFunction (Java UDF ze souboru JAR) |
Databricks Runtime 18 LTS |