Munkamenetszintű Scala- és Java UDF-ek

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.registerJavaFunction haszná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 a java.io.File fájlt, amit bármely könyvtárnak át tudsz adni, amely elfogad egy útvonalat.
  • open(): A fájl bájtjainak olvasásához egy java.io.InputStream objektumot 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) bytes elemet külső fájlként feltölti a destinationPath Volume elérési útjára, és visszaad egy hivatkozást.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Helyi fájlt tölt fel a destinationPath Volume-ú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:

  1. A Azure Databricks munkaterületen kattintson a Data icon.Catalog elemre a Catalog Explorer megnyitásához.
  2. Jelölje ki a katalógust, majd válassza ki a kötetet tartalmazó sémát.
  3. Kattintson a kötet nevére.
  4. Kattintson a Feltöltés erre a kötetre lehetőségre, és válassza ki a JAR-fájlt.
  5. Kattintson a Feltöltés gombra.
  6. 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
  • IF vagy CASE WHEN kifejezé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