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.
Tableau de bord des événements du pipeline
Un tableau de bord préconçu visualise l’état de santé du pipeline, le throughput, l’arriéré (backlog), la qualité des données et les erreurs de cette table système.


Pour importer et utiliser le tableau de bord, consultez la page Surveiller et query les événements.
Tables d’événements de pipeline disponibles
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 bêta. 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 détails du schéma peuvent également changer à ce moment-là. Les requêtes écrites pour le schéma bêta 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
WITH updates AS (
SELECT
workspace_id,
pipeline_id,
update_id,
MIN(period_start_time) AS update_start,
MAX(CASE WHEN result_state IS NOT NULL THEN period_end_time END) AS update_end
FROM
system.lakeflow.pipeline_update_timeline
WHERE
pipeline_id = '<your-pipeline-id>'
GROUP BY
workspace_id,
pipeline_id,
update_id
)
SELECT
u.update_start,
u.update_end,
e.event_time,
e.event_type,
e.level,
e.message
FROM
updates u
JOIN
system.lakeflow_pipeline_events_preview.pipeline_events e
ON e.workspace_id = u.workspace_id
AND e.update_id = u.update_id
ORDER BY
u.update_start 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 :
Alerte lorsqu’aucun événement n’est arrivé pour un pipeline en cours d’exécution 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 with a currently-running update that has emitted no event
-- in the last 30 minutes. The alert can trigger when this query returns any rows.
WITH running_pipelines AS (
SELECT
workspace_id,
pipeline_id
FROM
system.lakeflow.pipeline_update_timeline
QUALIFY
ROW_NUMBER() OVER (PARTITION BY workspace_id, pipeline_id ORDER BY period_end_time DESC) = 1
AND result_state IS NULL
),
last_event AS (
SELECT
workspace_id,
pipeline_id,
MAX(event_time) AS last_event_time
FROM
system.lakeflow_pipeline_events_preview.pipeline_events
GROUP BY
workspace_id,
pipeline_id
)
SELECT
r.workspace_id,
r.pipeline_id,
le.last_event_time
FROM
running_pipelines r
LEFT JOIN last_event le USING (workspace_id, pipeline_id)
WHERE
le.last_event_time IS NULL
OR le.last_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 active flows (latest status not terminal) that wrote fewer than 100 rows in the last hour.
-- The alert can trigger when this query returns any rows.
WITH flow_state AS (
SELECT
pipeline_id,
origin.flow_name,
event_time,
variant_get(details, '$.flow_progress.status', 'STRING') AS status,
variant_get(details, '$.flow_progress.metrics.num_output_rows', 'BIGINT') AS num_output_rows,
ROW_NUMBER() OVER (PARTITION BY pipeline_id, origin.flow_name ORDER BY event_time DESC) AS rn
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
),
active_flows AS (
SELECT
pipeline_id,
flow_name
FROM
flow_state
WHERE
rn = 1
AND status NOT IN ('COMPLETED', 'FAILED', 'SKIPPED', 'STOPPED', 'EXCLUDED')
)
SELECT
f.pipeline_id,
f.flow_name,
SUM(f.num_output_rows) AS rows_written_last_hour
FROM
flow_state f
JOIN active_flows a USING (pipeline_id, flow_name)
WHERE
f.num_output_rows IS NOT NULL
GROUP BY
f.pipeline_id,
f.flow_name
HAVING
SUM(f.num_output_rows) < 100
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 restauration 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. |