Informations de référence sur la table système des événements de pipeline

Important

Cette table système est en version bêta.

Cet article est une référence pour la pipeline_events table système, qui enregistre les entrées du journal des événements des pipelines Lakeflow pour les pipelines dans votre compte. Chaque ligne est un événement immuable du journal des événements de pipeline, la capture des 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 dans tous les pipelines et espaces de travail au sein d’une région.

Requirements

  • Pour accéder à cette table système, les utilisateurs doivent :
    • Être à la fois un administrateur de metastore et un administrateur de compte, ou
    • Disposez d’autorisations USE et SELECT sur les schémas système. Consultez Octroyer un accès aux tables système.

Tables d’événements de pipeline disponibles

La table du système d’événements du pipeline se trouve dans le lakeflow_pipeline_events_preview schéma pendant la bêta, et se déplace vers le lakeflow schéma en disponibilité générale :

Table Descriptif Prend en charge la diffusion en continu Période de rétention gratuite Inclut des données globales ou régionales
pipeline_events (Bêta) Enregistre les entrées du journal des événements de pipeline émises par les exécutions de pipeline Yes 13 mois Régional

Remarque

Le schéma est lakeflow_pipeline_events_preview en cours de Bêta. En disponibilité générale, la table se déplace vers le lakeflow schéma (le chemin final de la table sera system.lakeflow.pipeline_events). Les requêtes écrites sur le schéma Beta doivent être mises à jour lorsque la table bouge.

Informations de référence détaillées sur le schéma

Schéma de table d’événements de pipeline

La table des événements de pipeline est append-only. Chaque ligne enregistre un événement unique émis par une mise à jour de pipeline au moment où elle a été émise, et les lignes ne sont jamais modifiées ou supprimées en place.

Les champs renseignés sur une ligne dépendent du type d’événement. error, update_idet de nombreux origin.* sous-champs sont définis uniquement sur les événements où ils s’appliquent, et la structure du details champ varie également par event_type.

Utilisez cette table pour interroger l’activité du pipeline historique, générer des alertes sur les échecs de pipeline et mettre en corrélation le comportement du pipeline avec d’autres tables système Lakeflow.

Chemin d’accès de la table : system.lakeflow_pipeline_events_preview.pipeline_events

Clé primaire : (account_id, pipeline_event_id)

Nom de la colonne Type de données Descriptif Notes
account_id string ID du compte auquel appartient cet événement de pipeline
workspace_id string L’ID de l’espace de travail auquel appartient cet événement de pipeline
pipeline_id string ID du pipeline qui a émis l’événement
update_id string ID de la mise à jour du pipeline qui a émis l’événement
pipeline_event_id string Identificateur global unique pour l’événement
event_type string Type d’événement (par exemple, , flow_progressupdate_progress, create_update) Voir Valeurs de type d’événement pour l’ensemble complet des valeurs.
origin struct Métadonnées contextuelles sur l’origine de l’événement, telles que le fournisseur de cloud, la région, le type de pipeline, la table ou les noms de flux et d’autres identificateurs Voir champs de structures d’origine.
message string Description lisible par l’homme de l’événement Peut être vide pour certains événements.
level string Niveau de gravité de l’événement Un des INFO, WARN, ERROR, METRICS. Voir les valeurs de niveau.
maturity_level string Stabilité du schéma d’événement L’un des STABLE, EVOLVING, DEPRECATED. Voir Valeurs de niveau de maturité.
error struct Détails de l’erreur. Renseigné uniquement pour les événements qui contiennent des informations d’erreur Voir Champs de structure d’erreur.
details variante Charge utile spécifique à l’événement. Les champs qu’il contient dépendent de la event_type Voir le champ Détails.
event_time timestamp Heure à laquelle l’événement a été émis par le pipeline Fuseau horaire enregistré en tant que +00:00 (UTC).

Champs de structures d’origine

