Co jsou uživatelem definované funkce (UDF)?

Uživatelem definované funkce (UDF) umožňují opakovaně používat a sdílet kód, který rozšiřuje integrované funkce v Azure Databricks. Funkce definované uživatelem (UDFs) slouží k provádění konkrétních úloh, jako jsou složité výpočty, transformace nebo vlastní manipulace s daty.

Kdy použít funkci UDF vs. Apache Spark?

UDF používejte pro logiku, která se obtížně vyjadřuje pomocí integrovaných funkcí Apache Sparku. Integrované funkce Apache Sparku jsou optimalizované pro distribuované zpracování a nabízejí lepší výkon ve velkém měřítku. Další informace naleznete v tématu Functions.

Databricks doporučuje uživatelsky definované funkce (UDF) pro ad hoc dotazy, ruční čištění dat, explorační analýzu dat a operace nad malými až středně velkými datovými sadami. Mezi běžné případy použití uživatelsky definovaných funkcí patří šifrování dat, dešifrování, hashování, analýza JSON a ověřování.

Používejte metody Apache Spark pro operace s velkými datovými sadami a všechny úlohy, které běží pravidelně nebo nepřetržitě, včetně úloh ETL a operací streamování.

Porozumět typům uživatelem definovaným funkcím

Na následujících kartách vyberte typ UDF pro zobrazení popisu, příkladu a odkazu pro více informací.

Skalární UDF

Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek. Mohou být řízeny pomocí katalogu Unity nebo zaměřeny na relaci.

Následující příklad používá skalární UDF k výpočtu délky každého názvu ve name sloupci a přidání hodnoty do nového sloupce name_length.

+-------+-------+
| name  | score |
+-------+-------+
| alice |  10.0 |
| bob   |  20.0 |
| carol |  30.0 |
| dave  |  40.0 |
| eve   |  50.0 |
+-------+-------+
-- Create a SQL UDF for name length
CREATE OR REPLACE FUNCTION main.test.get_name_length(name STRING)
RETURNS INT
RETURN LENGTH(name);

-- Use the UDF in a SQL query
SELECT name, main.test.get_name_length(name) AS name_length
FROM your_table;
+-------+-------+-------------+
| name  | score | name_length |
+-------+-------+-------------+
| alice |  10.0 |      5      |
|  bob  |  20.0 |      3      |
| carol |  30.0 |      5      |
| dave  |  40.0 |      4      |
|  eve  |  50.0 |      3      |
+-------+-------+-------------+

Implementace tohoto v notebooku Azure Databricks s využitím PySparku:

from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType

@udf(returnType=IntegerType())
def get_name_length(name):
  return len(name)

df = df.withColumn("name_length", get_name_length(df.name))

# Show the result
display(df)

Viz SQL a Python uživatelem definované funkce (UDF) v katalogu Unity a Python skalární uživatelem definované funkce (UDF)

Skalární uživatele služby Batch

Zpracovávání dat v dávkách se zachováním rovnosti poměru 1:1 mezi vstupními a výstupními řádky. Tím se sníží náročnost operací zpracovávaných po jednotlivých řádcích pro zpracování velkých objemů dat. Uživatelsky definované funkce Batch také udržují stav mezi dávkovými zpracováními, aby se spouštěly efektivněji, opakovaně používaly prostředky a zpracovaly složité výpočty, které potřebují kontext napříč bloky.

Mohou být řízeny pomocí katalogu Unity nebo zaměřeny na relaci.

Následující funkce Batch Unity Catalog Python UDF vypočítá BMI při zpracování šarží řádků.

