Co to są funkcje zdefiniowane przez użytkownika (UDF)?

Funkcje zdefiniowane przez użytkownika (UDF) umożliwiają ponowne używanie i udostępnianie kodu, który rozszerza wbudowane możliwości na Azure Databricks. UDF-y pozwalają na wykonanie określonych zadań, takich jak złożone obliczenia, przekształcenia lub niestandardowe manipulacje danymi.

Kiedy używać funkcji UDF a Apache Spark?

Używaj funkcji zdefiniowanych przez użytkownika dla logiki, która jest trudna do wyrażenia za pomocą wbudowanych funkcji platformy Apache Spark. Wbudowane funkcje platformy Apache Spark są zoptymalizowane pod kątem przetwarzania rozproszonego i zapewniają lepszą wydajność na dużą skalę. Aby uzyskać więcej informacji, zobacz Functions.

Databricks zaleca funkcje UDF do zapytań ad hoc, ręcznego czyszczenia danych, eksploracji danych oraz operacji na małych i średnich zbiorach danych. Typowe przypadki użycia funkcji użytkownika obejmują szyfrowanie danych, odszyfrowywanie, haszowanie, analizowanie JSON i walidację.

Użyj metod platformy Apache Spark na potrzeby operacji na dużych zestawach danych i wszystkich obciążeniach, które są uruchamiane regularnie lub stale, w tym zadania ETL i operacje przesyłania strumieniowego.

Zrozumienie typów UDF

Wybierz typ funkcji zdefiniowanej przez użytkownika z poniższych zakładek, aby zobaczyć opis, przykład oraz link do dalszych informacji.

Skalarna UDF

Skalarne funkcje zdefiniowane przez użytkownika działają dla pojedynczego wiersza i zwracają pojedynczą wartość dla każdego wiersza. Mogą być zarządzane przez Unity Catalog lub obejmować zakres sesji.

W poniższym przykładzie użyto skalarnej funkcji UDF do obliczenia długości każdej nazwy w kolumnie name i dodania wartości w nowej kolumnie 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      |
+-------+-------+-------------+

Aby zaimplementować to w notatniku usługi Azure Databricks przy użyciu biblioteki PySpark:

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)

Zobacz funkcje użytkownika SQL i Python (UDF) w Unity Catalog oraz skalarne funkcje użytkownika w języku Python (UDF).

Funkcje zdefiniowane przez użytkownika wsadowe

Przetwarzanie danych w partiach przy zachowaniu parzystości wierszy wejściowych/wyjściowych 1:1. Zmniejsza to obciążenie operacji wiersz po wierszu na potrzeby przetwarzania danych na dużą skalę. Funkcje UDF usługi Batch również utrzymują stan między partiami, aby działać bardziej efektywnie, ponownie wykorzystywać zasoby i obsługiwać złożone obliczenia, które wymagają odniesienia między fragmentami danych.

Mogą być zarządzane przez Unity Catalog lub obejmować zakres sesji.

Następująca funkcja UDF Unity Catalog Batch w języku Python oblicza BMI podczas przetwarzania partii wierszy.

+-------------+-------------+
| 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  |
+--------+

Zobacz funkcje SQL i Python definiowane przez użytkownika (UDF) w Unity Catalog oraz wsadowe funkcje Python definiowane przez użytkownika (UDF) w Unity Catalog.

Nieskalarne funkcje zdefiniowane przez użytkownika

Nieskalowane funkcje zdefiniowane przez użytkownika działają na całych zestawach danych/kolumnach z elastycznymi współczynnikami wejściowymi/wyjściowymi (1:N lub wiele:wiele).

Funkcje zdefiniowane przez użytkownika (UDF) Panda w sesji mogą przyjmować następujące postacie:

  • Seria do serii
  • Iterator serii na iterator serii
  • Iterator wielu serii do iteratora serii
  • Szereg do skalara

Poniżej przedstawiono przykład funkcji UDF serii do serii pandas.

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()

Zobacz funkcje zdefiniowane przez użytkownika pandas.

UDAF

Funkcje UDAFs działają na wielu wierszach i zwracają jeden zagregowany wynik. UDAFs są ograniczone tylko do zakresu sesji.

Poniższy przykład UDAF agreguje wyniki według długości nazwy.

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    |
+-------------+-------------+

Zobacz funkcje zdefiniowane przez użytkownika biblioteki pandas dla funkcji agregujących zdefiniowanych przez użytkownika (UDAFs) Python i Scala.

