Fonctions définies par l'utilisateur dans Databricks Connect pour Python
Cet article couvre Databricks Connect pour Databricks Runtime 13.3 et versions ultérieures.
Databricks Connect pour Python prend en charge les fonctions définies par l'utilisateur (UDF). Lorsqu'une opération DataFrame qui inclut des UDF est exécutée, les UDF sont sérialisées par Databricks Connect et envoyées au serveur dans le cadre de la requête.
Pour plus d'informations sur les UDF pour Databricks Connect pour Scala, consultez Fonctions définies par l'utilisateur dans Databricks Connect pour Scala.
Étant donné que la fonction définie par l'utilisateur est sérialisée et désérialisée, la version Python du client doit correspondre à la version Python du compute Databricks. Pour les versions prises en charge, consultez la matrice de prise en charge des versions.
Définir une UDF
Pour créer une UDF dans Databricks Connect pour Python, utilisez l'une des fonctions prises en charge suivantes :
-
Fonctions définies par l'utilisateur PySpark
-
Fonctions de streaming PySpark
Par exemple, le Python suivant configure une simple UDF qui met au carré les valeurs d’une colonne.
from pyspark.sql.functions import col, udf
from pyspark.sql.types import IntegerType
from databricks.connect import DatabricksSession
@udf(returnType=IntegerType())
def double(x):
return x * x
spark = DatabricksSession.builder.getOrCreate()
df = spark.range(1, 2)
df = df.withColumn("doubled", double(col("id")))
df.show()
Gérer les dépendances UDF
Aperçu
Cette fonctionnalité est en Aperçu public et nécessite Databricks Connect pour Python 16.4 ou version ultérieure, ainsi qu'un cluster exécutant Databricks Runtime 16.4 ou version ultérieure. Pour utiliser cette fonctionnalité, activez l'aperçu **Enhanced Python UDFs in Unity Catalog** dans votre workspace.
Databricks Connect prend en charge la spécification des dépendances Python requises pour les UDF. Ces dépendances sont installées sur le compute Databricks dans le cadre de l'environnement Python des UDF.
Cette fonctionnalité permet aux utilisateurs de spécifier les dépendances dont l'UDF a besoin en plus des packages fournis dans l'environnement de base. Il peut également être utilisé pour installer une version du package différente de celle fournie dans l’ environnement de base.
Les dépendances peuvent être installées à partir des sources suivantes :
-
packages PyPI
- Les packages PyPI peuvent être spécifiés selon PEP 508, par exemple,
dice,pyjokes<1ousimplejson==3.19.*.
- Les packages PyPI peuvent être spécifiés selon PEP 508, par exemple,
-
Packages stockés dans les volumes Unity Catalog
- Les distributions intégrées (
.whl) et les distributions sources (.tar.gz) sont prises en charge. - Les packages de volumes Unity Catalog peuvent être spécifiés comme
dbfs:<path>, par exemple,dbfs:/Volumes/users/someone@example.com/wheels/my_private_dep-3.20.2-py3-none-any.whloudbfs:/Volumes/users/someone@example.com/tars/my_private_dep-4.0.0.tar.gz. - L'utilisateur doit disposer de l'autorisation
READ_FILEsur le fichier dans le volume re:[UC]. L'octroi de cette autorisation à tous les utilisateurs du compte l'active automatiquement pour les nouveaux utilisateurs.
- Les distributions intégrées (
-
Packages locaux, dossiers et fichiers Python
- Les distributions locales compilées (
.whl), les distributions source (.tar.gz), les dossiers et les fichiers Python peuvent être spécifiés commelocal:<path>, par exemple :local:/path/to/my_private_dep-3.20.2-py3-none-any.whl,local:/path/to/my_private_dep-4.0.0.tar.gz,local:/path/to/my_folder,local:/path/to/my_file.py. - Les chemins absolus et relatifs sont pris en charge, par exemple :
local:/path/to/my_file.pyoulocal:./path/to/my_file.py.
- Les distributions locales compilées (
Pour inclure des dépendances personnalisées dans votre UDF, spécifiez-les dans un environnement à l'aide de withDependencies, puis utilisez cet environnement pour créer une session Spark. Les dépendances sont installées sur votre compute Databricks et seront disponibles dans toutes les UDF qui utilisent cette session Spark.
Le code suivant déclare le package PyPI dice comme dépendance :
from databricks.connect import DatabricksSession, DatabricksEnv
env = DatabricksEnv().withDependencies("dice==3.1.0")
spark = DatabricksSession.builder.withEnvironment(env).getOrCreate()
Ou, pour spécifier une dépendance de wheel dans un volume :
from databricks.connect import DatabricksSession, DatabricksEnv
env = DatabricksEnv().withDependencies("/Volumes/users/someone@example.com/wheels/my_private_dep-3.20.2-py3-none-any.whl")
spark = DatabricksSession.builder.withEnvironment(env).getOrCreate()
Comportement dans les notebooks et jobs Databricks
Dans les Notebooks et les Jobs, les dépendances UDF doivent être installées directement dans le REPL. Databricks Connect valide l'environnement Python REPL en vérifiant que toutes les dépendances spécifiées sont déjà installées et lève une exception si certaines ne le sont pas. La validation de l'environnement Notebook s'exécute pour les dépendances de volume PyPI et Unity Catalog, mais pas pour les dépendances locales.
Limitations
- La prise en charge des dépendances UDF pour
pyspark.sql.streaming.DataStreamWriter.foreachnécessite Databricks Connect pour Python 18.0 ou version ultérieure, et un cluster exécutant Databricks Runtime 18.0 ou version ultérieure. - La prise en charge des dépendances UDF pour
pyspark.sql.streaming.DataStreamWriter.foreachBatchnécessite Databricks Connect pour Python 18.0 ou version ultérieure, et un cluster exécutant Databricks Runtime 18.0 ou version ultérieure. La fonctionnalité n'est pas prise en charge sur Serverless. - La prise en charge des dépendances UDF pour les packages locaux, les dossiers et les fichiers Python nécessite Databricks Connect pour Python 18,1 ou une version ultérieure, et un cluster exécutant Databricks Runtime 18,1 ou une version ultérieure.
- Les dépendances UDF ne sont pas prises en charge pour les UDF d’agrégation pandas sur les fonctions de fenêtre.
- Les packages de volumes Unity Catalog et les packages locaux doivent être empaquetés conformément aux spécifications d'empaquetage Python standard de PEP-427 ou ultérieures pour les distributions construites en roue (wheel) et PEP-241 ou ultérieures pour les distributions sources tar. Pour plus d'informations sur les normes d'empaquetage Python, consultez la documentation PyPA.
Exemples
L'exemple suivant définit les dépendances PyPI et de volumes dans un environnement, crée une session avec cet environnement, puis définit et appelle des fonctions UDF qui utilisent ces dépendances :
from databricks.connect import DatabricksSession, DatabricksEnv
from pyspark.sql.functions import udf, col, pandas_udf
from pyspark.sql.types import IntegerType, LongType, StringType
import pandas as pd
pypi_deps = ["pyjokes>=0.8,<1"]
volumes_deps = [
# Example library from: https://pypi.org/project/dice/#files
"/Volumes/main/someone@example.com/test/dice-4.0.0.tar.gz",
]
local_deps = [
# Example library from: https://pypi.org/project/simplejson/#files
"local:./test/simplejson-3.20.2-py3-none-any.whl",
]
env = DatabricksEnv().withDependencies(pypi_deps).withDependencies(volumes_deps).withDependencies(local_deps)
spark = DatabricksSession.builder.withEnvironment(env).getOrCreate()
# UDFs
@udf(returnType=StringType())
def get_joke():
from pyjokes import get_joke
return get_joke()
@udf(returnType=IntegerType())
def double_and_json_parse(x):
import simplejson
return simplejson.loads(simplejson.dumps(x * 2))
@pandas_udf(returnType=LongType())
def multiply_and_add_roll(a: pd.Series, b: pd.Series) -> pd.Series:
import dice
return a * b + dice.roll(f"1d10")[0]
df = spark.range(1, 10)
df = df.withColumns({
"joke": get_joke(),
"doubled": double_and_json_parse(col("id")),
"mutliplied_with_roll": multiply_and_add_roll(col("id"), col("doubled"))
})
df.show()
Gestion automatique des dépendances UDF
Aperçu
Cette fonctionnalité est en Préversion publique et nécessite Databricks Connect pour Python 18.1 ou version ultérieure, Python 3.12 sur votre machine locale et un cluster exécutant Databricks Runtime 18.1 ou version ultérieure. Pour utiliser cette fonctionnalité, activez l'aperçu **Enhanced Python UDFs in Unity Catalog** dans votre workspace.
L'API Databricks Connect withAutoDependencies() permet la découverte et l'upload automatiques des modules locaux et des dépendances PyPI publiques utilisées dans les instructions d'importation de vos UDF. Cela élimine le besoin de spécifier manuellement les dépendances.
Le code suivant active la gestion automatique des dépendances :
from databricks.connect import DatabricksSession, DatabricksEnv
env = DatabricksEnv().withAutoDependencies(upload_local=True, use_index=True)
spark = DatabricksSession.builder.withEnvironment(env).getOrCreate()
La méthode withAutoDependencies() accepte les paramètres suivants :
upload_local: Lorsque défini surTrue, les modules locaux importés dans vos UDF sont automatiquement découverts, mis en package et upload vers le sandbox UDF.use_index: lorsqu'elle est définie surTrue, les dépendances PyPI publiques utilisées dans vos UDF sont automatiquement découvertes et installées sur le compute Databricks. Le processus de découverte utilise les packages installés sur votre machine locale pour déterminer les versions, garantissant la cohérence entre votre environnement local et l'environnement d'exécution distant.
Limitations
- Les importations dynamiques (par exemple,
importlib.import_module("foo")) ne sont pas prises en charge. - Les packages d'espace de noms (par exemple,
azure.eventhubetgoogle.cloud.aiplatform) ne sont pas pris en charge. - Les dépendances installées à l'aide de références d'URL directes ne sont pas prises en charge. Cela inclut ceux installés à partir de fichiers wheel locaux.
- Les dépendances installées à partir d'index de packages privés ne sont pas prises en charge. Les packages installés de cette manière ne peuvent pas être distingués des packages installés depuis le PyPI public.
- La détection des dépendances ne fonctionne pas dans un Shell Python. Seuls les scripts Python, le shell IPython et les Notebooks Jupyter sont pris en charge.
Exemples
L'exemple suivant démontre la gestion automatique des dépendances avec des modules locaux et des packages PyPI. Cet exemple nécessite que vous ayez installé simplejson et dice (en utilisant pip install simplejson dice).
D'abord, créez des modules d'aide locaux :
# my_helper.py
def double(x):
return 2 * x
# my_json.py
import simplejson
def loads(x):
return simplejson.loads(x)
def dumps(x):
return simplejson.dumps(x)
Ensuite, dans votre script principal, importez ces modules et utilisez-les dans les UDF :
# main.py
import dice as dc
from databricks.connect import DatabricksSession, DatabricksEnv
from pyspark.sql.functions import col, udf
from pyspark.sql.types import IntegerType, FloatType
import my_json
from my_helper import double
env = DatabricksEnv().withAutoDependencies(upload_local=True, use_index=True)
spark = DatabricksSession.builder.withEnvironment(env).getOrCreate()
@udf(returnType=IntegerType())
def double_and_json_parse(x):
return my_json.loads(my_json.dumps(double(x)))
@udf(returnType=FloatType())
def sum_and_add_noise(x, y):
return x + y + (dc.roll("d6")[0] / 6)
df = spark.range(1, 10)
df = df.withColumns({
"doubled": double_and_json_parse(col("id")),
"summed_with_noise": sum_and_add_noise(col("id"), col("doubled")),
})
df.show()
Journalisation
Pour afficher les dépendances détectées, définissez la variable d’environnement SPARK_CONNECT_LOG_LEVEL sur info ou debug. Par ailleurs, configurez le framework de journalisation Python :
import logging
logging.basicConfig(level=logging.INFO)
Les Logs pertinents sont émis par le module databricks.connect.auto_dependencies, par exemple :
DEBUG:databricks.connect.auto_dependencies.discovery:Discovered local module: my_json
DEBUG:databricks.connect.auto_dependencies.discovery:Discovered local module: my_helper
DEBUG:databricks.connect.auto_dependencies.discovery:Discovered distribution: simplejson for module simplejson
DEBUG:databricks.connect.auto_dependencies.discovery:Discovered distribution: dice for module dice
INFO:databricks.connect.auto_dependencies.hook:Synced zip artifact for: my_json
INFO:databricks.connect.auto_dependencies.hook:Synced zip artifact for: my_helper
INFO:databricks.connect.auto_dependencies.hook:Updated simplejson with auto-detected version ==3.20.2
INFO:databricks.connect.auto_dependencies.hook:Updated dice with auto-detected version ==4.0.0
Environnement de base Python
Les fonctions UDF sont exécutées sur le compute Databricks et non sur le client. L'environnement Python de base dans lequel les UDF sont exécutées dépend du compute Databricks.
Pour les clusters, l'environnement Python de base est l'environnement Python de la version de Databricks Runtime exécutée sur le cluster. La version de Python et la liste des packages de cet environnement de base se trouvent dans les sections Environnement système et Bibliothèques Python installées des notes de publication de Databricks Runtime.
Pour le compute Serverless, l'environnement Python de base correspond à la version de l'environnement Serverless selon le tableau suivant. Les versions de Databricks Connect non répertoriées dans ce tableau ne prennent pas encore en charge le mode Serverless ou ont atteint la fin du support. Consultez la matrice de support des versions et les versions de Databricks Connect en fin de support.