Uživatelsky definované funkce v jazycích Scala a Java s platností pro relaci

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.registerJavaFunction vyž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í objekt java.io.File, který můžete předat libovolné knihovně, která akceptuje cestu.
  • open(): Vrací objekt java.io.InputStream pro č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): Nahraje bytes do cesty svazku destinationPath jako externí soubor a vrátí referenci.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Nahraje lokální soubor do cesty destinationPath Volume 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 (UDF1UDF22), 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:

  1. V pracovním prostoru Azure Databricks klikněte na ikonu Data. Klikněte na Katalog pro otevření Průzkumníka katalogu.
  2. Vyberte katalog a pak vyberte schéma, které obsahuje váš svazek.
  3. Klikněte na název svazku.
  4. Klikněte na Nahrát do tohoto svazku a vyberte svůj soubor JAR.
  5. Klikněte na tlačítko Odeslat.
  6. 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 IF nebo CASE WHEN vý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