Funkcje zdefiniowane przez użytkownika

Funkcja UDTF przyjmuje co najmniej jeden argument wejściowy i zwraca wiele wierszy (i ewentualnie wiele kolumn) dla każdego wiersza wejściowego. Mogą być zarządzane przez Unity Catalog lub obejmować zakres sesji.

Następujący program UDTF tworzy tabelę przy użyciu stałej listy dwóch argumentów liczb całkowitych:

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    |
+-----+------+

Aby zaimplementować tę funkcję w notesie Azure Databricks przy użyciu narzędzia PySpark:

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()

Zobacz UDTF-y katalogu Unity i UDTF-y w zakresie sesji.

Zarządzane funkcje UDF w Unity Catalog a funkcje UDF o zakresie sesji

Unity Catalog przechowuje funkcje UDF zarządzane przez Unity Catalog, aby zapewnić lepszy nadzór, ponowne wykorzystanie i łatwiejsze odnajdywanie. Definiujesz funkcje UDF o zakresie sesji w notatniku lub zadaniu, ograniczone do bieżącej sesji SparkSession. Funkcje UDF o zakresie sesji można definiować i uzyskiwać do nich dostęp przy użyciu języka SQL, Python lub Scala.

Skorzystaj z poniższej tabeli, aby wybrać jedną z dwóch kategorii, a następnie zapoznaj się ze skrótowymi zestawieniami dotyczącymi poszczególnych typów UDF w każdej z nich.

Rozważenie Funkcje zdefiniowane przez użytkownika zarządzane przez Unity Catalog Funkcje zdefiniowane przez użytkownika w zakresie sesji
Najlepsze dla Bezpieczne udostępnianie funkcji między zespołami, notesami, zadaniami i magazynami SQL Warehouse. Szybkie, iteracyjne programowanie w ramach jednego notesu lub zadania.
Languages SQL, Python, Scala i Java. SQL, Python i Scala.
Nadzór i udostępnianie Objęte uprawnieniami Unity Catalog i widoczne w Catalog Explorer. Ograniczone do bieżącej sesji SparkSession. Niezarządzane ani niewspółdzielone.
Wytrwałość Trwale zapisane w Unity Catalog i możliwe do ponownego użycia w różnych sesjach. Istnieją tylko w bieżącej sesji.

Ściągawka dotycząca zarządzanych przez katalog Unity funkcji UDF

UDF-y zarządzane przez Unity Catalog pozwalają na definiowanie, używanie, bezpieczne udostępnianie oraz zarządzanie funkcjami niestandardowymi w środowiskach obliczeniowych. Zobacz funkcje definiowane przez użytkownika (UDF) w językach SQL i Python w Unity Catalog.

Typ funkcji UDF Obsługiwane zasoby obliczeniowe Opis
Funkcja UDF w języku Python w Unity Catalog
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 13.3 LTS i nowsze)
  • SQL Warehouse (bezserwerowe i pro)
  • Potoki Lakeflow (klasyczne i bezserwerowe)
Zdefiniuj funkcję użytkownika w języku Python i zarejestruj ją w katalogu Unity w celu zarządzania.
Skalarne funkcje zdefiniowane przez użytkownika działają dla pojedynczego wiersza i zwracają pojedynczą wartość dla każdego wiersza.
Batch Unity Catalog funkcja UDF w Pythonie
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 16.3 lub nowszym)
  • SQL Warehouse (bezserwerowe i pro)
Zdefiniuj funkcję użytkownika w języku Python i zarejestruj ją w katalogu Unity w celu zarządzania.
Operacje wsadowe przetwarzają wiele wartości i zwracają je. Zmniejsza obciążenie operacji wiersz po wierszu na potrzeby przetwarzania danych na dużą skalę.
Unity Catalog — funkcja UDTF języka Python
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 17.1 lub nowszy)
  • SQL Warehouse (bezserwerowe i pro)
Zdefiniuj funkcję UDTF w języku Python i zarejestruj ją w Unity Catalog dla celów zarządzania.
Funkcja UDTF przyjmuje co najmniej jeden argument wejściowy i zwraca wiele wierszy (i ewentualnie wiele kolumn) dla każdego wiersza wejściowego.
Funkcja zdefiniowana przez użytkownika Scala lub Java w Unity Catalog
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia (standardowy dostęp i tryb dedykowanego dostępu)
  • Magazyny SQL (bezserwerowe, pro i klasyczne)
  • Deklaratywne potoki Spark w usłudze Lakeflow (klasyczne i bezserwerowe)
