Aller au contenu principal

Définir le monitoring personnalisé des pipelines avec des hooks d'événement

info

Aperçu

La prise en charge des hooks d'événements est en préversion publique.

Vous pouvez utiliser des *event hooks* pour ajouter des fonctions de rappel Python personnalisées qui s'exécutent lorsque des événements sont conservés dans le journal des événements d'un pipeline. Vous pouvez utiliser des event hooks pour implémenter un monitoring personnalisé et des solutions d'alerte. Par exemple, vous pouvez utiliser des event hooks pour envoyer des e-mails ou écrire dans un log lorsque des événements spécifiques se produisent, ou pour s'intégrer à des solutions tierces afin de surveiller les événements de pipeline.

Définissez un hook d'événement avec une fonction Python qui accepte un seul argument, où l'argument est un dictionnaire représentant un événement. Ensuite, incluez les hooks d'événement dans le cadre du code source d'un pipeline. Tous les hooks d'événement définis dans un pipeline tentent de traiter tous les événements générés lors de chaque mise à jour du pipeline. Si votre pipeline est composé de plusieurs fichiers de code source, tous les hooks d'événement définis s'appliquent à l'intégralité du pipeline. Bien que les hooks d'événement soient inclus dans le code source de votre pipeline, ils ne sont pas inclus dans le Graphe du pipeline.

Vous pouvez utiliser des hooks d'événements avec des pipelines qui publient vers le Hive metastore ou Unity Catalog.

remarque
  • Python est le seul langage pris en charge pour la définition de hooks d'événements. Pour définir des fonctions Python personnalisées qui traitent les événements dans un pipeline implémenté à l'aide de l'interface SQL, ajoutez les fonctions personnalisées dans un fichier source Python distinct qui s'exécute dans le cadre du pipeline. Les fonctions Python sont appliquées à l'ensemble du pipeline lorsque le pipeline s'exécute.
  • Les Event hooks sont déclenchés uniquement pour les événements où le maturity_level est STABLE.
  • Les event hooks sont exécutés de manière asynchrone par rapport aux mises à jour du pipeline, mais de manière synchrone avec les autres event hooks. Cela signifie qu'un seul hook d'événement s'exécute à la fois, et que les autres hooks d'événement attendent de s'exécuter jusqu'à ce que le hook d'événement en cours d'exécution se termine. Si un event hook s'exécute indéfiniment, il bloque tous les autres event hooks.
  • Lakeflow pipelines tente d'exécuter chaque hook d'événement sur chaque événement émis lors d'une mise à jour du pipeline. Pour s'assurer que les hooks d'événements en retard ont le temps de traiter tous les événements en file d'attente, le pipeline attend une période fixe non configurable avant de terminer son compute. Cependant, il n'est pas garanti que tous les hooks soient déclenchés sur tous les événements avant la fin du compute.

Surveiller le traitement des event hooks

Utilisez le type d'événement hook_progress dans le log des événements du pipeline pour surveiller l'état des crochets d'événements d'une mise à jour. Pour éviter les dépendances circulaires, les hooks d'événements ne sont pas Trigger pour les événements hook_progress.

Définir un hook d'événement

Pour définir un hook d'événement, utilisez le décorateur on_event_hook :

Python
@dp.on_event_hook(max_allowable_consecutive_failures=None)
def user_event_hook(event):
# Python code defining the event hook

Le max_allowable_consecutive_failures décrit le nombre maximal de fois consécutives qu’un event hook peut échouer avant qu’il ne soit désactivé. Un échec d’event hook est défini comme tout moment où l’event hook lève une exception. Si un hook d'événement est désactivé, il ne traite pas les nouveaux événements tant que le pipeline n'est pas redémarré.

max_allowable_consecutive_failures doit être un entier supérieur ou égal à 0 ou None. Une valeur de None (attribuée par default) signifie qu'il n'y a aucune limite au nombre d'échecs consécutifs autorisés pour le hook d'événement, et le hook d'événement n'est jamais désactivé.

Les échecs de hook d'événement et la désactivation des hooks d'événement peuvent être surveillés dans le log d'événements en tant qu'événements hook_progress.

La fonction de hook d'événement doit être une fonction Python qui accepte exactement un paramètre, une représentation sous forme de dictionnaire de l'événement qui a Trigger ce hook d'événement. Toute valeur de retour de la fonction de hook d'événement est ignorée.

Exemple : Sélectionner des événements spécifiques pour le traitement

L'exemple suivant présente un point d'entrée d'événement qui sélectionne des événements spécifiques pour le traitement. Plus précisément, cet exemple attend que les événements du pipeline STOPPING soient reçus, puis affiche un message dans les Logs du Driver stdout.

Python
@dp.on_event_hook
def my_event_hook(event):
if (
event['event_type'] == 'update_progress' and
event['details']['update_progress']['state'] == 'STOPPING'
):
print('Received notification that update is stopping: ', event)

Exemple : Envoyez tous les événements à une chaîne Slack

L'exemple suivant implémente un hook d'événement qui envoie tous les événements reçus à un canal Slack à l'aide de l'API Slack.

Cet exemple utilise un secret Databricks pour stocker de manière sécurisée un jeton requis pour s'authentifier auprès de l'API Slack.

Python
from pyspark import pipelines as dp
import requests

# Get a Slack API token from a Databricks secret scope.
API_TOKEN = dbutils.secrets.get(scope="<secret-scope>", key="<token-key>")

@dp.on_event_hook
def write_events_to_slack(event):
res = requests.post(
url='https://slack.com/api/chat.postMessage',
headers={
'Content-Type': 'application/json',
'Authorization': 'Bearer ' + API_TOKEN,
},
json={
'channel': '&lt;channel-id&gt;',
'text': 'Received event:\n' + event,
}
)

Exemple : configurez un hook d'événement à désactiver après quatre échecs consécutifs

L'exemple suivant montre comment configurer un hook d'événement qui est désactivé s'il échoue quatre fois consécutivement.

Python
from pyspark import pipelines as dp
import random

def run_failing_operation():
raise Exception('Operation has failed')

# Allow up to 3 consecutive failures. After a 4th consecutive
# failure, this hook is disabled.
@dp.on_event_hook(max_allowable_consecutive_failures=3)
def non_reliable_event_hook(event):
run_failing_operation()

Exemple : Pipeline avec un hook d'événement

L'exemple suivant illustre l'ajout d'un hook d'événement au code source d'un pipeline. Il s'agit d'un exemple simple mais complet d'utilisation de webhooks d'événements avec un pipeline.

Python
from pyspark import pipelines as dp
import requests
import json
import time

API_TOKEN = dbutils.secrets.get(scope="<secret-scope>", key="<token-key>")
SLACK_POST_MESSAGE_URL = 'https://slack.com/api/chat.postMessage'
DEV_CHANNEL = 'CHANNEL'
SLACK_HTTPS_HEADER_COMMON = {
'Content-Type': 'application/json',
'Authorization': 'Bearer ' + API_TOKEN
}

# Create a single dataset.
@dp.table
def test_dataset():
return spark.range(5)

# Definition of event hook to send events to a Slack channel.
@dp.on_event_hook
def write_events_to_slack(event):
res = requests.post(url=SLACK_POST_MESSAGE_URL, headers=SLACK_HTTPS_HEADER_COMMON, json = {
'channel': DEV_CHANNEL,
'text': 'Event hook triggered by event: ' + event['event_type'] + ' event.'
})