+-------------+-------------+
| weight_kg   | height_m    |
+-------------+-------------+
|      90     |     1.8     |
|      77     |     1.6     |
|      50     |     1.5     |
+-------------+-------------+
%sql
CREATE OR REPLACE FUNCTION main.test.calculate_bmi_pandas(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
AS $$
import pandas as pd
from typing import Iterator, Tuple

def handler_function(batch_iter: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
  for weight_series, height_series in batch_iter:
    yield weight_series / (height_series ** 2)
$$;

select main.test.calculate_bmi_pandas(cast(70 as double), cast(1.8 as double));
+--------+
|  BMI   |
+--------+
|  27.8  |
|  30.1  |
|  22.2  |
+--------+

Viz SQL a Python uživatelem definované funkce (UDF) v katalogu Unity a Batch Python uživatelem definované funkce (UDF) v katalogu Unity.

Jiné než skalární funkce definované uživatelem

Nes skalární funkce definované uživatelem pracují s celými datovými sadami a sloupci s flexibilními vstupními a výstupními poměry (1:N nebo N:N).

Uživatelem definované uživatelem knihovny pandas v rámci relace můžou být následující typy:

  • Serie na serii
  • Iterátor řady k iterátoru řad
  • Iterátor více sérií na iterátor série
  • Ze série na skalár

Následuje příklad pandas UDF typu Series na Series.

from pyspark.sql.functions import pandas_udf
import pandas as pd

df = spark.createDataFrame([(70, 1.75), (80, 1.80), (60, 1.65)], ["Weight", "Height"])

@pandas_udf("double")
def calculate_bmi_pandas(weight: pd.Series, height: pd.Series) -> pd.Series:
  return weight / (height ** 2)

df.withColumn("BMI", calculate_bmi_pandas(df["Weight"], df["Height"])).display()

Vizte funkce pandas definované uživatelem.

UDAF

UDAF pracují s více řádky a vrací jeden agregovaný výsledek. Funkce UDAFs jsou omezeny pouze na relaci.

Následující příklad UDAF agreguje skóre podle délky názvu.

from pyspark.sql.functions import pandas_udf
from pyspark.sql import SparkSession
import pandas as pd

# Define a pandas UDF for aggregating scores
@pandas_udf("int")
def total_score_udf(scores: pd.Series) -> int:
  return scores.sum()

# Group by name length and aggregate
result_df = (df.groupBy("name_length")
  .agg(total_score_udf(df["score"]).alias("total_score")))

display(result_df)
+-------------+-------------+
| name_length | total_score |
+-------------+-------------+
|      3      |     70.0    |
|      4      |     40.0    |
|      5      |     40.0    |
+-------------+-------------+

Viz uživatelsky definované funkce knihovny pandas pro Python a uživatelsky definované agregační funkce (UDAF) v jazyce Scala.

UDTF

UDTF přebírá jeden nebo více vstupních argumentů a vrací více řádků (a případně více sloupců) pro každý vstupní řádek. Mohou být řízeny pomocí katalogu Unity nebo zaměřeny na relaci.

Následující funkce UDTF vytvoří tabulku s pevným seznamem dvou celých argumentů:

CREATE OR REPLACE FUNCTION get_sum_diff(x INT, y INT)
RETURNS TABLE (sum INT, diff INT)
LANGUAGE PYTHON
HANDLER 'GetSumDiff'
AS $$
class GetSumDiff:
    def eval(self, x: int, y: int):
        yield x + y, x - y
$$;

SELECT * FROM get_sum_diff(10, 3);
+-----+------+
| sum | diff |
+-----+------+
| 13  | 7    |
+-----+------+

Implementace tohoto v notebooku Azure Databricks s využitím PySparku:

from pyspark.sql.functions import lit, udtf

@udtf(returnType="sum: int, diff: int")
class GetSumDiff:
    def eval(self, x: int, y: int):
        yield x + y, x - y

GetSumDiff(lit(1), lit(2)).show()

Viz UDTF katalogu Unity a uživatelsky definované funkce v rámci relace.

UDF spravované v Unity Catalogu vs. UDF s rozsahem relace

Unity Catalog uchovává UDF spravované katalogem Unity, což zlepšuje správu, opakované použití a dohledatelnost. Funkce definované uživatelem (UDF) definujete v poznámkovém bloku nebo úloze a jejich rozsah je omezen na aktuální SparkSession. Uživatelsky definované funkce (UDF) s platností v rámci relace můžete definovat a používat pomocí SQL, Pythonu nebo Scaly.

Pomocí následující tabulky určete, kterou z obou kategorií zvolit, a potom si projděte stručné přehledy konkrétních typů UDF v každé z nich.

Consideration Funkce definované uživatelem spravované v katalogu Unity UDF s rozsahem relace
Nejvhodnější pro Bezpečné sdílení funkcí napříč týmy, poznámkovými bloky, úlohami a službami SQL Warehouse Rychlý iterativní vývoj v rámci jednoho poznámkového bloku nebo úlohy
Languages SQL, Python, Scala a Java. SQL, Python a Scala.
Zásady správného řízení a sdílení Řídí se oprávněními katalogu Unity a zjistitelnými v Průzkumníku katalogu. Vymezený na aktuální SparkSession. Neřídí se ani nesdílí.
Perzistence Trvalé v katalogu Unity a opakovaně použitelné napříč relacemi. Existuje pouze pro aktuální relaci.

Stručný průvodce k řídícím uživatelem definovaným funkcím v Katalogu Unity

UDF, které se řídí službou Unity Catalog, umožňují definovat, používat, bezpečně sdílet a řídit vlastní funkce napříč výpočetními prostředími. Viz SQL a Python uživatelem definované funkce (UDF) v katalogu Unity.

Typ UDF Podporované výpočetní prostředky Popis
Uživatelská definovaná funkce Python v katalogu Unity
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředky se standardním režimem přístupu (Databricks Runtime 13.3 LTS a vyšší)
  • SQL Warehouse (bez serveru a pro)
  • Kanály Lakeflow (klasické i bezserverové)
Definujte uživatelem definované funkce v Pythonu a zaregistrujte ho v katalogu Unity pro správu.
Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek.
Dávkový katalog Unity pro Python UDF
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředí s standardním režimem přístupu (Databricks Runtime 16.3 a novější)
  • SQL Warehouse (bez serveru a pro)
Definujte uživatelem definované funkce v Pythonu a zaregistrujte ho v katalogu Unity pro správu.
Dávkové operace s více hodnotami a vrací více hodnot. Snižuje náklady spojené s operacemi po řádcích při zpracovávání velkých objemů dat.
Katalog Unity Python UDTF
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředky se standardním režimem přístupu (Databricks Runtime 17.1 a novější)
  • SQL Warehouse (bez serveru a pro)
Definujte UDTF v Pythonu a zaregistrujte ho pro správu v katalogu Unity.
UDTF přebírá jeden nebo více vstupních argumentů a vrací více řádků (a případně více sloupců) pro každý vstupní řádek.
Unity Catalog Scala nebo Java UDF
  • Bezserverové poznámkové bloky a úlohy
  • Klasický výpočetní výkon (standardní přístup a vyhrazený režim přístupu)
  • SQL Warehouses (bezserverové, pro a klasické)
  • Deklarativní kanály Sparku v Lakeflow (klasické i bezserverové)
Definujte UDF v jazyce Scala nebo Java a zaregistrujte ho v katalogu Unity pro zásady správného řízení.
Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek. Vyžaduje Scala 2.13.16, JDK 17 a prostředí verze 4.

Stručná příručka k UDF s omezeným oborem relace pro uživatelsky izolovaný výpočet

Funkce definované uživatelem (UDF) definujete v poznámkovém bloku nebo úloze a jejich rozsah je omezen na aktuální SparkSession. Uživatelsky definované funkce (UDF) s platností v rámci relace můžete definovat a používat pomocí SQL, Pythonu nebo Scaly.

Typ UDF Podporované výpočetní prostředky Popis
Skalární jazyk Python
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředky se standardním režimem přístupu (Databricks Runtime 13.3 LTS a vyšší)
  • Kanály Lakeflow (klasické i bezserverové)
Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek.
Neskalární Python
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředky se standardním režimem přístupu (Databricks Runtime 14.3 LTS a vyšší)
  • Kanály Lakeflow (klasické i bezserverové)
Mezi ne skalární funkce definované uživatelem patří pandas_udf, mapInPandas, mapInArrow, applyInPandas. Funkce Pandas UDFs používají Apache Arrow k přenosu dat a využívají knihovnu pandas pro práci s těmito daty. Funkce pandas UDF podporují vektorizované operace, které mohou výrazně zvýšit výkon oproti skalárním UDF zpracovávaným řádek po řádku.
Definované funkce Pythonu
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředky se standardním režimem přístupu (Databricks Runtime 14.3 LTS a vyšší)
  • Kanály Lakeflow (klasické i bezserverové)
UDTF přebírá jeden nebo více vstupních argumentů a vrací více řádků (a případně více sloupců) pro každý vstupní řádek.
Skalární uživatelem definované funkce Scala
  • Klasické výpočetní prostředky se standardním režimem přístupu (Databricks Runtime 13.3 LTS a vyšší)
Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek.
Scala nebo Java UDF z JAR
  • Bezserverové poznámkové bloky a úlohy
  • Klasické výpočetní prostředí s standardním režimem přístupu (Databricks Runtime 18 LTS a vyšší)
Zaregistrujte předkompilovanou třídu UDF ze souboru JAR pomocí spark.udf.registerJavaFunction. Viz Zaregistrujte Java UDF ze souboru JAR.
uživatelsky definované agregační funkce (UDAFs) ve Scale
  • Klasické výpočetní prostředí s vyhrazeným režimem přístupu (Databricks Runtime 14.2 LTS a vyšší)
UDAF pracují s více řádky a vrací jeden agregovaný výsledek.

Důležité informace o výkonu

  • Vestavěné funkce a SQL UDF jsou nejúčinnějšími možnostmi.

  • Scala UDF jsou obecně rychlejší než Python UDF.

    • Unisolated Scala UDF běží na virtuálním počítači Java Virtual Machine (JVM), aby se vyhnuli režijním nákladům na přesun dat do a z prostředí JVM.
    • Izolované UDF ve Scale musí přesouvat data do JVM a z JVM, ale přesto mohou být rychlejší než UDF v Pythonu, protože pracují s pamětí efektivněji.
  • Uživatelsky definované funkce Pythonu a uživatelsky definované funkce pandas bývají pomalejší než uživatelsky definované funkce ve Scale, protože je nutné serializovat data a přesouvat je z JVM do interpretru Pythonu.

    • Pandas UDFs jsou až 100x rychlejší než Python UDFs, protože používají Apache Arrow, aby snížily náklady na serializaci.