Sous-champ Type de données Descriptif
cloud string Fournisseur de cloud (par exemple, , AWSAZURE, GCP)
region string Région du fournisseur de cloud
org_id bigint ID d’organisation de l’espace de travail
pipeline_type string Type de pipeline
pipeline_name string Nom fourni par l’utilisateur du pipeline
cluster_id string ID de cluster de calcul qui sauvegarde la mise à jour du pipeline
maintenance_id string ID de la mise à jour de maintenance, si l’événement provient d’une exécution de maintenance
dataset_name string Le nom du jeu de données (table ou vue) auquel l’événement fait référence
sink_name string Le nom du récepteur auquel l’événement fait référence
catalog_name string Nom du catalogue Unity
schema_name string Nom du schéma du catalogue Unity
flow_id string L’ID du flux auquel l’événement fait référence
flow_name string Le nom du flux auquel l’événement fait référence
batch_id bigint ID de micro-lot pour les flux de streaming. bigint pour la compatibilité
request_id string ID de demande qui a lancé l’action
materialization_name string Nom de matérialisation
operation_id string ID d’opération
source_name string Nom de la source de données
uc_table_id string ID de la table de catalogue Unity
ingestion_source_type string Type de source d’ingestion (par exemple, SQL_SERVER, SALESFORCE)
ingestion_source_connection_name string Nom de connexion pour la source d’ingestion
ingestion_source_catalog_name string Nom du catalogue source dans le système en amont
ingestion_source_schema_name string Nom du schéma source dans le système en amont
ingestion_source_table_name string Nom de la table source dans le système en amont
ingestion_source_table_version string Version de la table source, le cas échéant

Champs de struct d’erreur

Sous-champ Type de données Descriptif
fatal boolean Indique si l’erreur a provoqué l’arrêt de la mise à jour
exceptions struct de tableau<> Chaîne d’exceptions associée à l’erreur (cause racine en dernier)
exceptions[].sql_state string Code SQLSTATE, le cas échéant
exceptions[].error_class string Classe d’erreur Databricks, si disponible

Champ de détails

La details colonne est un VARIANT, et les champs qu’il 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 pipeline. Utilisez la syntaxe de variant_get fonction ou de point pour lire les valeurs imbriquées. Consultez les exemples de requêtes ci-dessous pour connaître les modèles d’accès classiques.

La event_type clé enveloppe la charge utile. Par exemple, les métriques d’un flow_progress événement sont à $.flow_progress.metrics, et non $.metrics. Incluez la clé de type événement 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

Joindre la table pour filtrer par nom de pipelines pipeline

La pipelines table est une dimension à variation lente (SCD2). Prenez la dernière version de chaque pipeline avant de 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

Rejoindre avec pipeline_update_timeline on 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

Mise en place des alertes

Vous pouvez générer des alertes à pipeline_events l’aide d’alertes Databricks SQL. Écrivez une requête pipeline_events SQL sur (éventuellement jointe à d’autres tables système Lakeflow), planifiez-la sur un entrepôt SQL et configurez une destination de notification (e-mail, Slack, webhook, PagerDuty).

Quelques points de départ utiles :

Alerte quand aucun événement n’est arrivé pour un pipeline au cours des dernières minutes

Utilisez cette option pour détecter les pipelines bloqués ou défaillants en mode 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 d’un flux spécifique est trop élevé

