Aller au contenu principal

Fonctions définies par l'utilisateur dans Databricks Connect pour Python

remarque

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.

remarque

É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 :

Par exemple, le Python suivant configure une simple UDF qui met au carré les valeurs d’une colonne.

Python
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

info

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<1 ou simplejson==3.19.*.
  • 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.whl ou dbfs:/Volumes/users/someone@example.com/tars/my_private_dep-4.0.0.tar.gz.
    • L'utilisateur doit disposer de l'autorisation READ_FILE sur le fichier dans le volume re:[UC]. L'octroi de cette autorisation à tous les utilisateurs du compte l'active automatiquement pour les nouveaux utilisateurs.
  • 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 comme local:<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.py ou local:./path/to/my_file.py.

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 :

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

Python
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.foreach né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.foreachBatch né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 :

Python
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

info

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 :

Python
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 sur True, 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 sur True, 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.eventhub et google.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 :

Python
# my_helper.py
def double(x):
return 2 * x
Python
# 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 :

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

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.

Version de Databricks Connect

Environnement serverless UDF

18,0, Python 3.12

Version 5

de 17,2 à 17,3, Python 3.12

Version 4

16.4.1 à moins de 17, Python 3.12

Version 3

15.4.10 à inférieur à 16, Python 3.12

Version 3

De 15.4.10 à moins de 16, Python 3.11

Version 2

Version de Databricks Connect

Environnement serverless UDF

18,0, Python 3.12

Version 5

de 17,2 à 17,3, Python 3.12

Version 4

16.4.1 à moins de 17, Python 3.12

Version 3

15.4.10 à inférieur à 16, Python 3.12

Version 3

De 15.4.10 à moins de 16, Python 3.11

Version 2