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é |
|
|
|
|---|---|---|---|
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 ( | Autorisations Unity Catalog ( |
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 (pourpython-run-function) ou d'une fonction Unity Catalog (pouruc-udfetuc-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.yamlqui 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() :
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 portname.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 :
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 :
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 :
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) :
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():
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.
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 :
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 :
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 :
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 /* ... */) :
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) :
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 :
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 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.
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 :
# 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é :
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 :
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 :
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 :
test_df = spark.createDataFrame(
[("Alice", 100), ("Bob", 200)],
["name", "amount"]
)
result = run(
config={"column_name": "processed_at"},
inputs={"in": 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 :
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:
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 :
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 :
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 :
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 :
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
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 :
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 |
La validation du schéma échoue. | Vérifiez votre YAML par rapport au schéma officiel à |
Autorisation refusée. | Pour les opérateurs basés sur UC, vérifiez que les utilisateurs disposent de |
| Vérifiez que la signature de la fonction |
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 | Découvrir l'opérateur. |
**Accès en lecture** au fichier YAML ( | 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, |
Ressources supplémentaires
Explorez les tutoriels suivants :
Exemple | Type | Description |
|---|---|---|
| Envoyez des données DataFrame en tant que pièce jointe e-mail CSV via Gmail. | |
| Calculez les futures valeurs d'investissement à l'aide de la formule des intérêts composés. | |
| Segmenter les données en clusters à l'aide de scikit-learn. | |
| Envoyez des notifications aux canaux Slack via l'API. | |
| 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.