Aller au contenu principal
Page non répertoriée
Cette page n'est pas répertoriée. Les moteurs de recherche ne l'indexeront pas, et seuls les utilisateurs ayant un lien direct peuvent y accéder.

Intégrations Lakeflow

info

Bêta

Cette fonctionnalité est en version bêta.

Les intégrations sont des tâches personnalisées que vous pouvez ajouter aux Lakeflow Jobs. Un auteur écrit une intégration en Python et l'enregistre dans le workspace, la rendant disponible depuis la boîte de dialogue Ajouter une tâche . Tout autre utilisateur peut ensuite l'utiliser dans un job sans écrire de code.

Une intégration est soit une fonction, soit un capteur :

  • Les fonctions effectuent une opération unique, telle que l'envoi d'une notification.
  • Les capteurs attendent qu’une condition soit remplie en la vérifiant dans une boucle. Entre deux vérifications, un capteur libère son compute au lieu de le maintenir inactif, ce qui lui permet d’attendre efficacement.

Pour start, ajoutez une intégration ou utilisez une intégration enregistrée.

remarque

Pour envoyer vos commentaires ou poser des questions pendant l’aperçu, envoyez un e-mail à lakeflow-integrations-private-preview@databricks.com.

Ajouter une intégration​

Vous créez des intégrations à l’aide de Declarative Automation Bundles. Créez un bundle à partir du template Lakeflow Integrations , qui structure deux exemples d’intégrations que vous pouvez modifier pour créer les vôtres. Vous pouvez créer le bundle depuis le workspace ou depuis la CLI Databricks.

Pour créer et enregistrer une intégration depuis l'interface utilisateur du workspace :

  1. Créez un bundle à partir du template Lakeflow Integrations , en suivant le Tutoriel : Créer et déployer un bundle dans le Workspace.

  2. Une fois le bundle créé, cliquez sur l’icône de bundle (fusée) dans la barre latérale gauche, puis cliquez sur Deploy . Deployment upload the integration YAML and wheel files, and creates example jobs.

  3. Trouvez les wheels upload et les fichiers YAML sur /Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal.

  4. Enregistrez les intégrations importées afin qu'elles apparaissent dans l'interface utilisateur, à l'aide d'un fichier .lakeflow_integrations.yml :

    • Pour rendre une intégration disponible pour tout utilisateur, ajoutez-la à /Workspace/.lakeflow_integrations.yml et accordez aux utilisateurs du Workspace l'autorisation de lire le fichier.
    • Pour rendre une intégration visible uniquement par vous-même, ajoutez-la à /Workspace/Users/<user>/.lakeflow_integrations.yml.

    Par exemple :

    YAML
    integrations:
    - '/Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal/*.yml'
  5. Pour partager des intégrations entre les utilisateurs du workspace, ajoutez des autorisations de niveau supérieur dans databricks.yml:

    YAML
    permissions:
    - group_name: 'users'
    level: CAN_VIEW

Utiliser une intégration enregistrée​

Une fois une intégration enregistrée, tout utilisateur du workspace peut l’ajouter à un job :

  1. Refresh la page, ou quittez puis revenez sur la page Lakeflow Jobs, pour recharger la liste des intégrations disponibles.

  2. Lors de l'ajout d'une tâche à un job, cliquez sur Ajouter un autre type de tâche .

  3. Recherchez les intégrations enregistrées dans la section Integrations de la boîte de dialogue Ajouter une tâche et sélectionnez-en une.

  4. Configurez la tâche. Les champs du formulaire sont générés à partir de la configuration de l’intégration.

remarque

Les icônes personnalisées pour les intégrations ne sont pas encore prises en charge.

Référence de l’API​

Cette section décrit l’API Python pour la création d’intégrations, le type de tâche de bundle qui les exécute et le schéma YAML qui les enregistre.

API Python​

Importez ces objets depuis databricks.lakeflow.integrations pour définir des fonctions et des capteurs.

@integration​

Le décorateur @integration ajoute des métadonnées à une fonction ou à une classe Sensor. Il génère le fichier YAML d’intégration Lakeflow qui enregistre l’intégration dans le workspace ou dans la liste des intégrations utilisateur.

Python
from databricks.lakeflow.integrations import integration

