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.
Python skaláris UDF-ek lehetővé teszik, hogy egyedi Python logikát futtass SQL lekérdezéseken belül az Azure Databricks-en. Ez az oldal bemutatja, hogyan lehet regisztrálni és megidézni őket, hogyan lehet szolgáltatási hitelesítéseket és titkokat használni, valamint hogyan kezeljük az alexpresszió értékelési sorrendi figyelmeztetéseket a Spark SQL-ben.
Requirements
A Databricks Runtime 12.2 LTS és újabb verziókban a Python UDF-ek és a Pandas UDF-ek nem támogatottak a Standard hozzáférési módot használó Unity Catalog-számításban.
A Skaláris Python UDF-ek és Pandas UDF-ek a Databricks Runtime 13.3 LTS-ben és annál is támogatottak az összes hozzáférési mód esetében.
A Unity katalógusbarát fürtökön futó Python UDF-ek ARM-példánytámogatásához a Databricks Runtime 15.2-s vagy újabb verziója szükséges.
A Databricks Runtime 14.0-s és újabb verzióban a Python UDF-ek és a Pandas UDF-ek nem támogatottak a Standard hozzáférési módot használó Unity Catalog-fürtökön. A Skaláris Python UDF-ek és a Pandas UDF-ek a Databricks Runtime 14.1 és újabb verziók összes hozzáférési módjához támogatottak.
A Databricks Runtime 14.1-ben és újabb verziókban skaláris Python UDF-eket regisztrálhat a Unity Catalogban SQL-szintaxissal. Lásd: SQL és Python felhasználó által definiált függvények (UDF-ek) a Unity Catalogban.
Függvény regisztrálása felhasználó által definiált függvényként (UDF)
def squared(s):
return s * s
spark.udf.register("squaredWithPython", squared)
Az UDF visszatérési típusát igény szerint beállíthatja. Az alapértelmezett visszatérési típus a következő StringType.
from pyspark.sql.types import LongType
def squared_typed(s):
return s * s
spark.udf.register("squaredWithPython", squared_typed, LongType())
Az UDF meghívása a Spark SQL-ben
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, squaredWithPython(id) as id_squared from test
UDF használata DataFrame-ekkel
from pyspark.sql.functions import udf
from pyspark.sql.types import LongType
squared_udf = udf(squared, LongType())
df = spark.table("test")
display(df.select("id", squared_udf("id").alias("id_squared")))
Másik lehetőségként deklarálhatja ugyanezt az UDF-et széljegyzetszintaxissal:
from pyspark.sql.functions import udf
@udf("long")
def squared_udf(s):
return s * s
df = spark.table("test")
display(df.select("id", squared_udf("id").alias("id_squared")))
Változatok UDF-fel
A variáns PySpark típusa VariantType, és az értékek típusa VariantVal. A változatokról további információt a Lekérdezésvariáns adatok című témakörben talál.
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import VariantType, VariantVal
# Return Variant
@udf(returnType = VariantType())
def toVariant(jsonString):
return VariantVal.parseJson(jsonString)
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toVariant(col("json"))).display()
+---------------+
|toVariant(json)|
+---------------+
| {"a":1}|
+---------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import StructField, StructType, VariantType, VariantVal
# Return Struct<Variant>
@udf(returnType = StructType([StructField("v", VariantType(), True)]))
def toStructVariant(jsonString):
return {"v": VariantVal.parseJson(jsonString)}
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toStructVariant(col("json"))).display()
+---------------------+
|toStructVariant(json)|
+---------------------+
| {"v":{"a":1}}|
+---------------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import ArrayType, VariantType, VariantVal
# Return Array<Variant>
@udf(returnType = ArrayType(VariantType()))
def toArrayVariant(jsonString):
return [VariantVal.parseJson(jsonString)]
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toArrayVariant(col("json"))).display()
+--------------------+
|toArrayVariant(json)|
+--------------------+
| [{"a":1}]|
+--------------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import MapType, StringType, VariantType, VariantVal
# Return Map<String, Variant>
@udf(returnType = MapType(StringType(), VariantType(), True))
def toMapVariant(jsonString):
return {"v1": VariantVal.parseJson(jsonString), "v2": VariantVal.parseJson("[" + jsonString + "]")}
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toMapVariant(col("json"))).display()
+-----------------------------+
| toMapVariant(json)|
+-----------------------------+
|{"v2":[{"a":1}],"v1":{"a":1}}|
+-----------------------------+
Fájlok UDF-hez
Important
Ez a funkció bétaverzióban érhető el. A munkaterület rendszergazdái az Előnézetek lapon szabályozhatják a funkcióhoz való hozzáférést. Lásd: Az Azure Databricks előzetes verziójának kezelése.
A fájl PySpark típusa .FileType Használd paraméterként vagy visszaküldési típusként egy UDF-ben, akár felső szintű típusként, akár fészkelveként. A típusról, annak beágyazási szabályairól és az FileRef API-ról lásd: FileType.
Egy fájl tartalmának UDF-ben való olvasásához hívj file.as_local_file() , hogy helyi útvonalat nyitj meg, vagy file.open() hogy a bájtjait streamként olvasd fel. A Python, Scala és SQL példáihoz, beleértve a képfeldolgozást, fájltípus-felismerést és videóképkocka kibontását, lásd: Processzfájlok UDF-ekkel. A típusreferenciáért lásd: FILE típus.
A Unity Catalog-ban (CREATE FUNCTION) regisztrált UDF képes olvasni a ' FILEs metaadatait, de nem a tartalmát, és nem tud fájlokat létrehozni. Használj egy session-scope-s UDF-et a fájl tartalmának olvasásához (open, as_local_file) vagy létrehozz egyet (from_bytes, from_local_file).
Az UDF testben minden FILE érték egy FileRef objektum:
from pyspark.sql.functions import col, udf
from pyspark.sql.types import FileRef, StringType
@udf(returnType=StringType())
def file_content_type(file: FileRef) -> str:
return file.content_type
df = spark.table("documents")
display(df.select(col("file").uri, file_content_type(col("file"))))
Kiértékelési sorrend és nullérték-ellenőrzés
A Spark SQL (beleértve az SQL-t és a DataFrame- és Dataset API-t) nem garantálja az alexpressziók kiértékelésének sorrendjét. Az operátorok vagy függvények bemeneteit nem feltétlenül értékelik ki balról jobbra vagy más rögzített sorrendben. A logikai AND és OR kifejezési kifejezések például nem rendelkeznek balról jobbra "rövidzárolás" szemantikával.
Ezért veszélyes a logikai kifejezések mellékhatásaira vagy kiértékelési sorrendjére, valamint a záradékok sorrendjére WHEREHAVING támaszkodni, mivel az ilyen kifejezések és záradékok átrendezhetők a lekérdezésoptimalizálás és -tervezés során. Pontosabban, ha egy UDF az SQL rövidzárlatú szemantikára támaszkodik a nullellenőrzéshez, nincs garancia arra, hogy a null ellenőrzés az UDF meghívása előtt fog történni. Például,
spark.udf.register("strlen", lambda s: len(s), "int")
spark.sql("select s from test1 where s is not null and strlen(s) > 1") # no guarantee
Ez a WHERE záradék nem garantálja az UDF meghívását a strlen null értékek szűrése után.
A megfelelő nullellenőrzés végrehajtásához az alábbiak valamelyikét javasoljuk:
- Állítsa az UDF-et nullérzékenysé, és végezze el a null ellenőrzést magában az UDF-ben
- Használja a
IFvagy aCASE WHENkifejezéseket a null-ellenőrzéshez és az UDF feltételes ágban való meghívásához.
spark.udf.register("strlen_nullsafe", lambda s: len(s) if not s is None else -1, "int")
spark.sql("select s from test1 where s is not null and strlen_nullsafe(s) > 1") # ok
spark.sql("select s from test1 where if(s is not null, strlen(s), null) > 1") # ok
Hozzáférés Unity katalógus titkaihoz
Egy Unity Catalog titokhoz egy session-scoped Python UDF-ből való hozzáféréshez lásd: Titkot használj egy session-scoped Python UDF-ben. A deklarált titkokhoz való hozzáféréshez egy skaláris vagy Batch Unity Catalog Python UDF-ből lásd: Titkok használata a Python UDF-ben.
Szolgáltatási hitelesítések Python UDF-ekben
A session-scope-os skaláris Python UDF-ek és skaláris Unity Catalog Python UDF-ek Unity Catalog szolgáltatási hitelesítő adatokat használhatnak a külső felhőszolgáltatások biztonságos eléréséhez. Ez olyan műveletek integrálásához hasznos, mint a felhőalapú tokenizálás, titkosítás vagy titkos kulcskezelés közvetlenül az adatátalakításokba.
A követelmények UDF típusonként és számítási rendszertől függően változnak. Lásd: Használj szolgáltatási hitelesítést egy Python UDF-ben.
Szolgáltatás hitelesítő adatainak létrehozásáról a Szolgáltatás hitelesítő adatainak létrehozása című témakörben olvashat.
Használj szolgáltatási hitelminősítést egy session-scope-jú skaláris Python UDF-ben
A szolgáltatás hitelesítő adatainak eléréséhez használja az databricks.service_credentials.getServiceCredentialsProvider() UDF-logika segédprogramját a felhőbeli SDK-k inicializálásához a megfelelő hitelesítő adatokkal. Az összes kódot bele kell ágyazni az UDF törzsébe.
@udf
def use_service_credential():
from azure.mgmt.web import WebSiteManagementClient
# Assuming there is a service credential named 'testcred' set up in Unity Catalog
web_client = WebSiteManagementClient(subscription_id, credential = getServiceCredentialsProvider('testcred'))
# Use web_client to perform operations
Szolgáltatási hitelesítő adatok jogosultságai
A munkamenet-hatókörű UDF-ek a hívó jogosultságaival működnek. A szükséges jogosultságokért lásd: Használj szolgáltatási hitelesítő adatokat egy Python UDF-ben.
Számítási szintű alapértelmezett hitelesítő adatok munkamenet-hatókörű UDF-ekhez
Skaláris Python UDF-ekben a Databricks automatikusan használja az alapértelmezett szolgáltatási jogosultságot a számítási környezeti változóból. Ez a viselkedés lehetővé teszi, hogy biztonságosan hivatkozzon külső szolgáltatásokra anélkül, hogy explicit módon kezelne hitelesítő aliasokat az UDF-kódban. Lásd: Alapértelmezett szolgáltatás hitelesítő adatainak megadása számítási erőforráshoz
Az alapértelmezett hitelesítő adatok támogatása csak standard és dedikált hozzáférési módú fürtökben érhető el. Ez nem elérhető SQL raktárakban.
Ahhoz, hogy használni tudd a azure-identity szolgáltatót, telepítened kell a DefaultAzureCredential csomagot. A csomag telepítéséhez tekintse meg a jegyzetfüzet-hatókörű Python-kódtárakat vagy a számítási hatókörű kódtárakat.
@udf
def use_service_credential():
from azure.identity import DefaultAzureCredential
from azure.mgmt.web import WebSiteManagementClient
# DefaultAzureCredential is automatically using the default service credential for the compute
web_client_default = WebSiteManagementClient(DefaultAzureCredential(), subscription_id)
# Use web_client to perform operations
Használj szolgáltatási hitelesítő jogot egy skaláris Unity Catalog Python UDF-ben
A szolgáltatási jogosultságot CREDENTIALS az UDF definíciójának záradékában határozzuk meg. Jelölheted meg egy hitelesítést úgy, DEFAULT hogy a patchelt felhő SDK-k automatikusan használják. Klasszikus számítástechnikán ehhez a funkcióhoz Databricks Runtime 18.1 vagy annál magasabb verzió szükséges. A szerver nélküli számítási környezetben, valamint a pro és szerver nélküli SQL-adattárházakban állítsa be az UDF environment_version elemét kifejezetten 6 vagy magasabb értékre. A teljes számítási, hálózati és jogosultsági követelményekért lásd: Használj szolgáltatási hitelesítő jogot a Python UDF-ben.
Szerezze meg a feladat végrehajtási kontextust
Használja a TaskContext PySpark API-t a kontextus információk lekéréséhez, például a felhasználóazonosság, a klasztercímkék, a Spark feladat azonosító és további adatok megszerzéséhez. Lásd Feladat kontextusának megszerzése UDF-ben.
Korlátozások
A PySpark UDF-ekre a következő korlátozások vonatkoznak:
Fájlhozzáférés korlátozásai: A Databricks Runtime 14.2-ben és az alatta lévő PySpark UDF-ek megosztott fürtökön nem férnek hozzá a Git-mappákhoz, munkaterületfájlokhoz vagy Unity-katalóguskötetekhez.
Szórási változók: A PySpark UDF-jei a standard hozzáférési módú fürtökön és a kiszolgáló nélküli számításban nem támogatják a szórási változókat.
- Memóriakorlát kiszolgáló nélküli rendszeren: A kiszolgáló nélküli számítási pyspark-UDF-k memóriakorlátja PySpark UDF-enként 1 GB. A korlát túllépése UDF_PYSPARK_USER_CODE_ERROR típusú hibát eredményez . MEMORY_LIMIT_SERVERLESS.
- Memóriakorlát standard hozzáférési mód esetén: A PySpark UDF-jei normál hozzáférési módban a kiválasztott példánytípus rendelkezésre álló memóriája alapján memóriakorlátot biztosítanak. A rendelkezésre álló memória túllépése UDF_PYSPARK_USER_CODE_ERROR típusú hibát eredményez . MEMORY_LIMIT.
- Hálózati hozzáférés kiszolgáló nélküli SQL-raktárakban: Alapértelmezés szerint Python kiszolgáló nélküli SQL-raktárakban lévő UDF-ek nem tudnak kimenő hálózati kéréseket küldeni, és a hálózati hívásokat megkísérlő lekérdezések határozatlan ideig lefagynak. A kimenő hálózati hozzáférés engedélyezéséhez engedélyezze a Hálózatkezelés engedélyezése elkülönített számítási feladatokhoz kiszolgáló nélküli SQL-adattárházakban nyilvános előzetes verziójú funkciót a munkaterület Előzetes verziók lapján. Ellenkező esetben használjon kiszolgáló nélküli számítást vagy klasszikus számítást hálózati hozzáférést igénylő UDF-ekhez.