Référence de la table système des événements de pipeline
Bêta
Cette table système est en version Bêta.
Cet article est une référence pour la table système pipeline_events, qui enregistre les entrées du Log des événements LakeFlow Pipelines pour les pipelines de votre compte. Chaque ligne est un événement immuable du log des événements du pipeline, capturant les transitions de cycle de vie, la progression du flux, les métriques de qualité des données, les erreurs, les ressources de cluster et d'autres données opérationnelles à travers tous les pipelines et workspaces d'une région.
Exigences
- Pour accéder à cette table système, les utilisateurs doivent soit :
- Soyez à la fois administrateur du métastore et administrateur du compte, ou
- Avoir les permissions
USEetSELECTsur les schémas système. Voir Autoriser l'accès aux tables système.
Tables d'événements de pipeline disponibles
La table système des événements de pipeline réside dans le schéma lakeflow_pipeline_events_preview pendant la version Beta, et est déplacée vers le schéma lakeflow lors de la disponibilité générale :
Table | Description | Prend en charge le streaming | Période de rétention gratuite | Comprend des données mondiales ou régionales |
|---|---|---|---|---|
pipeline_events (bêta) | Enregistre les entrées du log des événements du pipeline émises par les exécutions de pipeline | Oui | 13 mois | Régional |
Le schéma est lakeflow_pipeline_events_preview pendant la version Beta. Lors de la disponibilité générale, la table est déplacée vers le schéma lakeflow (le chemin final de la table sera system.lakeflow.pipeline_events). Les requêtes écrites pour le schéma Beta doivent être mises à jour lorsque la table est déplacée.
Référence de schéma détaillée
Schéma de la table des événements du pipeline
La table des événements de pipeline est en ajout seulement. Chaque ligne enregistre un événement unique émis par une mise à jour de pipeline au moment où il a été émis, et les lignes ne sont jamais modifiées ou supprimées sur place.
Les champs renseignés sur une ligne dépendent du type d'événement. error, update_id et de nombreux sous-champs origin.* ne sont définis que sur les événements auxquels ils s'appliquent, et la structure du champ details varie également selon event_type.
Utilisez ce tableau pour query l'activité historique du pipeline, pour créer des alertes sur les défaillances du pipeline et pour corréler le comportement du pipeline avec d'autres tables système Lakeflow.
Chemin de la table : system.lakeflow_pipeline_events_preview.pipeline_events
Clé primaire : (account_id, pipeline_event_id)
Nom de colonne | Type de données | Description | Notes |
|---|---|---|---|
| chaîne | L'ID du compte auquel appartient cet événement de pipeline | |
| chaîne | L'ID du Workspace auquel appartient cet événement de pipeline | |
| chaîne | L'ID du pipeline qui a émis l'événement | |
| chaîne | L’ID de la mise à jour du pipeline qui a émis l’événement | |
| chaîne | Identifiant unique global de l'événement | |
| chaîne | Le type d'événement (par exemple, | Consultez Event type values pour obtenir la liste complète des valeurs. |
| structure | Métadonnées contextuelles sur l'origine de l'événement, telles que le fournisseur de cloud, la région, le type de pipeline, les noms de table ou de flux et d'autres identifiants | |
| chaîne | Description de l'événement lisible par un humain | Peut être vide pour certains événements. |
| chaîne | Niveau de gravité de l'événement | L’un des éléments |
| chaîne | Stabilité du schéma d'événement | L'un des éléments suivants : |
| structure | Détails de l'erreur. Renseigné uniquement pour les événements qui contiennent des informations d'erreur | |
| variante | Charge utile spécifique à l'événement. Les champs qu'il contient dépendent de l' | Voir le champ Details. |
| Horodatage | L'heure à laquelle l'événement a été émis par le pipeline. | Fuseau horaire enregistré comme |
Champs de structure d’origine
Sous-champ | Type de données | Description |
|---|---|---|
| chaîne | Fournisseur de cloud (par exemple, |
| chaîne | Région du fournisseur de cloud |
| bigint | ID d'organisation du Workspace |
| chaîne | Le type de pipeline |
| chaîne | Le nom du pipeline fourni par l'utilisateur |
| chaîne | L'ID du cluster compute qui soutient la mise à jour du pipeline |
| chaîne | L'ID de la mise à jour de maintenance, si l'événement provient d'une exécution de maintenance |
| chaîne | Le nom du dataset (table ou vue) auquel l'événement fait référence |
| chaîne | Le nom du puits auquel l'événement fait référence |
| chaîne | Le nom du catalogue Unity Catalog |
| chaîne | Le nom du schéma Unity Catalog |
| chaîne | L'ID du flux auquel l'événement fait référence |
| chaîne | Le nom du flux auquel l'événement se réfère |
| bigint | L’ID de micro-batch pour les flux de streaming. |
| chaîne | L'ID de requête qui a initié l'action |
| chaîne | Le nom de la matérialisation |
| chaîne | L'identifiant de l'opération |
| chaîne | Le nom de la source de données |
| chaîne | L'ID de table Unity Catalog |
| chaîne | Le type de source d'ingestion (par exemple, |
| chaîne | Le nom de la connexion pour la source d’ingestion |
| chaîne | Le nom du catalogue source dans le système amont |
| chaîne | Le nom du schéma source dans le système amont |
| chaîne | Le nom de la table source dans le système en amont |
| chaîne | La version de la table source, le cas échéant |
Champs de structure d’erreur
Sous-champ | Type de données | Description |
|---|---|---|
| booléen | Si l'erreur a provoqué l'arrêt de la mise à jour. |
| array<struct> | Chaîne d'exceptions associée à l'erreur (cause principale en dernier) |
| chaîne | Code SQLSTATE, si disponible |
| chaîne | Classe d'erreur Databricks, si disponible |
Champ de détails
La colonne details est un VARIANT, et les champs qu’elle contient dépendent du event_type. Pour connaître les champs disponibles sous chaque type d’événement, consultez le schéma du journal des événements du pipeline. Utilisez la fonction variant_get ou la syntaxe de points pour lire les valeurs imbriquées. Voir les exemples de requêtes ci-dessous pour les modèles d'accès typiques.
La clé event_type enveloppe la charge utile. Par exemple, les métriques d'un événement flow_progress se trouvent à $.flow_progress.metrics, et non à $.metrics. Incluez la clé event-type dans chaque chemin.
-- Using variant_get (lets you cast to a specific type)
SELECT
pipeline_id,
event_time,
variant_get(details, '$.flow_progress.status', 'STRING') AS flow_status,
variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT') AS rows_written,
variant_get(details, '$.flow_progress.metrics.backlog_bytes', 'BIGINT') AS backlog_bytes
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
-- Using dot syntax (returns VARIANT, cast when needed)
SELECT
pipeline_id,
event_time,
details:flow_progress.status::STRING AS flow_status,
details:flow_progress.metrics.num_output_rows::BIGINT AS rows_written
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
Exemples de requêtes
-- Flow throughput for a specific pipeline
SELECT
origin.flow_name,
date_trunc('HOUR', event_time) AS hour,
SUM(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT')) AS rows_written
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 7 DAYS
GROUP BY
origin.flow_name,
date_trunc('HOUR', event_time)
ORDER BY
hour DESC,
rows_written DESC
-- The latest error for each pipeline that has errored in the last 7 days, with the outermost exception.
-- The exception chain is ordered with the root cause last, so read element -1 for the root cause.
-- On many errors only the first element carries error_class and sql_state.
SELECT
workspace_id,
pipeline_id,
event_time,
event_type,
message,
error.exceptions[0].error_class AS exception_error_class,
error.exceptions[0].sql_state AS exception_sql_state
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
level = 'ERROR'
AND event_time >= current_timestamp() - INTERVAL 7 DAYS
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id ORDER BY event_time DESC) = 1
ORDER BY
event_time DESC
-- Data quality: failed expectations by dataset, per update, in the last 1 day
SELECT
pipeline_id,
update_id,
origin.dataset_name,
expectation.name AS expectation_name,
SUM(expectation.failed_records) AS failed_records
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
LATERAL VIEW explode(variant_get(details, '$.flow_progress.data_quality.expectations', 'ARRAY<STRUCT<name:STRING,dataset:STRING,passed_records:BIGINT,failed_records:BIGINT>>')) AS expectation
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 DAY
GROUP BY
pipeline_id,
update_id,
origin.dataset_name,
expectation.name
HAVING
SUM(expectation.failed_records) > 0
ORDER BY
failed_records DESC
Modèles de jointure courants
Joignez-vous à la table pipelines pour filtrer par nom de pipeline
La table pipelines est une dimension à évolution lente (SCD2). Prenez la dernière version de chaque pipeline avant de les joindre.
WITH latest_pipelines AS (
SELECT *
FROM system.lakeflow.pipelines
QUALIFY ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id ORDER BY change_time DESC) = 1
)
SELECT
p.name AS pipeline_name,
e.event_time,
e.event_type,
e.level,
e.message
FROM
system.lakeflow_pipeline_events_preview.pipeline_events e
JOIN
latest_pipelines p
ON e.workspace_id = p.workspace_id
AND e.pipeline_id = p.pipeline_id
WHERE
e.level = 'ERROR'
AND e.event_time >= current_timestamp() - INTERVAL 24 HOURS
ORDER BY
e.event_time DESC
Joindre avec pipeline_update_timeline sur update_id
SELECT
u.period_start_time AS update_start,
u.period_end_time AS update_end,
e.event_time,
e.event_type,
e.level,
e.message
FROM
system.lakeflow.pipeline_update_timeline u
JOIN
system.lakeflow_pipeline_events_preview.pipeline_events e
ON e.update_id = u.update_id
WHERE
u.pipeline_id = '<your-pipeline-id>'
AND u.period_start_time >= current_timestamp() - INTERVAL 7 DAYS
ORDER BY
u.period_start_time DESC,
e.event_time ASC
Configuration des alertes
Vous pouvez créer des alertes sur pipeline_events en utilisant les alertes Databricks SQL. Écrivez une query SQL sur pipeline_events (éventuellement jointe à d'autres tables système LakeFlow), planifiez-la sur un SQL Warehouse, et configurez une destination de notification (e-mail, Slack, webhook, PagerDuty).
Quelques points de départ utiles :
Alerter lorsqu'aucun événement n'est arrivé pour un pipeline au cours des N dernières minutes
Utilisez ceci pour détecter les pipelines bloqués ou en échec silencieux.
-- Returns one row per pipeline that has not emitted any event in the last 30 minutes.
-- The alert can trigger when this query returns any rows.
SELECT
workspace_id,
pipeline_id,
MAX(event_time) AS last_event_time
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_time >= current_timestamp() - INTERVAL 24 HOURS
GROUP BY
workspace_id,
pipeline_id
HAVING
MAX(event_time) < current_timestamp() - INTERVAL 30 MINUTES
Alerte lorsque le backlog pour un flux spécifique est trop élevé
L’arriéré est rapporté sur flow_progress événements comme backlog_bytes, et pour les sources de fichiers également comme backlog_files. Trigger lorsque la lecture la plus récente dépasse un threshold (par exemple, 100 Mo de travail non traité). Toutes les sources ne rapportent pas chaque métrique, donc filtrez selon celle que votre source remplit.
-- Returns the most recent backlog reading per flow for a given pipeline.
-- The alert can trigger when backlog_bytes exceeds the threshold for any flow.
WITH latest_flow_progress AS (
SELECT
workspace_id,
pipeline_id,
origin.flow_name,
event_time,
CASE
WHEN variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED' THEN 0
ELSE variant_get(details, '$.flow_progress.metrics.backlog_bytes', 'BIGINT')
END AS backlog_bytes
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
AND (
variant_get(details, '$.flow_progress.metrics.backlog_bytes', 'BIGINT') IS NOT NULL
OR variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED'
)
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id, origin.flow_name ORDER BY event_time DESC) = 1
)
SELECT *
FROM latest_flow_progress
WHERE backlog_bytes > 100000000
Pour voir quelle source est en retard, $.flow_progress.metrics.source_metrics est un tableau de lectures par source, chacune avec source_name aux côtés de backlog_bytes, backlog_records ou backlog_files de cette source.
Alerte sur les baisses de qualité des données dans un pipeline
Chaque événement flow_progress indique le nombre de lignes supprimées par les attentes EXPECT … DROP. Additionnez-les par dataset sur une fenêtre de mise à jour et déclenchez une alerte lorsque le total dépasse un threshold.
-- Returns datasets where more than 100 rows were dropped by expectations, per update, in the last hour.
SELECT
pipeline_id,
update_id,
origin.dataset_name,
SUM(variant_get(details, '$.flow_progress.data_quality.dropped_records', 'BIGINT')) AS dropped_records
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
GROUP BY
pipeline_id,
update_id,
origin.dataset_name
HAVING
SUM(variant_get(details, '$.flow_progress.data_quality.dropped_records', 'BIGINT')) > 100
Alerte lorsqu'un flux traite trop peu de lignes
flow_progress les événements rapportent metrics.num_output_rows sous forme de comptage par micro-batch, donc sommer les événements dans une fenêtre donne le nombre de lignes écrites sur cette fenêtre. Créez une alerte pour le moment où le throughput tombe en dessous d'un seuil attendu. Par exemple, un flux qui écrit normalement des milliers de lignes par heure mais qui en produit près de zéro peut indiquer une source mal configurée.
Cette query ne rapporte que les flux ayant émis un événement flow_progress avec un nombre de lignes dans la fenêtre. Un flux totalement bloqué n’émet aucun événement ; associez donc cette alerte à l’alerte d’événements manquants ci-dessus.
-- Returns flows that wrote fewer than 100 rows in the last hour.
-- The alert can trigger when this query returns any rows.
SELECT
pipeline_id,
origin.flow_name,
SUM(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT')) AS rows_written_last_hour
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
AND variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT') IS NOT NULL
GROUP BY
pipeline_id,
origin.flow_name
HAVING
SUM(variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT')) < 100
Alerter lorsque la latence des nouvelles données est trop élevée
Pour les flux de streaming, les événements flow_progress indiquent une latence en streaming_metrics. stream_latency_ms est le temps écoulé entre l'arrivée des données en amont et la validation du micro-batch dans la table Delta. Vous pouvez définir un Trigger pour le moment où la lecture la plus récente dépasse un threshold (par exemple, 5 minutes).
Seuls les flux en streaming avec une heure d’événement marquée signalent stream_latency_ms, et uniquement lorsque les métriques temporelles SDP sont activées. Les autres flux renvoient NULL à chaque événement, et cette alerte ne se déclenche jamais pour eux.
-- Returns the most recent new-data latency per flow for a given pipeline.
-- The alert can trigger when stream_latency_ms > 300000 (5 minutes) for any flow.
WITH latest_latency AS (
SELECT
workspace_id,
pipeline_id,
origin.flow_name,
event_time,
CASE
WHEN variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED' THEN 0
ELSE variant_get(details, '$.flow_progress.streaming_metrics.stream_latency_ms', 'BIGINT')
END AS stream_latency_ms
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
WHERE
pipeline_id = '<your-pipeline-id>'
AND event_type = 'flow_progress'
AND event_time >= current_timestamp() - INTERVAL 1 HOUR
AND (
variant_get(details, '$.flow_progress.streaming_metrics.stream_latency_ms', 'BIGINT') IS NOT NULL
OR variant_get(details, '$.flow_progress.status', 'STRING') = 'COMPLETED'
)
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id, origin.flow_name ORDER BY event_time DESC) = 1
)
SELECT *
FROM latest_latency
WHERE stream_latency_ms > 300000
Conseils pour les alertes de production
- Filtrez par
pipeline_id(etworkspace_idsi vous gérez des alertes par Workspace) afin que chaque alerte cible une portée spécifique plutôt que l'ensemble du compte. - Choisissez une cadence d'évaluation qui correspond à la sensibilité de l'alerte. Utilisez un intervalle court (par exemple, toutes les 5 minutes) pour les signaux de défaillance rapide et un intervalle plus long (par exemple, toutes les heures) pour le backlog et les tendances de qualité des données. Une condition de type « Trigger lorsque la query renvoie plus de 0 lignes » fonctionne dans la plupart des cas.
Valeurs de référence
Valeurs de niveau
Valeur | Description |
|---|---|
| Activité normale du pipeline (progression du flux, mises à jour des transitions du cycle de vie, modifications de configuration). |
| Problèmes non fatals dont le pipeline s'est rétabli, ou qui peuvent nécessiter une attention. |
| Échecs qui ont empêché le pipeline de progresser sur un flux ou une mise à jour. |
| Mesures quantitatives émises pendant l'exécution (nombre de lignes, throughput, latence). |
Valeurs de niveau de maturité
Valeur | Description |
|---|---|
| Le schéma d’événement est stable. Aucune modification majeure n'est attendue. Fiable pour créer des requêtes et des alertes de production. |
| Le schéma d'événement peut changer dans les prochaines versions. Utiliser avec prudence. |
| Le type d'événement ou le schéma est obsolète et sera supprimé dans les futures versions. Migrez-en. |
Valeurs de type d'événement
Le champ event_type est une énumération. L'ensemble complet des valeurs :
Valeur | Description |
|---|---|
| Une nouvelle mise à jour de pipeline a été demandée. |
| Une mise à jour de pipeline a transité par un état de cycle de vie. |
| Un flux (dataset) au sein d'une mise à jour est passé par un état. |
| Métadonnées statiques sur un flux. |
| Métadonnées statiques concernant un dataset. |
| Métadonnées statiques sur un récepteur de sortie. |
| Une fonctionnalité obsolète a été utilisée par le pipeline. |
| Décision de dimensionnement automatique des clusters. |
| Une Opération non prise en charge dans la configuration actuelle. |
| Indicateurs d'emplacement de tâche et de mise à l'échelle automatique pour le compute sous-jacent. |
| Informations de la phase de planification pour la mise à jour. |
| Pression de la collecte de mémoire sur le Driver ou les exécuteurs. |
| La mise à jour s'est terminée anormalement. |
| Pression sur l’espace disque du cluster. |
| Progression du cycle de vie d'un hook de pipeline. |
| Un événement de cycle de vie de dataset. |
| Une Opérations en arrière-plan est passée par un état. |
| Le pipeline a effectué un appel API sortant. |
| Progression pour une opération générique. |
| Progression d'une query de streaming qui prend en charge un flux. |
| Résumé d'une opération de retour de pipeline. |
| Un message consultatif du moteur. |
| Configuration détaillée de l'environnement d'exécution. |
| Informations sur les Ressources (cluster, type d'instance, etc.). |
| Statut de la configuration des notifications de fichier (pour les sources de fichiers cloud). |
| Notification de changement de comportement sous Spark Connect. |
| Une action initiée par l'utilisateur sur le pipeline. |
| Contexte concernant le code utilisateur associé à l'événement. |