Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
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 |
|
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 |
|
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 |
|
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 |
|
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 |
|
A skaláris UDF-ek egyetlen sorban működnek, és minden sorhoz egyetlen eredményértéket ad vissza. |
| Python nem skaláris |
|
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 |
|
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 |
|
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 |
|
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 |
|
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.