Aller au contenu principal

Fonctions définies par l'utilisateur (UDF) Python en batch dans Unity Catalog

info

Aperçu

Cette fonctionnalité est en aperçu public.

Les UDF Python de Unity Catalog par batch étendent les capacités des UDF de Unity Catalog en vous permettant d'écrire du code Python pour opérer sur des lots de données, ce qui améliore considérablement l'efficacité en réduisant la surcharge associée aux UDF ligne par ligne. Ces optimisations rendent les UDF Python batch de Unity Catalog idéales pour le traitement de données à grande échelle.

Exigences

Les UDF Python batch Unity Catalog nécessitent Databricks Runtime versions 16.3 et supérieures.

Créer un UDF Python batch Unity Catalog

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': Ceci spécifie la fonction de gestionnaire qui est appelée pour traiter les batchs.

L'exemple suivant vous montre comment créer un UDF Python de Unity Catalog en mode batch :

Python
%sql
CREATE OR REPLACE TEMPORARY FUNCTION 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 :

SQL
SELECT person_id, 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 Unity Catalog de batch nécessitent une fonction de gestionnaire qui traite les batchs et produit des résultats. Vous devez spécifier le nom de la fonction de gestionnaire lorsque vous créez l'UDF en utilisant la clé HANDLER.

La fonction de gestionnaire effectue les opérations suivantes :

  1. Accepte un argument itérateur qui itère sur un ou plusieurs pandas.Series. Chaque pandas.Series contient les paramètres d'entrée de l'UDF.
  2. Itère sur le générateur et traite les données.
  3. 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.

Les UDF de batch peuvent accepter des paramètres uniques ou multiples

Paramètre unique : Lorsque la fonction gestionnaire utilise un seul paramètre d'entrée, elle reçoit un itérateur qui itère sur un pandas.Series pour chaque lot.

Python
%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.

Python
%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 :

Python
%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

remarque

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 :

SQL
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 batch Unity Catalog peuvent utiliser les informations d'identification de service Unity Catalog pour accéder aux services cloud externes. Ceci est particulièrement utile pour l'intégration de fonctions cloud comme les tokeniseurs de sécurité dans les workflows de traitement des données.

remarque

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 :

SQL
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

Le créateur d'UDF doit disposer de l'autorisation ACCESS sur les informations d'identification du service Unity Catalog. Cependant, pour les appelants d'UDF, il suffit de leur accorder l'autorisation EXECUTE sur l'UDF. En particulier, les appelants d'UDF n'ont pas besoin d'accéder aux informations d'identification de service sous-jacentes, car l'UDF s'exécute en utilisant les autorisations d'informations d'identification du créateur de l'UDF.

Pour les fonctions temporaires, le créateur est toujours l'invocateur. Les UDF qui s'exécutent dans l'étendue No-PE , également appelée clusters dédiés, utilisent plutôt les autorisations de l'appelant.

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.

Python
from databricks.service_credentials import getServiceCredentialsProvider
import boto3

# Assuming credential definition: CREDENTIALS(`aws-cred` AS testcred)
boto3_session = boto3.Session(botocore_session=getServiceCredentialsProvider('testcred'))
s3 = boto3_session.client('s3')

Exemple d'identifiant de service - fonction AWS Lambda

L'exemple suivant utilise un identifiant de service pour appeler une fonction AWS Lambda à partir d'une UDF Python de Unity Catalog par batch :

Python
%sql
CREATE OR REPLACE FUNCTION main.test.call_lambda_func(data STRING, debug BOOLEAN) RETURNS STRING LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'batchhandler'
CREDENTIALS (
`batch-udf-service-creds-example-cred` DEFAULT
)
AS $$
import boto3
import json
import pandas as pd
import base64
from pyspark.taskcontext import TaskContext


def batchhandler(it):
# Automatically picks up DEFAULT credential:
session = boto3.Session()

client = session.client("lambda", region_name="us-west-2")

# Propagate TaskContext information to lambda context:
user_ctx = {"custom": {"user": TaskContext.get().getLocalProperty("user")}}

for vals, is_debug in it:
payload = json.dumps({"values": vals.to_list(), "is_debug": bool(is_debug[0])})

res = client.invoke(
FunctionName="HashValuesFunction",
InvocationType="RequestResponse",
ClientContext=base64.b64encode(json.dumps(user_ctx).encode("utf-8")).decode(
"utf-8"
),
Payload=payload,
)

response_payload = json.loads(res["Payload"].read().decode("utf-8"))
if "errorMessage" in response_payload:
raise Exception(str(response_payload))

yield pd.Series(response_payload["values"])
$$;

Appelez l'UDF après son enregistrement :

SQL
SELECT main.test.call_lambda_func(data, false)
FROM VALUES
('abc'),
('def')
AS t(data)

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.

By default, les UDTF Python batch Unity Catalog sont considérées comme non déterministes sauf si elles sont explicitement déclarées. Les exemples de fonctions non déterministes incluent : générer des valeurs aléatoires, accéder aux heures ou aux dates actuelles, ou effectuer des appels d’API externes.

Voir CREATE FUNCTION (SQL, Python, Scala et Java)

Limitations

  • Les fonctions Python doivent gérer les valeurs NULL indé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.

  • Les identifiants de service sont disponibles uniquement dans les UDF Python Unity Catalog Batch et les UDF Python scalaires. Ils ne sont pas pris en charge dans les UDF Python Unity Catalog standards.

  • Sur les clusters dédiés et pour les fonctions temporaires, l'appelant de fonction doit disposer des autorisations ACCESS sur les identifiants de service. Consultez Accorder des autorisations pour utiliser un identifiant de service afin d'accéder à un service cloud externe.

  • Activez la fonctionnalité en Public Preview **Enable networking for isolated workloads in Serverless SQL Warehouses** dans la page d'aperçus de votre Workspace pour effectuer des appels d'UDF Python Unity Catalog en batch vers des services externes sur le compute de SQL warehouse serverless.

  • Pour effectuer des appels Batch Unity Catalog Python UDF sur un Notebook ou un compute Serverless, vous devez configurer le contrôle de sortie Serverless.