Aller au contenu principal

Opérateurs définis par l’utilisateur dans Lakeflow Designer

Lakeflow Designer vous permet de créer des opérateurs définis par l'utilisateur qui apparaissent directement sur le canevas à côté des opérateurs intégrés. Utilisez-les pour étendre Lakeflow Designer avec votre propre logique métier, vos calculs ou vos intégrations.

Il existe trois types d'opérateurs définis par l'utilisateur :

  • python-run-function : Un fichier YAML autonome avec du Python intégré stocké dans le Workspace. Idéal pour les transformations au niveau du DataFrame et les intégrations externes. Les autorisations sont gérées au niveau des fichiers du Workspace.
  • uc-udf : Encapsule une fonction scalaire Unity Catalog. Idéal pour les transformations au niveau des colonnes. L'accès est régi par les autorisations d'Unity Catalog.
  • uc-udtf : Encapsule une fonction table Unity Catalog. Idéal pour les transformations au niveau de la table comme le clustering ML et l'agrégation. L'accès est régi par les autorisations d'Unity Catalog.

Fonctionnalité

python-run-function

uc-udf

uc-udtf

Exemple de cas d'usage

Transformations de DataFrame, intégrations d'API, notifications e-mail

Calculs au niveau des colonnes (IMC, taux d'intérêt)

Clustering ML, agrégation de lignes

Entrée

DataFrames

Valeurs uniques

Table entière, ligne par ligne

Résultat

DataFrames

Valeur unique

Table (plusieurs lignes)

Requiert la fonction Unity Catalog

Non

Oui

Oui

Gouvernance de l'accès

Permissions des fichiers Workspace

Autorisations Unity Catalog (EXECUTE, USE SCHEMA)

Autorisations Unity Catalog (EXECUTE, USE SCHEMA)

Langues prises en charge

Python uniquement

SQL ou Python dans un wrapper SQL.

SQL ou Python dans un wrapper SQL.

Fonctionnalité

python-run-function

uc-udf

uc-udtf

Exemple de cas d'usage

Transformations de DataFrame, intégrations d'API, notifications e-mail

Calculs au niveau des colonnes (IMC, taux d'intérêt)

Clustering ML, agrégation de lignes

Entrée

DataFrames

Valeurs uniques

Table entière, ligne par ligne

Résultat

DataFrames

Valeur unique

Table (plusieurs lignes)

Requiert la fonction Unity Catalog

Non

Oui

Oui

Gouvernance de l'accès

Permissions des fichiers Workspace

Autorisations Unity Catalog (EXECUTE, USE SCHEMA)

Autorisations Unity Catalog (EXECUTE, USE SCHEMA)

Langues prises en charge

Python uniquement

SQL ou Python dans un wrapper SQL.

SQL ou Python dans un wrapper SQL.

Fonctionnement des opérateurs définis par l'utilisateur

Un opérateur défini par l'utilisateur se compose de :

  • **Logique d'opérateur** : Le code qui s'exécute lorsque l'opérateur est exécuté. Il peut s'agir d'une fonction Python run() en ligne (pour python-run-function) ou d'une fonction Unity Catalog (pour uc-udf et uc-udtf).
  • Configuration YAML : Indique à Lakeflow Designer comment présenter l'opérateur dans l'interface utilisateur, notamment son nom, sa description, ses paramètres d'entrée, ses widgets d'interface utilisateur et ses ports. Tous les types d'opérateur utilisent le schéma user-defined-operator-v0.1.0.
  • Fichier d'enregistrement : une entrée dans .user_defined_operators.yaml qui permet à Lakeflow Designer de découvrir l'opérateur.

Logique de l'opérateur

Logique d'opérateur défini par l'utilisateur de fonction d'exécution Python

Chaque opérateur python-run-function doit définir une fonction run() :

Python
def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
  • config : Valeurs configurées par l'utilisateur à partir de l'interface utilisateur, indexées par le nom de la propriété.
  • inputs : DataFrames d’entrée, indexés par le port name.
  • spark : La SparkSession active.
  • Retourne : un dictionnaire associant les valeurs du port de sortie name à des DataFrames.

L’exemple suivant filtre les lignes d’un DataFrame d’entrée :

Python
def run(config, inputs, spark):
df = inputs["in"]
filtered = df.filter(config["filter_expression"])
return {"out": filtered}

Si votre opérateur nécessite des packages pip externes, ajoutez le champ environment au YAML :

YAML
environment:
environment_version: '4'
dependencies:
- requests==2.31.0
- beautifulsoup4==4.12.0

Logique de l'opérateur UDF et UDTF

Vous pouvez écrire des fonctions UC en SQL ou en Python. Les fonctions Python sont encapsulées dans une instruction SQL CREATE FUNCTION :

Fonction SQL :

SQL
CREATE OR REPLACE FUNCTION my_catalog.my_schema.calculate_bmi(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE SQL
RETURN
SELECT weight_kg / (height_m * height_m);

Fonction Python (encapsulée dans SQL) :

SQL
CREATE OR REPLACE FUNCTION my_catalog.my_schema.calculate_bmi(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
AS $$
return weight_kg / (height_m ** 2)
$$;

Les UDF traitent une seule valeur à la fois et renvoient une valeur calculée. Les UDTF traitent les tables ligne par ligne et peuvent maintenir l'état sur toutes les lignes. Utilisez uc-udf pour les transformations au niveau des colonnes et uc-udtf pour les Opérations telles que le clustering ML ou l'agrégation.

En outre, les UDTF vous obligent à définir trois méthodes clés : __init__(), eval() et terminate():

Python
class MyOperator:
def __init__(self):
# Called before processing - initialize any values needed.

def eval(self, row, id_column, columns, k):
# Called one time per input row - accumulate data here.

def terminate(self):
# Called after all rows - perform final calculations and yield results.
remarque

Les tables de retour UDTF doivent avoir des types fixes et explicites. Vous ne pouvez pas référencer les types de colonnes d'entrée dans la configuration de retour.

Configuration YAML

La configuration YAML indique à Lakeflow Designer comment présenter l'opérateur dans l'interface utilisateur. Il définit le nom de l'opérateur, la description, les paramètres d'entrée, les widgets de l'interface utilisateur et les ports. Chaque champ de configuration est une propriété avec un type, un titre et des indices de widget x-ui facultatifs :

YAML
config:
type: object
properties:
my_param:
type: string
title: My Parameter
x-ui:
widget: input
my_expression:
type: string
title: Column
format: expression
x-ui:
widget: expression
port: in
my_number:
type: number
title: Count
default: 10
minimum: 0
maximum: 100
required:
- my_param
- my_expression

Pour plus de détails sur le schéma YAML, y compris tous les types de widgets et les options de configuration, consultez la référence YAML de l'opérateur défini par l'utilisateur.

Ports

Les ports définissent les entrées et les sorties pour votre opérateur :

YAML
ports:
input:
- name: in
title: Input Data
mime: application/vnd.databricks.dataframe
required: true
allowMultiple: false
output:
- name: out
title: Output Data

YAML pour les opérateurs de fonctions d'exécution Python

Pour les opérateurs python-run-function, le fichier YAML est autonome et inclut un champ run_function avec du code Python inline :

YAML
schema: user-defined-operator-v0.1.0
type: python-run-function
name: Filter Rows
id: filter_rows
version: '1.0.0'
description: Filters rows based on a SQL expression.
config:
type: object
properties:
filter_expression:
type: string
title: Filter Expression
x-ui:
widget: input
required:
- filter_expression
ports:
input:
- name: in
title: Input
output:
- name: out
title: Output
run_function:
type: inline
code: |
def run(config, inputs, spark):
df = inputs["in"]
filtered = df.filter(config["filter_expression"])
return {"out": filtered}

YAML pour les fonctions Unity Catalog

Pour les opérateurs basés sur UC, intégrez la configuration YAML sous forme de commentaire ou de docstring dans votre fonction.

Dans SQL (utilisez le commentaire /* ... */) :

SQL
RETURN(/*
schema: user-defined-operator-v0.1.0
type: uc-udf
name: Calculate BMI
id: calculate_bmi
version: "1.0.0"
description: Calculates BMI from weight and height.
config:
type: object
properties:
weight_kg:
type: string
title: Weight (in kg)
format: expression
x-ui:
widget: expression
port: in
height_m:
type: string
title: Height (in meters)
format: expression
x-ui:
widget: expression
port: in
required:
- weight_kg
- height_m
ports:
input:
- name: in
title: Input Data
output:
- name: out
title: Output
*/
SELECT weight_kg / (height_m * height_m)
);

En Python (utilisez la """ ... """ docstring) :

SQL
AS $$
"""
schema: user-defined-operator-v0.1.0
type: uc-udf
name: Calculate BMI
id: calculate_bmi
version: "1.0.0"
description: Calculates BMI from weight and height.
config:
type: object
properties:
weight_kg:
type: string
title: Weight (in kg)
format: expression
x-ui:
widget: expression
port: in
height_m:
type: string
title: Height (in meters)
format: expression
x-ui:
widget: expression
port: in
required:
- weight_kg
- height_m
ports:
input:
- name: in
title: Input Data
output:
- name: out
title: Output
"""

return weight_kg / (height_m ** 2)
$$;

Enregistrer et déployer votre opérateur dans Lakeflow Designer

Pour que votre opérateur apparaisse dans Lakeflow Designer, enregistrez-le dans un fichier .user_defined_operators.yaml :

  • Niveau Workspace : Placez le fichier à la racine de votre Workspace pour rendre l'opérateur visible à tous les utilisateurs.
  • Niveau utilisateur : Placez le fichier dans votre dossier personnel (/Workspace/Users/<user-name>/.user_defined_operators.yaml) pour que les opérateurs ne soient visibles que pour vous.

La section operators: prend en charge les chemins de fichiers, les références de fonction Unity Catalog et les modèles glob. Vous pouvez mélanger les types d'entrée :

YAML
operators:
# File path (python-run-function operators)
- /Workspace/Users/me/udos/my_operator.yaml
# Glob pattern (registers all matching files)
- /Workspace/Users/me/udos/transforms/*.yaml
# UC function reference (uc-udf and uc-udtf operators)
- catalog: my_catalog
schema: my_schema
functionName: my_function

Mettre à jour ou supprimer un opérateur

Lorsque vous modifiez le code d'un opérateur, refresh vos opérateurs définis par l’utilisateur pour charger la modification. Dans le tab **Opérateurs** du menu, cliquez Icône refresh. sur.

  • Si l'opérateur conserve le même version, l'actualisation charge le code mis à jour.
  • Si l’opérateur possède un nouveau version, l’opérateur sur la toile vous invite à le mettre à niveau (ou à conserver la version actuelle) après le refresh.

Pour supprimer un opérateur de Lakeflow Designer, supprimez son entrée de .user_defined_operators.yaml. Pour les opérateurs uc-udf et uc-udtf, vous pouvez également supprimer la fonction Unity Catalog sous-jacente avec DROP FUNCTION si vous n'en avez plus besoin.

Configurations avancées

Mode aperçu

Lakeflow Designer prend en charge les aperçus en mode conception. Pour les opérateurs qui appellent des APIs externes ou écrivent dans des systèmes externes, ajoutez une propriété de configuration is_preview afin que vous puissiez ignorer les effets secondaires pendant l'aperçu. Lorsque le mode d'aperçu est activé, les utilisateurs doivent cliquer explicitement sur Exécuter pour exécuter l'opérateur avec des effets secondaires.

YAML
config:
type: object
properties:
is_preview:
type: boolean
format: is_preview
default: false

Lakeflow Designer définit automatiquement cette valeur à true pendant l'aperçu. Vérifiez-le dans votre logique pour ignorer les effets secondaires :

Python
# In a python-run-function
if config.get("is_preview"):
return {"out": inputs["in"]}

# In a UC function (SQL)
CASE WHEN is_preview THEN 'preview' ELSE /* actual work */ END

Connexions Unity Catalog

Pour les opérateurs SQL basés sur UC qui appellent des APIs externes, utilisez les connexions HTTP Unity Catalog pour stocker les informations d'identification en toute sécurité :

SQL
CREATE CONNECTION my_api_connection TYPE HTTP OPTIONS (
host 'https://api.example.com',
port '443',
base_path '/v1/',
bearer_token 'your-token-here'
);

Utilisez ensuite la connexion dans votre UDF SQL avec la fonction http_request(). Pour plus de détails, voir Se connecter aux services HTTP externes.

WorkspaceClient

Pour les opérateurs python-run-function, vous pouvez utiliser le Databricks WorkspaceClient pour accéder aux ressources du workspace et aux APIs externes :

Python
def run(config, inputs, spark):
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
# Use w to access workspace resources

Créez un opérateur défini par l'utilisateur python-run-function complet

Les étapes suivantes décrivent la création d'un opérateur python-run-function à partir de zéro.

Étape 1 : définir la logique

Écrivez votre fonction run() dans un Notebook :

Python
from typing import Dict, Any

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
from pyspark.sql import functions as F
df = inputs["in"]
result = df.withColumn(config["column_name"], F.current_timestamp())
return {"out": result}

Étape 2 : tester la fonction

Tester la fonction interactivement avec des exemples de données :

Python
test_df = spark.createDataFrame(
[("Alice", 100), ("Bob", 200)],
["name", "amount"]
)

result = run(
config={&quot;column_name&quot;: &quot;processed_at&quot;},
inputs={&quot;in&quot;: test_df},
spark=spark
)

result["out"].show()

Étape 3 : Créer la configuration YAML

Définissez les métadonnées de l'opérateur, les champs de configuration et les ports dans un fichier YAML :

YAML
schema: user-defined-operator-v0.1.0
type: python-run-function
name: Add Timestamp
id: transforms.add_timestamp
version: '1.0.0'
description: Adds a timestamp column to the input DataFrame.
config:
type: object
properties:
column_name:
type: string
title: Column Name
default: processed_at
x-ui:
widget: input
required:
- column_name

Étape 4 : Combinez la logique et YAML

Ajoutez les champs run_function et ports pour créer le fichier YAML complet. Enregistrez-le dans votre Workspace, par exemple /Workspace/Users/<user-name>/udos/add_timestamp.yaml:

YAML
schema: user-defined-operator-v0.1.0
type: python-run-function
name: Add Timestamp
id: transforms.add_timestamp
version: '1.0.0'
description: Adds a timestamp column to the input DataFrame.
config:
type: object
properties:
column_name:
type: string
title: Column Name
default: processed_at
x-ui:
widget: input
required:
- column_name
ports:
input:
- name: in
title: Input
output:
- name: out
title: Output
run_function:
type: inline
code: |
from typing import Dict, Any

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
from pyspark.sql import functions as F
df = inputs["in"]
result = df.withColumn(config["column_name"], F.current_timestamp())
return {"out": result}

Étape 5 : enregistrer l'opérateur

Ajoutez le chemin d’accès au fichier à votre fichier .user_defined_operators.yaml :

YAML
operators:
- /Workspace/Users/<user-name>/udos/add_timestamp.yaml

Étape 6 : Utiliser l'opérateur dans Lakeflow Designer

Ouvrez Lakeflow Designer et vérifiez que l'opérateur apparaît dans la palette d'opérateurs. Faites-le glisser sur la toile, connectez une entrée, configurez le nom de la colonne et exécutez un aperçu.

Créer un opérateur défini par l'utilisateur UC complet

Les étapes suivantes décrivent la création d'un opérateur uc-udf basé sur Unity Catalog.

Étape 1 : définir la logique

Écrivez et testez la logique de votre fonction dans un Notebook :

Python
def double_value(input_value: float) -> float:
if input_value is None:
return None
return input_value * 2

Étape 2 : Créez la configuration YAML

Définissez les métadonnées de l'opérateur, les champs de configuration et les ports :

YAML
schema: user-defined-operator-v0.1.0
type: uc-udf
name: Double Value
id: math.double_value
version: '1.0.0'
description: Doubles the input value
config:
type: object
properties:
input_value:
type: string
title: Input Value
format: expression
x-ui:
widget: expression
port: input_data
required:
- input_value
ports:
input:
- name: input_data
title: Input
output:
- name: out
title: Output

Étape 3 : combiner la logique et le YAML

Créez la fonction Unity Catalog avec le YAML intégré en tant que docstring :

SQL
CREATE OR REPLACE FUNCTION main.my_schema.double_value(input_value DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
AS $$
"""
schema: user-defined-operator-v0.1.0
type: uc-udf
name: Double Value
id: math.double_value
version: "1.0.0"
description: Doubles the input value
config:
type: object
properties:
input_value:
type: string
title: Input Value
format: expression
x-ui:
widget: expression
port: input_data
required:
- input_value
ports:
input:
- name: input_data
title: Input
output:
- name: out
title: Output
"""

def double_value(input_value: float) -> float:
if input_value is None:
return None
return input_value * 2

return double_value(input_value)
$$

Étape 4 : Testez la fonction

SQL
SELECT main.my_schema.double_value(5) AS result;
-- Should return: 10

Étape 5 : enregistrer l'opérateur

Ajoutez la référence de la fonction Unity Catalog à votre fichier .user_defined_operators.yaml :

YAML
operators:
- catalog: main
schema: my_schema
functionName: double_value

Étape 6 : Utiliser l'opérateur dans Lakeflow Designer

Ouvrez Lakeflow Designer et vérifiez que l'opérateur apparaît dans la palette d'opérateurs. Faites-le glisser sur la toile, connectez une entrée et exécutez un aperçu.

Dépannage

Problème

Solutions

L'opérateur n'apparaît pas dans Lakeflow Designer.

Vérifiez que .user_defined_operators.yaml existe et répertorie votre fonction ou votre chemin de fichier. Pour les opérateurs python-run-function, vérifiez le chemin d’accès du fichier et que le fichier YAML est accessible.

La validation du schéma échoue.

Vérifiez votre YAML par rapport au schéma officiel à https://your-workspace.cloud.databricks.com/static/schemas/user-defined-operator-v0.1.0.json.

Autorisation refusée.

Pour les opérateurs basés sur UC, vérifiez que les utilisateurs disposent de EXECUTE sur la fonction et de USE SCHEMA sur le schéma. Pour les opérateurs python-run-function, vérifiez que les utilisateurs disposent d'un accès en lecture au fichier YAML.

python-run-function L'opérateur échoue au moment de l'exécution.

Vérifiez que la signature de la fonction run() correspond à def run(config, inputs, spark). Vérifiez que les noms de port dans le code correspondent au YAML et que les clés du dictionnaire de retour correspondent aux valeurs du port de sortie name.

UDTF renvoie des types incorrects.

Les types de retour UDTF doivent être explicites ; vous ne pouvez pas référencer les types de colonne d'entrée.

Problème

Solutions

L'opérateur n'apparaît pas dans Lakeflow Designer.

Vérifiez que .user_defined_operators.yaml existe et répertorie votre fonction ou votre chemin de fichier. Pour les opérateurs python-run-function, vérifiez le chemin d’accès du fichier et que le fichier YAML est accessible.

La validation du schéma échoue.

Vérifiez votre YAML par rapport au schéma officiel à https://your-workspace.cloud.databricks.com/static/schemas/user-defined-operator-v0.1.0.json.

Autorisation refusée.

Pour les opérateurs basés sur UC, vérifiez que les utilisateurs disposent de EXECUTE sur la fonction et de USE SCHEMA sur le schéma. Pour les opérateurs python-run-function, vérifiez que les utilisateurs disposent d'un accès en lecture au fichier YAML.

python-run-function L'opérateur échoue au moment de l'exécution.

Vérifiez que la signature de la fonction run() correspond à def run(config, inputs, spark). Vérifiez que les noms de port dans le code correspondent au YAML et que les clés du dictionnaire de retour correspondent aux valeurs du port de sortie name.

UDTF renvoie des types incorrects.

Les types de retour UDTF doivent être explicites ; vous ne pouvez pas référencer les types de colonne d'entrée.

Autorisations

Autorisation

Objectif

Accès en lecture .user_defined_operators.yamlà.

Découvrir l'opérateur.

**Accès en lecture** au fichier YAML (python-run-function uniquement).

Chargez la définition de l'opérateur.

EXÉCUTER sur la fonction Unity Catalog (opérateurs basés sur UC uniquement).

Exécutez l'opérateur.

**USE SCHEMA** sur le schéma (opérateurs basés sur UC uniquement).

Accédez au schéma où la fonction est créée.

Autres autorisations

Selon votre opérateur, les utilisateurs peuvent nécessiter d'autres autorisations. Par exemple, USE CONNECTION sur une connexion Unity Catalog pour les appels d'API HTTP.

Autorisation

Objectif

Accès en lecture .user_defined_operators.yamlà.

Découvrir l'opérateur.

**Accès en lecture** au fichier YAML (python-run-function uniquement).

Chargez la définition de l'opérateur.

EXÉCUTER sur la fonction Unity Catalog (opérateurs basés sur UC uniquement).

Exécutez l'opérateur.

**USE SCHEMA** sur le schéma (opérateurs basés sur UC uniquement).

Accédez au schéma où la fonction est créée.

Autres autorisations

Selon votre opérateur, les utilisateurs peuvent nécessiter d'autres autorisations. Par exemple, USE CONNECTION sur une connexion Unity Catalog pour les appels d'API HTTP.

Ressources supplémentaires

Explorez les tutoriels suivants :

Exemple

Type

Description

Expéditeur d'e-mail Gmail

python-run-function

Envoyez des données DataFrame en tant que pièce jointe e-mail CSV via Gmail.

Calculateur d'intérêts composés

uc-udf

Calculez les futures valeurs d'investissement à l'aide de la formule des intérêts composés.

Clustering K-means

uc-udtf

Segmenter les données en clusters à l'aide de scikit-learn.

Envoyer un message Slack

uc-udf

Envoyez des notifications aux canaux Slack via l'API.

Tous les widgets d'interface utilisateur

uc-udf

Opérateur de référence présentant tous les widgets d'interface utilisateur disponibles.

Exemple

Type

Description

Expéditeur d'e-mail Gmail

python-run-function

Envoyez des données DataFrame en tant que pièce jointe e-mail CSV via Gmail.

Calculateur d'intérêts composés

uc-udf

Calculez les futures valeurs d'investissement à l'aide de la formule des intérêts composés.

Clustering K-means

uc-udtf

Segmenter les données en clusters à l'aide de scikit-learn.

Envoyer un message Slack

uc-udf

Envoyez des notifications aux canaux Slack via l'API.

Tous les widgets d'interface utilisateur

uc-udf

Opérateur de référence présentant tous les widgets d'interface utilisateur disponibles.

Pour une référence complète du schéma YAML, consultez Référence YAML de l'opérateur défini par l'utilisateur.