Les paramètres de fonction et les paramètres de constructeur de classe deviennent des paramètres de tâche. Leurs valeurs de chaîne sont interprétées comme du JSON, donc int, str, bool, dict et list sont tous pris en charge.

Sensor​

Un protocole pour les objets qui interrogent une condition externe et qui se terminent ou sont reportés. Un Sensor est recréé entre les appels d’interrogation ; par conséquent, tout état devant survivre à un report doit être conservé en externe, par exemple dans des valeurs de tâche, des fichiers Workspace ou Lakebase.

Python
from databricks.lakeflow.integrations import Context, Sensor, SensorResult

Comme pour les fonctions, vous pouvez ajouter des paramètres de tâche via la méthode __init__.

Méthodes

  • poll(self, ctx: Context) -> SensorResult: Appelé une fois par tentative. Renvoyez SensorResult.completed() lorsque la condition est remplie, ou SensorResult.deferred(duration) pour libérer le compute et réessayer ultérieurement.

Context​

Transmis à Sensor.poll, et fournit des informations sur l'exécution de la tâche.

Python
from databricks.lakeflow.integrations import Context

Attributs

Attribut

Type

Description

main

str

Le point d'entrée principal de l'intégration ; la fonction ou la classe Sensor que la tâche exécute.

task_key

str

La clé de la tâche exécutant l’intégration.

job_run_id

int

L'ID de l'exécution du job en cours.

task_run_id

int

L’ID de l’exécution de la tâche en cours.

job_id

int

L’ID du job auquel appartient la tâche.

Attribut

Type

Description

main

str

Le point d'entrée principal de l'intégration ; la fonction ou la classe Sensor que la tâche exécute.

task_key

str

La clé de la tâche exécutant l’intégration.

job_run_id

int

L'ID de l'exécution du job en cours.

task_run_id

int

L’ID de l’exécution de la tâche en cours.

job_id

int

L’ID du job auquel appartient la tâche.

SensorResult​

Renvoyé par Sensor.poll pour indiquer si la tâche est terminée ou doit être différée.

Python
from databricks.lakeflow.integrations import SensorResult

Champ

Type

Description

status

"completed" OU "deferred"

Le résultat du sondage.

defer_for

datetime.timedelta

Durée du report avant le prochain sondage. Défini uniquement en cas de report.

Champ

Type

Description

status

"completed" OU "deferred"

Le résultat du sondage.

defer_for

datetime.timedelta

Durée du report avant le prochain sondage. Défini uniquement en cas de report.

Méthodes

  • SensorResult.completed(): La condition a été remplie et la tâche s’est terminée avec succès.
  • SensorResult.deferred(duration): La condition n'a pas été remplie. La tâche est replanifiée après duration, et le compute est libéré entre-temps.

Type de tâche de bundle​

Les intégrations Lakeflow s'exécutent en tant que type de tâche appelé python_operator_task:

  • main: la fonction principale, ou une classe qui étend Sensor.
  • parameters: un tableau de paramètres de fonction ou de constructeur de classe.

Par exemple :