Zdefiniuj funkcję UDF w języku Scala lub Java i zarejestruj ją w Unity Catalog na potrzeby nadzoru.
Skalarne funkcje zdefiniowane przez użytkownika działają dla pojedynczego wiersza i zwracają pojedynczą wartość dla każdego wiersza. Wymaga języka Scala 2.13.16, JDK 17 i środowiska w wersji 4.

Ściągawka funkcji zdefiniowanych przez użytkownika w zakresie sesji dla obliczeń izolowanych przez użytkownika

Funkcje UDF o zakresie sesji definiuje się w notatniku lub zadaniu i są one ograniczone do bieżącej sesji SparkSession. Możesz definiować UDF-y o zakresie sesji i uzyskiwać do nich dostęp za pomocą języków SQL, Python lub Scala.

Typ funkcji UDF Obsługiwane zasoby obliczeniowe Opis
Skalarny języka Python
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 13.3 LTS i nowsze)
  • Potoki Lakeflow (klasyczne i bezserwerowe)
Skalarne funkcje zdefiniowane przez użytkownika działają dla pojedynczego wiersza i zwracają pojedynczą wartość dla każdego wiersza.
Python nieskalarny
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 14.3 LTS i nowsze)
  • Potoki Lakeflow (klasyczne i bezserwerowe)
Nieskalowane funkcje zdefiniowane przez użytkownika obejmują pandas_udf, mapInPandas, mapInArrow, applyInPandas. Funkcje zdefiniowane przez użytkownika biblioteki Pandas używają narzędzia Apache Arrow do przesyłania danych i biblioteki pandas do pracy z danymi. Funkcje użytkownika Pandas obsługują operacje wektoryzowane, które mogą znacznie zwiększyć wydajność w porównaniu do skalarnych funkcji wykonywanych wiersz po wierszu.
Funkcje zdefiniowane przez użytkownika języka Python
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 14.3 LTS i nowsze)
  • Potoki Lakeflow (klasyczne i bezserwerowe)
Funkcja UDTF przyjmuje co najmniej jeden argument wejściowy i zwraca wiele wierszy (i ewentualnie wiele kolumn) dla każdego wiersza wejściowego.
Scala skalarne funkcje zdefiniowane przez użytkownika
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 13.3 LTS i nowsze)
Skalarne funkcje zdefiniowane przez użytkownika działają dla pojedynczego wiersza i zwracają pojedynczą wartość dla każdego wiersza.
Funkcja UDF Scala lub Java z pliku JAR
  • Notesy i zadania bezserwerowe
  • Klasyczne obliczenia ze standardowym trybem dostępu (Databricks Runtime 18.3 lub nowszym)
Zarejestruj wstępnie skompilowaną klasę UDF z pliku JAR przy użyciu polecenia spark.udf.registerJavaFunction. Zobacz Jak zarejestrować funkcję UDF języka Java z pliku JAR.
Funkcje agregujące definiowane przez użytkownika (UDAFs) w języku Scala
  • Klasyczne obliczenia z dedykowanym trybem dostępu (Środowisko Databricks Runtime 14.2 LTS i nowsze)
Funkcje UDAFs działają na wielu wierszach i zwracają jeden zagregowany wynik.

Zagadnienia dotyczące wydajności

  • Wbudowane funkcje i funkcje zdefiniowane przez użytkownika SQL to najbardziej wydajne opcje.

  • Funkcje zdefiniowane przez użytkownika języka Scala są zwykle szybsze niż funkcje zdefiniowane przez użytkownika języka Python.

    • Unisolated Scala UDF działają na maszynie wirtualnej Java (JVM), dzięki czemu unikają obciążenia związanego z przenoszeniem danych do i z JVM.
    • Izolowane funkcje UDF Scala muszą przesyłać dane do i z JVM, ale nadal mogą być szybsze niż funkcje UDF języka Python, ponieważ wydajniej zarządzają pamięcią.
  • Funkcje UDF Pythona i funkcje UDF pandas są zwykle wolniejsze niż funkcje UDF języka Scala, ponieważ wymagają serializacji danych i przesyłania ich poza maszynę JVM do interpretera Pythona.

    • UDF Pandas są do 100 razy szybsze niż UDF Python, ponieważ wykorzystują Apache Arrow do zmniejszenia kosztów serializacji.