Mik azok a felhasználó által definiált függvények (UDF-ek)?

A felhasználó által definiált függvények (UDF-ek) lehetővé teszik a beépített képességeket Azure Databricks kiterjesztő kód újrafelhasználását és megosztását. Az UDF-ek használatával konkrét feladatokat hajthat végre, például összetett számításokat, átalakításokat vagy egyéni adatmanipulációkat.

Mikor érdemes UDF és Apache Spark függvényt használni?

A beépített Apache Spark-függvényekkel nehezen kifejezhető logikai UDF-eket használjon. A beépített Apache Spark-függvények elosztott feldolgozásra vannak optimalizálva, és nagyobb léptékben jobb teljesítményt nyújtanak. További információ: Functions.

A Databricks alkalmi lekérdezésekhez, manuális adattisztításhoz, feltáró adatelemzéshez és kis- és közepes méretű adathalmazokon végzett műveletekhez javasolja az UDF-eket. Az UDF-ek gyakori használati esetei közé tartozik az adattitkosítás, a visszafejtés, a kivonatolás, a JSON-elemzés és az ellenőrzés.

Apache Spark-metódusok használata nagy adathalmazokon és rendszeresen vagy folyamatosan futó számítási feladatokhoz, beleértve az ETL-feladatokat és a streamelési műveleteket.

UDF-típusok ismertetése

Válasszon ki egy UDF-típust az alábbi lapok közül, és további információért tekintse meg a leírást, a példát és a hivatkozást.

Skaláris UDF

A skaláris UDF-ek egyetlen sorban működnek, és minden sorhoz egyetlen eredményértéket ad vissza. Ezek lehetnek a Unity Catalog által szabályozott vagy munkamenet-hatókörűek.

Az alábbi példa egy skaláris UDF használatával számítja ki az egyes nevek hosszát egy name oszlopban, és hozzáadja az értéket egy új oszlophoz 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      |
+-------+-------+-------------+

Ha ezt egy Azure Databricks-jegyzetfüzetben szeretné implementálni a PySpark használatával:

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)

Lásd: SQL és Python felhasználó által definiált függvények (UDF-ek) a Unity Catalogban, és Python skaláris felhasználó által definiált függvényeket (UDF-eket).

Batch Skaláris UDF-ek

Adatok feldolgozása kötegekben az 1:1 bemeneti/kimeneti sor paritás fenntartása mellett. Ez csökkenti a sorról sorra végzett műveletek többletterhelését a nagy léptékű adatfeldolgozáshoz. A Batch UDF-ek a kötegek között is fenntartják az állapotot a hatékonyabb futtatás, az erőforrások újrafelhasználása és az adattömbök környezetét igénylő összetett számítások kezelése érdekében.

Ezek lehetnek a Unity Catalog által szabályozott vagy munkamenet-hatókörűek.

A következő Batch Unity Catalog Python UDF kiszámítja a BMI-t a sorok kötegeinek feldolgozása során:

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

Lásd: SQL és Python felhasználó által definiált függvények (UDF-ek) a Unity Katalógusban és a Batch Python felhasználó által definiált függvények (UDF-ek) a Unity Katalógusban.

Nem skaláris UDF-ek

A nem skaláris UDF-ek teljes adatkészleteken/oszlopokon működnek rugalmas bemeneti/kimeneti arányokkal (1:N vagy több:több).

A munkamenet-hatókörű batch pandas UDF-ek a következő típusúak lehetnek:

  • Sorozatról sorozatra
  • Sorozat iterátora a sorozat iterátorához
  • Több sorozat iterátora a sorozat iterátorához
  • Sorozatból skalárisra

Az alábbiakban egy példa látható a Series to Series pandas UDF-ra.

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

Lásd a pandas felhasználó által definiált függvényeit.

UDAF

Az UDAF-k több sorban működnek, és egyetlen összesített eredményt ad vissza. Az UDAF-ek csak munkamenet-hatókörűek.

Az alábbi UDAF-példa névhossz szerint összesíti a pontszámokat.

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