YAML
resources:
jobs:
my_function:
name: 'my_function'
tasks:
- task_key: slack
environment_key: my_environment
python_operator_task:
main: my_lakeflow_integrations.my_function
parameters:
- name: 'conn_id'
value: 'CHANGEME'
environments:
- environment_key: my_environment
spec:
environment_version: '5'
dependencies:
- ../dist/*.whl

YAML d'intégration Lakeflow​

Le template de bundle génère automatiquement le YAML d’intégration Lakeflow pour toute fonction ou classe annotée avec @integration. L’interface utilisateur découvre les intégrations disponibles en inspectant /Workspace/.lakeflow_integrations.yml et /Workspace/Users/<user>/.lakeflow_integrations.yml.

Par exemple :

YAML
schema: lakeflow-integration-v0.1.0
name: Slack message
description: Post a message to a Slack channel
icon:
name: send # A known Databricks icon. See the list of available options below.
library: databricks # The only available option currently.
main: slack_operator.integrations.slack.send_message
environment:
environment_version: '5'
dependencies: # Must include the Python wheel that contains the sensor or function.
- /Workspace/Shared/integrations/slack_operator-0.1.0.whl
config:
type: object
properties:
channel: # The name of your parameter.
type: string
title: Channel # The title to render in the UI instead of the raw parameter name.
description: Channel to post to.
default: '#alerts' # Default value for new task instances.
examples: ['#my-channel'] # The first example is used as the placeholder if the field is empty.
x-ui:
widget: input # The widget to render: input, textarea, or number.
message:
type: string
title: Message
description: Message body.
examples: ['Pipeline :white_check_mark: completed']
x-ui:
widget: textarea
required: # Parameters listed here are marked as required in the UI.
- channel
- message

Pour obtenir la liste complète des valeurs d'icônes, développez les détails ci-dessous.

Valeurs d'icône disponibles

  • app
  • arrow-in
  • at
  • backup
  • bar-chart
  • beaker
  • binary
  • book
  • bookmark
  • brackets-curly
  • brackets-square
  • branch
  • briefcase
  • brush
  • bug
  • calendar
  • camera
  • catalog
  • chain
  • chart-line
  • check-circle
  • checklist
  • chip
  • clipboard
  • clock
  • cloud
  • cloud-database
  • code
  • columns
  • compass
  • connect
  • copy
  • dag
  • dashboard
  • database
  • decimal
  • dollar
  • download
  • erd
  • face-smile
  • file
  • filter
  • flag
  • flow
  • folder
  • fork
  • function
  • gear
  • gift
  • globe
  • grid
  • hash
  • history
  • home
  • image
  • ingestion
  • key
  • layer
  • leaf
  • letters
  • lightbulb
  • lightning
  • link
  • list
  • lock
  • mail
  • map
  • measure
  • megaphone
  • models
  • moon
  • notebook
  • notification
  • numbers
  • office
  • pencil
  • pie-chart
  • pipeline
  • play
  • plug
  • puzzle
  • query
  • refresh
  • robot
  • rocket
  • rows
  • school
  • search
  • send
  • share
  • shield
  • sliders
  • sparkle
  • speech-bubble
  • speedometer
  • star
  • storefront
  • stream
  • sun
  • sync
  • table
  • tag
  • target
  • terminal
  • trash
  • tree
  • trending
  • upload
  • user
  • user-group
  • visible
  • workflows
  • wrench
  • zoom-in

Facturation et compute​

Les intégrations s'exécutent sur un compute serverless, et la facturation est similaire à celle de l'exécution de notebooks. L'utilisation apparaît dans les tables système sous le nom d'exécution lakeflow_integrations:

SQL
SELECT *
FROM system.billing.usage
WHERE usage_metadata.job_name = 'lakeflow_integrations'

Le compute est réutilisé et libéré pour réduire les coûts lorsqu'une tâche remplit l'une des conditions suivantes :

  • La durée de report est supérieure à une minute.
  • Plusieurs capteurs parallèles s'exécutent pour le compte du même utilisateur ou Service Principal « exécuter en tant que », de sorte que le compute sous-jacent peut être partagé entre eux.
  • La tâche utilise la version 5 ou ultérieure du client serverless.

Lorsque plusieurs exécutions de job réutilisent le même compute, l'utilisation est attribuée à la première exécution de job ayant acquis le compute. Par exemple, l'exécution de 10 capteurs parallèles, chacun avec 20 itérations et un report de cinq minutes, produit environ 22 enregistrements totalisant environ 0,48 DBU. Le chiffre exact varie en fonction de la charge de travail.

SQL
SELECT
usage_metadata.job_name,
SUM(usage_quantity) AS usage_quantity,
COUNT(*) AS records
FROM system.billing.usage
WHERE usage_metadata.job_name = 'lakeflow_integrations'
GROUP BY ALL

Questions fréquemment posées​

Les questions suivantes traitent des problèmes courants.

Comment créer une connexion Unity Catalog pour un service HTTP externe ?​

Suivez Connect to external HTTP services.

Puis-je exécuter une charge de travail de compute lourde ?​

Les charges de travail intensives en ressources peuvent affecter d'autres charges de travail qui réutilisent le même compute ; Databricks ne recommande donc pas d'utiliser cette fonctionnalité avec des charges de travail gourmandes en compute.

Ressources supplémentaires​