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.
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 |
|
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 |
|
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 |
|
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 |
|
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 |
|
Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek. |
| Neskalární Python |
|
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 |
|
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 |
|
Skalární UDF pracují na jednom řádku a vracejí jednu výslednou hodnotu pro každý řádek. |
| Scala nebo Java UDF z JAR |
|
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 |
|
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.