Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Gli UDF scalari Python permettono di eseguire logica Python personalizzata all'interno di query SQL su Azure Databricks. Questa pagina mostra come registrarli e invocarli, utilizzare credenziali di servizio e segreti, e gestire le avvertenze sull'ordine di valutazione delle sottoespressioni in Spark SQL.
Requisiti
In Databricks Runtime 12.2 LTS e versioni precedenti, le funzioni definite dall'utente Python e Pandas non sono supportate nelle risorse di calcolo di Unity Catalog che usano la modalità di accesso standard.
Le funzioni utente scalari Python e le funzioni utente Pandas sono supportate in Databricks Runtime 13.3 LTS e versioni successive per tutte le modalità di accesso.
Il supporto delle istanze ARM per le funzioni definite dall'utente Python nei cluster abilitati per Unity Catalog richiede Databricks Runtime 15.2 o versione successiva.
In Databricks Runtime 14.0 e versioni precedenti, le UDF (funzioni definite dall'utente) Python e le UDF Pandas non sono supportate nei cluster Unity Catalog che usano la modalità di accesso standard. Le funzioni scalari definite dall'utente in Python e le funzioni definite dall'utente in Pandas sono supportate per tutte le modalità di accesso in Databricks Runtime 14.1 e versioni successive.
È possibile registrare UDF Python scalari in Unity Catalog usando la sintassi SQL in Databricks Runtime 14.1 e versioni successive. Vedi le funzioni definite dall'utente (UDF) in SQL e Python in Unity Catalog.
Registrare una funzione come UDF
def squared(s):
return s * s
spark.udf.register("squaredWithPython", squared)
È possibile facoltativamente impostare il tipo di ritorno della funzione definita dall'utente. Il tipo restituito predefinito è StringType.
from pyspark.sql.types import LongType
def squared_typed(s):
return s * s
spark.udf.register("squaredWithPython", squared_typed, LongType())
Richiama l’UDF in Spark SQL
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, squaredWithPython(id) as id_squared from test
Usare le UDF con i DataFrame
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")))
In alternativa, è possibile dichiarare la stessa UDF usando la sintassi di annotazione.
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")))
Varianti con UDF
Il tipo PySpark per variant è VariantType e i valori sono di tipo VariantVal. Per informazioni sulle varianti, vedere Eseguire query dei dati delle varianti.
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}}|
+-----------------------------+
File con UDF
Importante
Questa funzionalità è in versione beta.
Il tipo PySpark per un file è FileType. Usarlo come parametro o come tipo restituito in una UDF, sia come tipo di primo livello sia come tipo annidato. Per il tipo, le sue regole di annidamento e l'API FileRef , vedi FileType.
Per leggere il contenuto di un file in un UDF, chiama file.as_local_file() per ottenere un percorso locale che puoi aprire, oppure file.open() per leggere i suoi byte come un stream. Per esempi in Python, Scala e SQL, inclusi l'elaborazione delle immagini, il rilevamento dei tipi di file e l'estrazione dei fotogrammi video, vedi File di processo con UDF. Per il riferimento tipografico, vedi FILE tipo.
Un UDF registrato nel Catalogo Unity (CREATE FUNCTION) può leggere i metadati di un FILE', ma non i suoi contenuti, e non può creare file. Usa un UDF con ambito di sessione per leggere il contenuto di un file (open, as_local_file) o crearne uno (from_bytes, from_local_file).
Nel corpo UDF, ogni FILE valore è un FileRef oggetto:
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"))))
Ordine di valutazione e controllo del valore nullo
Spark SQL (incluso SQL e l'API DataFrame e Dataset) non garantisce l'ordine di valutazione delle sottoespressioni. In particolare, gli input di un operatore o di una funzione non vengono necessariamente valutati da sinistra a destra o in qualsiasi altro ordine fisso. Ad esempio, le espressioni logiche AND e OR non hanno la semantica di “cortocircuito” da sinistra a destra.
Pertanto, è pericoloso basarsi sugli effetti collaterali o sull'ordine di valutazione delle espressioni booleane e sull'ordine delle clausole WHERE e HAVING, poiché tali espressioni e clausole possono essere riordinate durante l'ottimizzazione e la pianificazione delle query. In particolare, se una UDF si basa sulla semantica di corto circuito in SQL per il controllo dei valori null, non c'è garanzia che il controllo venga eseguito prima di richiamare la UDF. Ad esempio,
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
Questa clausola WHERE non garantisce che la UDF strlen venga invocata dopo il filtraggio dei valori nulli.
Per eseguire un controllo null appropriato, è consigliabile eseguire una delle operazioni seguenti:
- Rendere la funzione definita dall'utente null-aware ed eseguire il controllo dei valori null all'interno di questa stessa funzione.
- Utilizzare le espressioni
IFoCASE WHENper effettuare il controllo su valori nulli e invocare una funzione definita dall'utente all'interno di un ramo condizionale.
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
Accedi ai segreti di Unity Catalog
Per accedere a un segreto del Catalogo Unity da un UDF Python con ambito di sessione, vedi Usa un segreto in un UDF Python con ambito di sessione. Per accedere ai secreti dichiarati da una UDF Python scalare o batch di Unity Catalog, consulta Usare i segreti in una UDF Python.
Credenziali di servizio nei UDF Python
Le UDF Python scalari con ambito di sessione e le UDF Python scalari di Unity Catalog possono utilizzare le credenziali di servizio di Unity Catalog per accedere in modo sicuro ai servizi cloud esterni. Ciò è utile per l'integrazione di operazioni quali la tokenizzazione basata sul cloud, la crittografia o la gestione dei segreti direttamente nelle trasformazioni dei dati.
I requisiti variano in base al tipo di UDF e al calcolo. Vedi Usa una credenziale di servizio in un UDF Python.
Per creare credenziali del servizio, vedere Creare credenziali del servizio.
Usa una credenziale di servizio in un UDF scalar Python con ambito di sessione
Per accedere alle credenziali del servizio, utilizzare l'utilità nella logica UDF databricks.service_credentials.getServiceCredentialsProvider() per inizializzare gli SDK cloud con le credenziali appropriate. Tutto il codice deve essere incapsulato nel corpo della funzione definita dall'utente.
@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
Autorizzazioni per le credenziali del servizio
Gli UDF con ambito di sessione utilizzano i permessi del chiamante. Vedi Usa una credenziale di servizio in un UDF Python per i privilegi richiesti.
Credenziali predefinite del livello di elaborazione per UDF con ambito limitato alla sessione
Quando viene utilizzato negli UDF scalari Python, Databricks utilizza automaticamente la credenziale di servizio predefinita della variabile dell'ambiente di calcolo. Questo comportamento consente di fare riferimento in modo sicuro ai servizi esterni senza gestire in modo esplicito gli alias delle credenziali nel codice UDF. Vedere Specificare una credenziale del servizio predefinita per una risorsa di calcolo
Il supporto predefinito delle credenziali è disponibile solo nei cluster in modalità di accesso standard e dedicato. Non è disponibile nei magazzini SQL.
È necessario installare il azure-identity pacchetto per usare il DefaultAzureCredential provider. Per installare il pacchetto, vedere Librerie Python con ambito notebook o librerie con ambito calcolo.
@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
Usa una credenziale di servizio in un UDF Python scalare di Unity Catalog
Specifica la credenziale del servizio nella CREDENTIALS clausola della definizione UDF. Puoi segnare una credenziale come DEFAULT così che gli SDK cloud patchati la utilizzino automaticamente. Nel calcolo classico, questa caratteristica richiede Databricks Runtime 18.1 o superiore. Nel calcolo serverless e nei warehouse SQL Pro e Serverless, imposta esplicitamente environment_version dell'UDF su 6 o versione successiva. Per requisiti completi di calcolo, rete e permessi, vedi Usa una credenziale di servizio in un UDF Python.
Ottenere il contesto di esecuzione delle attività
Usare l'API PySpark TaskContext per ottenere informazioni di contesto, ad esempio l'identità dell'utente, i tag del cluster, l'ID processo Spark e altro ancora. Vedi Ottieni il contesto del compito in una UDF.
Limiti
Le seguenti limitazioni si applicano alle UDF di PySpark.
Restrizioni di accesso ai file: In Databricks Runtime 14.2 e versioni precedenti le funzioni definite dall'utente PySpark nei cluster condivisi non possono accedere a cartelle Git, file dell'area di lavoro o volumi del catalogo Unity.
Variabili di trasmissione: le UDF PySpark nei cluster in modalità di accesso standard e nel calcolo serverless non supportano le variabili di trasmissione.
- Limite di memoria per serverless: le PySpark UDFs nel calcolo serverless hanno un limite di memoria di 1 GB per ciascuna PySpark UDF. Il superamento di questo limite genera un errore di tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT_SERVERLESS.
- Limite di memoria in modalità di accesso standard: le UDF di PySpark in modalità di accesso standard hanno un limite di memoria basato sulla memoria disponibile del tipo di istanza scelto. Il superamento della memoria disponibile genera un errore di tipo UDF_PYSPARK_USER_CODE_ERROR. MEMORY_LIMIT.
- Accesso alla rete nei data warehouse SQL serverless: per impostazione predefinita, le UDF Python nei data warehouse SQL serverless non possono effettuare richieste di rete in uscita e le query che tentano di effettuare chiamate di rete rimangono bloccate indefinitamente. Per abilitare l'accesso alla rete in uscita, abilitare la funzionalità Anteprima pubblica Abilitare la rete per i carichi di lavoro isolati in SQL Warehouse serverless nella pagina Anteprime dell'area di lavoro. In caso contrario, utilizzare l'elaborazione serverless o l'elaborazione classica per le UDF che richiedono l'accesso alla rete.