L’arriéré est rapporté sur flow_progress les événements comme backlog_bytes, et pour les sources de fichiers également comme backlog_files. Déclenche lorsque la lecture la plus récente dépasse un seuil (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 il y a un ensemble de lectures par source, chacune avec source_name à côté de la backlog_bytessource , backlog_records ou backlog_files.

Alerte sur les baisses de qualité des données dans un pipeline

Chaque flow_progress événement signale le nombre de lignes supprimées par EXPECT … DROP les attentes. Additionnez ces données par ensemble de données sur une fenêtre de mise à jour et alertez-les lorsque le total dépasse un seuil.

-- 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

Alerter lorsqu’un flux traite trop peu de lignes

flow_progress les événements rapportent metrics.num_output_rows par micro-batch, donc en sommant les événements dans une fenêtre, on obtient les lignes écrites au-dessus de cette fenêtre. Créez une alerte lorsque le débit descend en dessous d’un plancher attendu. Par exemple, un flux qui écrit normalement des milliers de lignes par heure mais produit près de zéro peut indiquer une source mal configurée.

Cette requête ne rapporte que les flux qui ont émis un flow_progress événement avec un nombre de lignes dans la fenêtre. Un flux complètement arrêté n’émet aucun événement, donc associez 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 flow_progress événements rapportent la latence dans streaming_metrics. stream_latency_ms est le temps entre l’arrivée des données en amont et le moment où le micro-batch s’est engagé sur la table Delta. Vous pouvez définir un déclencheur lorsque la lecture la plus récente dépasse un seuil (par exemple, 5 minutes).

Seul le streaming s’écoule avec un rapport stream_latency_msd’heure d’événement balisé, et uniquement lorsque les métriques SDP sont activées. D’autres flux reviennent 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 (et workspace_id si vous conservez des alertes par espace de travail) afin que chaque alerte cible une étendue 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 à échec rapide et un intervalle plus long (par exemple, toutes les heures) pour les tendances de backlog et de qualité des données. Un « déclencheur lorsque la requête retourne plus de 0 lignes » fonctionne dans la plupart des cas.

Valeurs de référence

Valeurs de niveau

Valeur Descriptif
INFO Activité de pipeline normale (progression du flux, transitions de cycle de vie de mise à jour, modifications de configuration).
WARN Problèmes non irrécupérables sur lesquels le pipeline a été récupéré ou qui peut nécessiter une attention particulière.
ERROR Échecs qui empêchaient le pipeline d’avancer sur un flux ou une mise à jour.
METRICS Mesures quantitatives émises pendant l’exécution (nombres de lignes, débit, latence).

Valeurs de niveau d’échéance

Valeur Descriptif
STABLE Le schéma d’événement est stable. Les changements cassants ne sont pas attendus. Coffre-fort pour générer des requêtes et des alertes de production.
EVOLVING Le schéma d’événement peut changer dans les versions ultérieures. Utilisez avec soin.
DEPRECATED Le type d’événement ou le schéma est déconseillé et sera supprimé dans les versions ultérieures. Migrez-la.

Valeurs de type d’événement

Le event_type champ est une énumération. Ensemble complet de valeurs :

Valeur Descriptif
create_update Une nouvelle mise à jour de pipeline a été demandée.
update_progress Une mise à jour de pipeline est passée à un état de cycle de vie.
flow_progress Flux (jeu de données) dans une mise à jour passée par un état.
flow_definition Métadonnées statiques sur un flux.
dataset_definition Métadonnées statiques sur un jeu de données.
sink_definition Métadonnées statiques sur un récepteur de sortie.
deprecation Une fonctionnalité déconseillée a été utilisée par le pipeline.
autoscale Décision de mise à l’échelle automatique du cluster.
unsupported_operation Opération qui n’est pas prise en charge dans la configuration actuelle.
cluster_resources Emplacement de tâche et métriques de mise à l’échelle automatique pour le calcul de stockage.
planning_information Informations de phase de planification pour la mise à jour.
gc_pressure Pression sur le garbage collection sur les pilotes ou les exécuteurs.
abnormal_termination La mise à jour s’est arrêtée anormalement.
disk_space Pression de l’espace disque sur le cluster.
hook_progress Progression du cycle de vie d’un hook de pipeline.
dataset_life_cycle Événement de cycle de vie du jeu de données.
background_operation Une opération en arrière-plan est passée à un état.
remote_api_usage Le pipeline a effectué un appel d’API sortant.
operation_progress Progression d’une opération générique.
stream_progress Progression d’une requête de streaming qui sauvegarde un flux.
rewind_summary Résumé d’une opération de rembobinage de pipeline.
advisory Message consultatif du moteur.
runtime_details Configuration détaillée du runtime.
resource_info Informations sur les ressources (cluster, type d’instance, etc.).
file_notification_set_up État de configuration de la notification de fichier (pour les sources cloud-files).
behavior_change_in_spark_connect Notification de modification de comportement sous Spark Connect.
user_action Action initiée par l’utilisateur sur le pipeline.
user_code_context Contexte sur le code utilisateur associé à l’événement.