Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Important
A Scala és Java UDF-ek regisztrálhatók a Unity Katalógusban az irányítás, az újrafelhasználás és a felderíthetőség érdekében. Lásd: Scala és Java felhasználó által definiált függvények (UDF-ek) a Unity Catalogban.
Ez a lap bemutatja, hogyan hozhat létre munkamenet-hatókörű Scalát és Java UDF-eket Azure Databricks. A munkamenet-hatókörű UDF-ek egy jegyzetfüzetben vagy feladatban vannak definiálva, és csak az aktuális SparkSession-ra vonatkoznak. Az SQL nyelvi referenciáját a külső, felhasználó által definiált skaláris függvények (UDF-ek) című témakörben talál.
A megközelítés kiválasztása
A Scala vagy Java UDF az alábbi módokon határozható meg. Az összes UDF-típus nyelv, irányítás és számítás összehasonlításához tekintse meg a Unity Catalog által szabályozott és a munkamenet-hatókörű UDF-eket.
| Approach | Leírás |
|---|---|
| Beágyazott Scala UDF | Definiáljon egy UDF-et egy jegyzetfüzetben Scala-függvény vagy lambda használatával. Munkamenetre korlátozott. A kiszolgáló nélküli számítás nem támogatott. |
| Java UDF JAR-fájlból | Regisztráljon egy előre lefordított UDF-osztályt egy JAR-fájlból a(z) spark.udf.registerJavaFunction használatával. Munkamenetre korlátozott. Kiszolgáló nélküli számítással támogatott. |
| Unity Catalog által szabályozott Scala vagy Java UDF | Regisztráljon egy UDF-t a Unity Katalógusban a szabályozáshoz, az újrafelhasználáshoz és a felderíthetőséghez. Kiszolgáló nélküli számítással támogatott. |
Requirements
- A Unity catalog-kompatibilis, standard hozzáférési móddal rendelkező Scala UDF-ekhez a Databricks Runtime 14.2-s vagy újabb verziójára van szükség.
- A Scala UDF-ek ARM-példányainak támogatása Unity katalógusbarát fürtökön a Databricks Runtime 15.2-s vagy újabb verziójára van szükség.
- A Java UDF JAR-fájlból történő regisztrálásához a(z)
spark.udf.registerJavaFunctionhasználatával Databricks Runtime 18 LTS vagy újabb verzió szükséges. Lásd: Java UDF regisztrálása JAR-ból.
Important
A JAR-t ugyanazokkal a Scala- és Apache Spark-verziókkal készítheti el, mint az azt futtató számítás. Az eltérés miatt az UDF regisztráláskor vagy meghíváskor sikertelen lehet.
- Klasszikus számítás: Egyezik a Databricks Runtime-verzió Scala és Spark verziójával. Tekintse meg a Databricks Runtime kiadási megjegyzésverzióinakRendszerkörnyezet című szakaszát, és tekintse meg a verzió kompatibilitását. A Databricks Runtime 18 LTS például a Scala 2.13.16-ot és az Apache Spark 4.0-t használja.
- Kiszolgáló nélküli számítás: Egyezik a környezeti verzió Scala-verziójával. Lásd: Környezeti verziók.
Jelölje meg az Apache Spark-függőséget provided úgy, hogy az ne legyen a JAR-ba csomagolva. Csak az UDF által használt külső függőségeket foglalja magában.
Függvény regisztrálása UDF-ként
Scala-függvény regisztrálása UDF-ként a(z) spark.udf.register használatával:
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
Az UDF meghívása a Spark SQL-ben
Hozzon létre egy ideiglenes nézetet, majd hívja meg az UDF-et egy SQL-lekérdezésben:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
UDF használata DataFrame-ekkel
UDF-et a DataFrame API használatával is meghívhat:
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"))
Fájlok UDF-hez
Important
Ez a funkció bétaverzióban érhető el. A munkaterület rendszergazdái az Előnézetek lapon szabályozhatják a funkcióhoz való hozzáférést. Lásd: Az Azure Databricks előzetes verziójának kezelése.
A fájl Scala típusa .FileRef Használd paraméterként vagy visszaküldési típusként egy UDF-ben, akár felső szintű típusként, akár fészkelveként. A típusról és annak beágyazási szabályairól lásd: FILE type.
Egy fájl tartalmának UDF-ben történő beolvasásához hívja meg az alábbi metódusok egyikét egy FileRef objektumon:
-
asLocalFile(): Visszaadja ajava.io.Filefájlt, amit bármely könyvtárnak át tudsz adni, amely elfogad egy útvonalat. -
open(): A fájl bájtjainak olvasásához egyjava.io.InputStreamobjektumot ad vissza. A hívó bezárja.
Új FileRef előállításához egy UDF-ben hívja meg az alábbi statikus metódusok egyikét:
-
FileRef.create(uri): Hivatkozást hoz létre a következő helyen található fájlra:uri. -
FileRef.fromBytes(bytes, destinationPath, contentType): A(z)byteselemet külső fájlként feltölti adestinationPathVolume elérési útjára, és visszaad egy hivatkozást. -
FileRef.fromLocalFile(localFile, destinationPath, contentType): Helyi fájlt tölt fel adestinationPathVolume-útra külső fájlként, és visszaad egy hivatkozást.
A FileRef visszaadása olyan UDF-ből, amely egy FILE MANAGED oszlopba ír, nem támogatott.
A Python, Scala és SQL példáihoz, beleértve a képfeldolgozást, fájltípus-felismerést és videóképkocka kibontását, lásd: Processzfájlok UDF-ekkel.
Java UDF regisztrálása JAR-ból
Csomagolja a UDF-et JAR-fájlba, adja hozzá a munkamenetéhez a(z) spark.addArtifact használatával, és regisztrálja a UDF-osztályt a(z) spark.udf.registerJavaFunction használatával.
Megjegyzés
A Databricks Runtime 18 LTS vagy újabb verziójában a standard hozzáférési mód és a kiszolgáló nélküli számítás támogatott. A regisztrált függvény munkamenet-hatókörű, és nincs regisztrálva a Unity Katalógusban.
Az alábbi lépések végigvezetik a projekt létrehozásán, egy UDF-osztály megírásán, egy fat JAR elkészítésén és annak regisztrálásán.
1. lépés: A projekt létrehozása
Projekt beállítása a Scalában vagy Java.
Scala
Hozzon létre egy új Scala-projektet a következővel sbt:
sbt new scala/scala-seed.g8
Cserélje le a build.sbt fájl tartalmát a következőre. Állítsa be scalaVersion és a spark-sql verziót a számításnak megfelelően:
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"
)
Engedélyezze az sbt-assembly beépülő modult egy kövér JAR létrehozásához. Hozza létre vagy szerkessze a(z) project/assembly.sbt elemet, és adja hozzá:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Java
Hozzon létre egy új Maven-projektet a quickstart archetípus használatával:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
Ez a parancs létrehozza a szabványos Maven-projektstruktúrát a src/main/java- és src/test/java könyvtárakkal.
A létrehozott pom.xml fájlban, a <project></project> címkék között adjon hozzá egy <properties> blokkot, és konfigurálja úgy a maven-shade-plugin elemet, hogy fat JAR-t építsen:
<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>
2. lépés: Az UDF-osztály írása
Az UDF-osztálynak implementálnia kell az org.apache.spark.sql.api.java.UDF egyik interfészt (UDF1 ezen keresztül UDF22), ahol a szám azt jelzi, hogy az UDF hány bemeneti argumentumot vesz fel. Valósítsa meg a call() metódust a saját logikája szerint.
A kezelőnek Java osztálynak kell lennie.
spark.udf.registerJavaFunction az osztályt tükröződés alapján tölti be, ezért egy legfelső szintű (vagy static beágyazott) nyilvános osztálynak kell lennie egy nyilvános no-arg konstruktorral. A Scala class vagy object nem felel meg ennek a követelménynek, és híváskor meghiúsul. Az sbt-vel elkészítheti a JAR-fájlt, de magát az UDF-osztályt Java nyelven kell megírni.
Hozza létre 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;
}
}
3. lépés: Önálló JAR létrehozása
Csomagolja a lefordított UDF-et egy fat JAR-ba.
Scala
A projektje gyökérkönyvtárában futtassa a következő parancsot:
sbt clean assembly
A fat JAR a(z) target/scala-2.13/ helyen jön létre, my-udf-assembly-0.1.0-SNAPSHOT.jar névhez hasonló néven.
Java
A projektje gyökérkönyvtárában futtassa a következő parancsot:
mvn clean package
A fat JAR a(z) target/ helyen jön létre, my-udf-1.0-SNAPSHOT.jar névhez hasonló néven.
4. lépés: JAR feltöltése Unity-katalóguskötetbe
Töltse fel a JAR-t egy Unity-katalóguskötetbe , hogy a számítása hozzáférhessen hozzá. Ha még nincs kötete, hozzon létre egyet:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Töltse fel JAR-fájlját a kötetre a Catalog Explorer használatával:
- A Azure Databricks munkaterületen kattintson a
Catalog elemre a Catalog Explorer megnyitásához.
- Jelölje ki a katalógust, majd válassza ki a kötetet tartalmazó sémát.
- Kattintson a kötet nevére.
- Kattintson a Feltöltés erre a kötetre lehetőségre, és válassza ki a JAR-fájlt.
- Kattintson a Feltöltés gombra.
- A feltöltés befejezése után kattintson a JAR-fájl nevére, majd kattintson az Elérési út másolása lehetőségre a kötet elérési útjának másolásához. Például:
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. A következő lépésben szüksége lesz erre az útvonalra.
5. lépés: A UDF regisztrálása és meghívása
Adja hozzá a JAR-fájlt a munkamenethez a kötet elérési útjának használatával, regisztrálja az UDF-osztályt, majd hívja meg a Spark SQL-ből:
# 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()
Kiszolgáló nélküli és standard hozzáférési módú számítás esetén explicit visszatérési típust kell megadnia. A visszatérési típus elhagyása UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE hibát eredményez. A felhasználó által definiált összesítő függvények (UDAF-ek) nem támogatottak.registerJavaFunction
A lekérdezés az UDF kimenetét adja vissza, megerősítve, hogy a függvény regisztrálva van és meghívható:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
Kiértékelési sorrend és nulla érték ellenőrzése
A Spark SQL (beleértve az SQL-t és a DataFrame- és Adathalmaz API-kat) nem garantálja a szubexpressziós kiértékelés sorrendjét. A Spark nem értékeli ki egy operátor vagy függvény bemeneteit balról jobbra. A logikai AND és OR kifejezések nem rendelkeznek balról jobbra rövidzárolású szemantikával.
Ne támaszkodjon a logikai kifejezések mellékhatásaira vagy kiértékelési sorrendjére, illetve a záradékok sorrendjére WHEREHAVING . A lekérdezésoptimalizáló átrendezheti ezeket a kifejezéseket és záradékokat. Ha egy UDF a null ellenőrzéshez rövidzárlatú szemantikára támaszkodik, a Spark nem garantálja, hogy a null ellenőrzés az UDF előtt fut. Például:
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
Ez a WHERE záradék nem garantálja, hogy a Spark a null értékek kiszűrése után meghívja az strlen UDF-et.
A null értékű ellenőrzés kezeléséhez a Databricks az alábbiak valamelyikét javasolja:
- Állítsa az UDF-et nullérzékenysé, és végezze el a null ellenőrzést az UDF-ben
-
IFvagyCASE WHENkifejezések használata a null-ellenőrzéshez és az UDF meghívása feltételes ágban
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
Tipizált adatbázis API-k
Megjegyzés
Ez a funkció a Unity Catalog-kompatibilis fürtök esetében támogatott standard hozzáférési móddal a Databricks Runtime 15.4-ben és újabb verziókban.
A gépelt Adathalmaz API-k segítségével olyan átalakításokat futtathat, mint a leképezés, a szűrés és az összesítések egy felhasználó által definiált függvénnyel rendelkező adathalmazokon.
Az alábbi példa az map() API-val módosít egy eredményt tartalmazó oszlopban lévő számot egy előtaggal rendelkező sztringre:
spark.range(3).map(f => s"row-$f").show()
Ez a példa használjamap(), de ugyanez a minta vonatkozik más típusú adathalmaz API-kra, például filter(), , mapPartitions()foreach(), foreachPartition(), reduce()és flatMap().
A Scala UDF szolgáltatásai és a Databricks Runtime kompatibilitása
Az alábbi funkciók a standard (megosztott) hozzáférési módú, Unity Catalogot használó fürtök esetén a Databricks Runtime minimális verzióit igénylik.
| Tulajdonság | A Databricks runtime minimális verziója |
|---|---|
| Skaláris UDF-ek | Databricks Runtime 14.2 |
Dataset.map, Dataset.mapPartitions, Dataset.filter, Dataset.reduceDataset.flatMap |
Databricks Runtime 15.4 verzió |
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups |
Databricks Runtime 15.4 verzió |
(Streamelés) foreachWriter Sink |
Databricks Runtime 15.4 verzió |
(Streamelés) foreachBatch |
Databricks Runtime 16.1 |
(Streamelés) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks Runtime 16.2 |
spark.udf.registerJavaFunction (Java UDF JAR-fájlból) |
Databricks Runtime 18 LTS |