Lásd a Pythonhoz készült, felhasználó által definiált pandas-függvényeket és a Scala felhasználó által definiált aggregátumfüggvényeit (UDAF-eket).

UDTF-ek

Az UDTF egy vagy több bemeneti argumentumot vesz fel, és több sort (és esetleg több oszlopot) ad vissza minden egyes bemeneti sorhoz. Ezek lehetnek a Unity Catalog által szabályozott vagy munkamenet-hatókörűek.

A következő UDTF két egész argumentum rögzített listájával hoz létre egy táblát:

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

Ha ezt egy Azure Databricks-jegyzetfüzetben szeretné implementálni a PySpark használatával:

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

Lásd a Unity Catalog UDTF-eit és a munkamenet-hatókörű UDTF-eket.

A Unity Katalógus által szabályozott és munkamenet-hatókörű UDF-ek

A Unity Catalog megőrzi a Unity Katalógus által szabályozott UDF-eket a jobb szabályozás, az újrafelhasználás és a felderíthetőség érdekében. A munkamenet-hatókörű UDF-eket egy jegyzetfüzetben vagy feladatban határozhatja meg, amely az aktuális SparkSession-hatókörre terjed ki. A munkamenet-hatókörű UDF-eket SQL, Python vagy Scala használatával definiálhatja és érheti el.

Az alábbi táblázat segítségével döntsön a két kategória között, majd tekintse meg az ezt követő összefoglaló lapokat az azokban szereplő konkrét UDF-típusokhoz.

Consideration Unity Catalog által szabályozott UDF-ek Munkamenet-hatókörű UDF-ek
Legjobb a számára A megosztás biztonságosan működik a csapatok, jegyzetfüzetek, feladatok és SQL-raktárak között. Gyors, iteratív fejlesztés egyetlen jegyzetfüzetben vagy feladatban.
Languages SQL, Python, Scala és Java. SQL, Python és Scala.
Irányítás és megosztás A Unity Catalog engedélyeinek hatálya alá tartozó és a Catalog Explorerben felderíthető. Az aktuális SparkSession hatóköre. Nincs szabályozva vagy megosztva.
Kitartás A Unity Katalógusban megmarad, és a munkamenetek között újra felhasználható. Csak az aktuális munkamenet során léteznek.

Unity Catalog által szabályozott UDF-ek gyorstalpalója

A Unity Catalog által szabályozott UDF-ek lehetővé teszik az egyéni függvények meghatározását, használatát, biztonságos megosztását és szabályozását a számítási környezetekben. Lásd: SQL és Python felhasználó által definiált függvények (UDF-ek) a Unity Catalogban.

UDF-típus Támogatott számítás Leírás
Unity Catalog Python UDF
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 13.3 LTS és újabb)
  • SQL Warehouse (kiszolgáló nélküli és pro)
  • Lakeflow-folyamatok (klasszikus és kiszolgáló nélküli)
Definiáljon egy UDF-t a Pythonban, és regisztrálja a Unity Katalógusban a szabályozáshoz.
A skaláris UDF-ek egyetlen sorban működnek, és minden sorhoz egyetlen eredményértéket ad vissza.
Batch Unity Katalógus Python UDF
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 16.3 vagy újabb)
  • SQL Warehouse (kiszolgáló nélküli és pro)
Definiáljon egy UDF-t a Pythonban, és regisztrálja a Unity Katalógusban a szabályozáshoz.
Batch-műveletek több értéken, és több értéket ad vissza. Csökkenti a sorról sorra végzett műveletek többletterhelését a nagy léptékű adatfeldolgozáshoz.
Unity Catalog Python UDTF
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 17.1 vagy újabb)
  • SQL Warehouse (kiszolgáló nélküli és pro)
Definiáljon egy UDTF-t a Pythonban, és regisztrálja a Unity Katalógusban a szabályozáshoz.
Az UDTF egy vagy több bemeneti argumentumot vesz fel, és több sort (és esetleg több oszlopot) ad vissza minden egyes bemeneti sorhoz.
Unity Catalog Scala vagy Java UDF
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás (standard hozzáférés és dedikált hozzáférési mód)
  • SQL-raktárak (kiszolgáló nélküli, pro és klasszikus)
  • Spark Deklaratív folyamatok a Lakeflow-on (klasszikus és kiszolgáló nélküli)
