Fonctions définies par l'utilisateur (UDF) Python en batch dans Unity Catalog
Les UDF Python Batch Unity Catalog sont généralement disponibles. Elles traitent les données par batch au lieu d'une ligne à la fois.
Exigences
Sur le compute classique, les UDF Python Unity Catalog par batch nécessitent Databricks Runtime 16.3 ou une version ultérieure. Ils sont également pris en charge sur le compute serverless ainsi que sur les warehouses SQL pro et serverless.
Les capacités supplémentaires ont leurs propres exigences en matière de compute et de version. Consultez la page Python UDF feature requirements.
Créer une UDF Python Unity Catalog par batch
La création d'un UDF Python de Unity Catalog par batch est similaire à la création d'un UDF Unity Catalog régulier, avec les ajouts suivants :
PARAMETER STYLE PANDAS: Ceci indique que l'UDF traite les données par batchs à l'aide d'itérateurs pandas.HANDLER 'handler_function': spécifie la fonction de gestionnaire qui traite les batches.
L’exemple suivant crée une UDF Python batch persistante dans Unity Catalog. Remplacez my_catalog et my_schema par votre catalogue et votre schéma :
%sql
CREATE OR REPLACE FUNCTION my_catalog.my_schema.calculate_bmi_pandas(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
DETERMINISTIC
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)
$$;
Après avoir enregistré l'UDF, vous pouvez l'appeler en utilisant SQL ou Python :
SELECT person_id, my_catalog.my_schema.calculate_bmi_pandas(weight_kg, height_m) AS bmi
FROM (
SELECT 1 AS person_id, CAST(70.0 AS DOUBLE) AS weight_kg, CAST(1.75 AS DOUBLE) AS height_m UNION ALL
SELECT 2 AS person_id, CAST(80.0 AS DOUBLE) AS weight_kg, CAST(1.80 AS DOUBLE) AS height_m
);
Fonction de gestionnaire d'UDF batch
Les UDF Python Batch Unity Catalog nécessitent une fonction de gestion qui traite les batchs et renvoie des résultats. Vous devez spécifier le nom de la fonction de gestionnaire lorsque vous créez l’UDF en utilisant la clause HANDLER.
La fonction de gestionnaire effectue les opérations suivantes :
- Accepte un argument d’itérateur qui itère sur un ou plusieurs
pandas.Series. Chaquepandas.Seriescontient les parameter d’entrée de l’UDF. - Itère sur le générateur et traite les données.
- Renvoie un itérateur de générateur.
Les UDF Python Batch Unity Catalog doivent renvoyer le même nombre de lignes que l'entrée. La fonction de gestionnaire garantit cela en produisant un pandas.Series de la même longueur que la série d'entrée pour chaque batch.
Installer les dépendances personnalisées
Vous pouvez étendre les fonctionnalités des UDF Python de Batch Unity Catalog au-delà de l'environnement Databricks Runtime en définissant des dépendances personnalisées pour les bibliothèques externes.
Consultez Étendre les UDF à l'aide de dépendances personnalisées.
Accéder aux secrets Unity Catalog
Les UDF Python Unity Catalog par batch peuvent accéder aux secrets déclarés dans la clause SECRETS. Vous devez explicitement définir environment_version sur 6 ou une version supérieure. L’invocation directe d’une UDF qui utilise cette clause n’est pas prise en charge sur un compute à mode d’accès dédié. Pour en savoir plus sur l’exception de masquage de colonne, consultez la page Utiliser des UDF compatibles avec les secrets dans des masques de colonnes sur un compute dédié.
Les UDF de batch peuvent accepter des paramètres uniques ou multiples
Single parameter : lorsque la fonction de gestion utilise un paramètre d’entrée unique, elle reçoit un itérateur sur un objet pandas.Series pour chaque batch.
%sql
CREATE OR REPLACE TEMPORARY FUNCTION one_parameter_udf(value INT)
RETURNS STRING
LANGUAGE PYTHON
DETERMINISTIC
PARAMETER STYLE PANDAS
HANDLER 'handler_func'
AS $$
import pandas as pd
from typing import Iterator
def handler_func(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
for value_batch in batch_iter:
d = {"min": value_batch.min(), "max": value_batch.max()}
yield pd.Series([str(d)] * len(value_batch))
$$;
SELECT one_parameter_udf(id), count(*) from range(0, 100000, 3, 8) GROUP BY ALL;
Plusieurs paramètres : Pour plusieurs paramètres d'entrée, la fonction de gestionnaire reçoit un itérateur qui itère sur plusieurs pandas.Series. Les valeurs de la série sont dans le même ordre que les paramètres d'entrée.
%sql
CREATE OR REPLACE TEMPORARY FUNCTION two_parameter_udf(p1 INT, p2 INT)
RETURNS INT
LANGUAGE PYTHON
DETERMINISTIC
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 p1, p2 in batch_iter: # same order as arguments above
yield p1 + p2
$$;
SELECT two_parameter_udf(id , id + 1) from range(0, 100000, 3, 8);
Optimisez les performances en séparant les Opérations coûteuses
Vous pouvez optimiser les opérations coûteuses en calcul en séparant ces opérations de la fonction de gestionnaire. Ceci garantit qu'ils sont exécutés une seule fois plutôt qu'à chaque itération sur des batchs de données.
L'exemple suivant montre comment s'assurer qu'un calcul coûteux n'est effectué qu'une seule fois :
%sql
CREATE OR REPLACE TEMPORARY FUNCTION expensive_computation_udf(value INT)
RETURNS INT
LANGUAGE PYTHON
DETERMINISTIC
PARAMETER STYLE PANDAS
HANDLER 'handler_func'
AS $$
def compute_value():
# expensive computation...
return 1
expensive_value = compute_value()
def handler_func(batch_iter):
for batch in batch_iter:
yield batch * expensive_value
$$;
SELECT expensive_computation_udf(id), count(*) from range(0, 100000, 3, 8) GROUP BY ALL
Isolation de l'environnement
Les environnements d'isolation partagés nécessitent Databricks Runtime 17.1 et versions ultérieures. Dans les versions antérieures, toutes les UDF Python de Unity Catalog (batch) s'exécutent en mode d'isolation strict.
Les UDF Python Unity Catalog de batch avec le même propriétaire et la même session peuvent partager un environnement d'isolation par default. Cela peut améliorer les performances et réduire la consommation de mémoire en diminuant le nombre d'environnements séparés à lancer.
Isolement strict
Pour s'assurer qu'une UDF s'exécute toujours dans son propre environnement entièrement isolé, ajoutez la clause caractéristique STRICT ISOLATION.
La plupart des UDF n'ont pas besoin d'une isolation stricte. Les UDF de traitement de données standard bénéficient de l'environnement d'isolation partagé default et s'exécutent plus rapidement avec une consommation de mémoire inférieure.
Ajoutez la clause caractéristique STRICT ISOLATION aux UDF qui :
- Exécuter l'entrée en tant que code à l'aide de
eval(),exec()ou de fonctions similaires - Écrire des fichiers vers le système de fichiers local
- Modifier les variables globales ou l'état du système
- Modifier les variables d'environnement
L'exemple suivant montre une UDF qui exécute l'entrée comme du code et nécessite une isolation stricte :
CREATE OR REPLACE TEMPORARY FUNCTION eval_string(input STRING)
RETURNS STRING
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_func'
STRICT ISOLATION
AS $$
import pandas as pd
from typing import Iterator
def handler_func(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
for code_series in batch_iter:
def eval_func(code):
try:
return str(eval(code))
except Exception as e:
return f"Error: {e}"
yield code_series.apply(eval_func)
$$;
Identifiants de service dans les UDF Python du catalogue Unity pour le traitement par batch
Les UDF Python Unity Catalog par batch peuvent utiliser les credentials de service Unity Catalog pour accéder aux services cloud externes. C’est particulièrement utile pour intégrer des fonctions cloud telles que les tokeniseurs de sécurité dans les workflows de traitement des données.
API spécifique aux UDF pour les identifiants de service :
Dans les UDF, utilisez databricks.service_credentials.getServiceCredentialsProvider() pour accéder aux identifiants de service.
Cela diffère de la fonction dbutils.credentials.getServiceCredentialsProvider() utilisée dans les notebooks, qui n'est pas disponible dans les contextes d'exécution UDF.
Pour créer un identifiant de service, consultez Créer des identifiants de service.
Spécifiez les informations d'identification du service que vous souhaitez utiliser dans la clause CREDENTIALS dans la définition de l'UDF :
CREATE OR REPLACE TEMPORARY FUNCTION example_udf(data STRING)
RETURNS STRING
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
CREDENTIALS (
`credential-name` DEFAULT,
`complicated-credential-name` AS short_name,
`simple-cred`,
cred_no_quotes
)
AS $$
# Python code here
$$;
Autorisations des identifiants de service
Pour connaître les exigences d’autorisation de création et d’appel selon les types de compute, consultez la page Utiliser des informations d’identification de service dans une UDF Python.
Identifiants default et alias
Vous pouvez inclure plusieurs identifiants dans la clause CREDENTIALS, mais un seul peut être marqué comme DEFAULT. Vous pouvez créer un alias pour les identifiants non-default à l'aide du mot-clé AS. Chaque identifiant doit avoir un alias unique.
Les SDK cloud corrigés reprennent automatiquement les informations d'identification default. L'identifiant default a préséance sur tout paramètre default spécifié dans la configuration Spark du calcul et persiste dans la définition UDF de Unity Catalog.
Exemple d'identifiant de service - Google Cloud Storage
L'exemple suivant utilise un identifiant de service pour accéder à Google Cloud Storage à partir d'une UDF Python Batch de Unity Catalog :
%sql
CREATE OR REPLACE FUNCTION main.test.read_gcs_blob(blob_name STRING) RETURNS STRING LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'batchhandler'
CREDENTIALS (
`batch-udf-service-creds-example-cred` DEFAULT
)
ENVIRONMENT (
dependencies = '["google-auth", "google-cloud-storage"]', environment_version = 'None'
)
AS $$
import google.auth # This import is required to enable SDK credential integration
import pandas as pd
from google.cloud import storage
def batchhandler(it):
# The client automatically uses the DEFAULT service credential
client = storage.Client(project="your-project")
bucket = client.bucket("your-bucket")
for blob_names in it:
results = []
for name in blob_names:
blob = bucket.blob(name)
try:
content = blob.download_as_text()
results.append(content)
except Exception as e:
results.append(f"Error: {e}")
yield pd.Series(results)
$$;
Appelez l'UDF après son enregistrement :
SELECT main.test.read_gcs_blob(blob_name)
FROM VALUES
('config/settings.json'),
('data/input.txt')
AS t(blob_name)
Obtenir le contexte d'exécution de la tâche
Utilisez l'API PySpark TaskContext pour obtenir des informations contextuelles telles que l'identité de l'utilisateur, les cluster tags, l'ID du job Spark et plus encore. Consultez Obtenir le contexte de tâche dans une UDF.
Définissez DETERMINISTIC si votre fonction produit des résultats cohérents
Ajoutez DETERMINISTIC à la définition de votre fonction si elle produit les mêmes sorties pour les mêmes entrées. Cela permet aux optimisations de query d'améliorer les performances.
Par default, les UDTF Python Batch Unity Catalog sont supposées être non déterministes, sauf déclaration explicite. Les exemples de fonctions non déterministes incluent : la génération de valeurs aléatoires, l’accès aux heures ou dates actuelles, ou l’exécution d’appels d’API externes.
Voir CREATE FUNCTION (SQL, Python, Scala et Java)
Limitations
- Les fonctions Python doivent gérer les valeurs
NULLindépendamment, et toutes les mises en correspondance de types doivent suivre les mises en correspondance linguistiques de Databricks SQL. - Les UDF Python Unity Catalog Batch s'exécutent dans un environnement sécurisé et isolé, et n'ont pas accès à un système de fichiers partagé ou à des services internes.
- Plusieurs invocations d'UDF au sein d'une étape sont sérialisées et les résultats intermédiaires sont matérialisés et peuvent spill sur le disque.