Definiáljon egy UDF-t a Scalában vagy Java, és regisztrálja azt a Unity Katalógusban irányítás céljából.
A skaláris UDF-ek egyetlen sorban működnek, és minden sorhoz egyetlen eredményértéket ad vissza. A Scala 2.13.16-os, JDK 17-es és 4-es környezeti verziót igényel.

Munkamenet-hatókörű UDF-ek csalólapja a felhasználó által izolált számításhoz

A munkamenet-hatókörű UDF-eket egy jegyzetfüzetben vagy feladatban határozhatja meg, amely az aktuális SparkSession-hatókörre terjed ki. A munkamenet-hatókörű UDF-eket SQL, Python vagy Scala használatával definiálhatja és érheti el.

UDF-típus Támogatott számítás Leírás
Python skalár
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 13.3 LTS és újabb)
  • Lakeflow-folyamatok (klasszikus és kiszolgáló nélküli)
A skaláris UDF-ek egyetlen sorban működnek, és minden sorhoz egyetlen eredményértéket ad vissza.
Python nem skaláris
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 14.3 LTS és újabb)
  • Lakeflow-folyamatok (klasszikus és kiszolgáló nélküli)
A nem skaláris UDF-ek közé tartozik pandas_udfa , mapInPandas, mapInArrow. applyInPandas A Pandas UDF-ek az Apache Arrow használatával továbbítják az adatokat, a pandas pedig az adatokkal dolgozik. A Pandas UDF-ek olyan vektoros műveleteket támogatnak, amelyek jelentősen növelhetik a teljesítményt a sorról sorra történő skaláris UDF-eken.
Python UDTF-ek
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 14.3 LTS és újabb)
  • Lakeflow-folyamatok (klasszikus és kiszolgáló nélküli)
Az UDTF egy vagy több bemeneti argumentumot vesz fel, és több sort (és esetleg több oszlopot) ad vissza minden egyes bemeneti sorhoz.
Scala skaláris UDF-ek
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 13.3 LTS és újabb)
A skaláris UDF-ek egyetlen sorban működnek, és minden sorhoz egyetlen eredményértéket ad vissza.
Scala- vagy Java-UDF JAR-fájlból
  • Kiszolgáló nélküli jegyzetfüzetek és feladatok
  • Klasszikus számítás standard hozzáférési móddal (Databricks Runtime 18 LTS vagy újabb)
Regisztráljon egy előre lefordított UDF-osztályt egy JAR-fájlból a(z) spark.udf.registerJavaFunction használatával. Lásd: Java UDF regisztrálása JAR-ból.
Scala UDAF-ek
  • Klasszikus számítás dedikált hozzáférési móddal (Databricks Runtime 14.2 LTS és újabb)
Az UDAF-k több sorban működnek, és egyetlen összesített eredményt ad vissza.

Teljesítménnyel kapcsolatos szempontok

  • A beépített függvények és az SQL UDF-ek a leghatékonyabb lehetőségek.

  • A Scala UDF-ek általában gyorsabbak, mint a Python UDF-k.

    • Az nem oldott Scala UDF-ek a Java virtuális gépen (JVM) futnak, így elkerülik az adatok JVM-be és kifelé történő áthelyezésének többletterhelését.
    • Az izolált Scala UDF-eknek át kell helyezniük az adatokat a JVM-be és onnan, de még mindig gyorsabbak lehetnek, mint Python UDF-ek, mert hatékonyabban kezelik a memóriát.
  • Python UDF-ek és pandas UDF-ek általában lassabbak, mint a Scala UDF-ek, mert szerializálniuk kell az adatokat, és át kell helyezniük őket a JVM-ből a Python értelmezőbe.

    • A Pandas UDF-ek akár 100-szor gyorsabbak, mint a Python UDF-ek, mivel az Apache Arrow használatával csökkentik a szerializálási költségeket.