Referencia de la tabla del sistema de eventos de canalización

Importante

Esta tabla del sistema está en Beta.

Este artículo es una referencia para la tabla del pipeline_events sistema, que registra entradas del registro de eventos de canalizaciones de Lakeflow para las canalizaciones de la cuenta. Cada fila es un evento inmutable del registro de eventos de canalización, la captura de transiciones del ciclo de vida, el progreso del flujo, las métricas de calidad de datos, los errores, los recursos del clúster y otros datos operativos en todas las canalizaciones y áreas de trabajo de una región.

Requirements

  • Para acceder a esta tabla del sistema, los usuarios deben:

Tablas de eventos de canalización disponibles

La tabla del sistema de eventos de la tubería permanece en el lakeflow_pipeline_events_preview esquema durante la Beta y se traslada al lakeflow esquema en disponibilidad general:

Table Descripción Admite transmisión en directo Período gratuito de retención Incluye datos globales o regionales
pipeline_events (Beta) Registra entradas del registro de eventos de canalización emitidas por ejecuciones de canalización 13 meses Regional

Note

El esquema está lakeflow_pipeline_events_preview durante la Beta. En disponibilidad general, la tabla se traslada al lakeflow esquema (el camino final de la tabla será system.lakeflow.pipeline_events). Las consultas escritas contra el esquema Beta deben actualizarse cuando la tabla se mueve.

Referencia detallada del esquema

Esquema de tabla de eventos de canalización

La tabla de eventos de canalización es de solo anexión. Cada fila registra un único evento emitido por una actualización de canalización en el momento en que se emitió y las filas nunca se modifican ni eliminan en su lugar.

Los campos que se rellenan en una fila dependen del tipo de evento. error, update_idy muchos origin.* subcampos solo se establecen en eventos en los que se aplican, y la estructura del details campo también varía según event_type.

Use esta tabla para consultar la actividad de canalización histórica, crear alertas sobre errores de canalización y correlacionar el comportamiento de la canalización con otras tablas del sistema de Lakeflow.

Ruta de acceso de tabla: system.lakeflow_pipeline_events_preview.pipeline_events

Clave principal: (account_id, pipeline_event_id)

Nombre de la columna Tipo de dato Descripción Notas
account_id string El identificador de la cuenta a la que pertenece este evento de canalización
workspace_id string El identificador del área de trabajo al que pertenece este evento de canalización
pipeline_id string Identificador de la canalización que emitió el evento
update_id string Identificador de la actualización de canalización que emitió el evento
pipeline_event_id string Identificador único global del evento
event_type string Tipo de evento (por ejemplo, flow_progress, update_progress, create_update) Consulta los valores de tipo de evento para el conjunto completo de valores.
origin estructura Metadatos contextuales sobre el origen del evento, como proveedor de nube, región, tipo de canalización, nombres de tabla o flujo, y otros identificadores Véase campos estructurados de origen.
message string Descripción legible del evento Puede estar vacío para algunos eventos.
level string Nivel de gravedad del evento Uno de INFO, WARN, ERROR, METRICS. Ver valores de nivel.
maturity_level string Estabilidad del esquema de eventos Uno de STABLE, EVOLVING, DEPRECATED. Ver Valores de nivel de madurez.
error estructura Detalles del error. Rellenado solo para eventos que contienen información de error Véase Campos de estructuras de error.
details variante Carga específica del evento. Los campos que contiene dependen de event_type Ver campo Detalles.
event_time timestamp Hora en que la canalización emitió el evento Zona horaria registrada como +00:00 (UTC).

Campos estructurales de origen

Subcampo Tipo de dato Descripción
cloud string Proveedor de nube (por ejemplo, AWS, AZURE, GCP)
region string Región del proveedor de nube
org_id bigint Id. de organización del área de trabajo
pipeline_type string Tipo de canalización
pipeline_name string Nombre proporcionado por el usuario de la canalización
cluster_id string Identificador del clúster de proceso que respalda la actualización de la canalización
maintenance_id string Identificador de la actualización de mantenimiento, si el evento procede de una ejecución de mantenimiento
dataset_name string Nombre del conjunto de datos (tabla o vista) al que hace referencia el evento.
sink_name string El nombre del receptor al que hace referencia el evento
catalog_name string Nombre del catálogo de Unity
schema_name string El nombre del esquema del catálogo de Unity
flow_id string El identificador del flujo al que hace referencia el evento
flow_name string Nombre del flujo al que hace referencia el evento
batch_id bigint Identificador de microproceso para flujos de streaming. bigint para compatibilidad
request_id string Identificador de solicitud que inició la acción
materialization_name string Nombre de materialización
operation_id string El identificador de la operación
source_name string El nombre del origen de datos
uc_table_id string Identificador de tabla del catálogo de Unity
ingestion_source_type string Tipo de origen de ingesta (por ejemplo, SQL_SERVER, SALESFORCE)
ingestion_source_connection_name string Nombre de conexión para el origen de ingesta
ingestion_source_catalog_name string Nombre del catálogo de origen en el sistema ascendente
ingestion_source_schema_name string Nombre del esquema de origen en el sistema ascendente
ingestion_source_table_name string Nombre de la tabla de origen en el sistema ascendente
ingestion_source_table_version string La versión de la tabla de origen, si procede.

Campos de estructura de error

Subcampo Tipo de dato Descripción
fatal boolean Si el error provocó que la actualización finalizara
exceptions estructura de matriz<> Cadena de excepciones asociadas al error (última causa principal)
exceptions[].sql_state string Código SQLSTATE, si está disponible
exceptions[].error_class string Clase de error de Databricks, si está disponible

Campo de detalles

La details columna es un VARIANTy los campos que contiene dependen de event_type. Para los campos disponibles en cada tipo de evento, consulte Esquema del registro de eventos de canalización. Use la función o la variant_get sintaxis de puntos para leer valores anidados. Consulte las consultas de ejemplo siguientes para ver los patrones de acceso típicos.

La event_type llave envuelve la carga útil. Por ejemplo, las métricas de un flow_progress evento están en $.flow_progress.metrics, no $.metricsen . Incluye la clave de tipo evento en cada ruta.

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

Consultas de ejemplo

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

Patrones de combinación comunes

Combinación con la pipelines tabla para filtrar por nombre de canalización

La pipelines tabla es una dimensión de variación lenta (SCD2). Tome la versión más reciente de cada canalización antes de unirse.

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

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

Configuración de alertas

Puede crear alertas sobre pipeline_events el uso de alertas de SQL de Databricks. Escriba una consulta SQL en pipeline_events (opcionalmente unida a otras tablas del sistema de Lakeflow), programándola en un almacén de SQL y configure un destino de notificación (correo electrónico, Slack, webhook, PagerDuty).

Algunos puntos de partida útiles:

Alerta cuando no llegó ningún evento para una canalización en los últimos N minutos

Úselo para detectar canalizaciones bloqueadas o con errores silenciosas.

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

Alerta cuando el trabajo pendiente de un flujo específico es demasiado alto

El retraso se reporta en flow_progress los eventos como backlog_bytes, y para fuentes de archivo también como backlog_files. Se activa cuando la lectura más reciente supera un umbral (por ejemplo, 100 MB de trabajo sin procesar). No todas las fuentes reportan todas las métricas, así que filtra por la que tu fuente incluye.

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

Para ver qué fuente está detrás, $.flow_progress.metrics.source_metrics es una serie de lecturas por fuente, cada una junto source_name a la de backlog_bytesesa fuente , backlog_records o backlog_files.

Alerta sobre caídas de calidad de datos en una canalización

Cada flow_progress evento informa del número de filas eliminadas por EXPECT … DROP las expectativas. Suma estos datos por conjunto de datos durante una ventana de actualización y alerta cuando el total supere un umbral.

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

Alerta cuando un flujo procesa demasiadas filas

flow_progress los eventos informan metrics.num_output_rows como un recuento por microlote, así que sumar los eventos en una ventana da las filas escritas sobre esa ventana. Crea una alerta para cuando el rendimiento baje por debajo de un mínimo esperado. Por ejemplo, un flujo que normalmente escribe miles de filas por hora pero produce casi cero puede indicar una fuente mal configurada.

Esta consulta solo informa de flujos que emitieron un flow_progress evento con un conteo de filas en la ventana. Un flujo completamente estancado no emite ningún evento, así que empareja esta alerta con la alerta de eventos perdidos que aparece arriba.

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

Alerta cuando la latencia de nuevos datos sea demasiado alta

Para flujos de streaming, flow_progress los eventos informan de latencia en streaming_metrics. stream_latency_ms es el tiempo desde que los datos llegaron aguas arriba hasta que el microlote se comprometió en la tabla Delta. Puedes configurar un disparador para cuando la lectura más reciente supere un umbral (por ejemplo, 5 minutos).

Solo fluye el streaming con un informe stream_latency_msde tiempo de evento etiquetado, y solo cuando las métricas de tiempo SDP están habilitadas. Otros flujos regresan NULL en cada evento, y esta alerta nunca se activa para ellos.

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

Sugerencias para alertas de producción

  • Filtre por pipeline_id (y workspace_id si mantiene alertas por área de trabajo), por lo que cada alerta tiene como destino un ámbito específico en lugar de toda la cuenta.
  • Elija una cadencia de evaluación que coincida con la confidencialidad de la alerta. Use un intervalo corto (por ejemplo, cada 5 minutos) para las señales de error rápido y un intervalo más largo (por ejemplo, cada hora) para las tendencias de calidad de datos y trabajos pendientes. Una condición "desencadenador cuando la consulta devuelve más de 0 filas" funciona en la mayoría de los casos.

Valores de referencia

Valores de nivel

Value Descripción
INFO Actividad de canalización normal (progreso del flujo, transiciones del ciclo de vida de actualización, cambios de configuración).
WARN Problemas no irrecuperables de los que se recuperó la canalización o que pueden requerir atención.
ERROR Errores que impedían que la canalización progresase en un flujo o actualización.
METRICS Medidas cuantitativas emitidas durante la ejecución (recuentos de filas, rendimiento, latencia).

Valores de nivel de madurez

Value Descripción
STABLE El esquema de eventos es estable. No se esperan cambios importantes. Seguro para crear consultas y alertas de producción.
EVOLVING El esquema de eventos puede cambiar en futuras versiones. Úselo con cuidado.
DEPRECATED El tipo de evento o el esquema están en desuso y se quitarán en futuras versiones. Migración fuera de él.

Valores de tipo de evento

El event_type campo es una enumeración. Conjunto completo de valores:

Value Descripción
create_update Se solicitó una nueva actualización de canalización.
update_progress Una actualización de canalización ha pasado a través de un estado de ciclo de vida.
flow_progress Flujo (conjunto de datos) dentro de una actualización que pasa a través de un estado.
flow_definition Metadatos estáticos sobre un flujo.
dataset_definition Metadatos estáticos sobre un conjunto de datos.
sink_definition Metadatos estáticos sobre un receptor de salida.
deprecation La canalización usó una característica en desuso.
autoscale Decisión de escalado automático del clúster.
unsupported_operation Operación que no se admite en la configuración actual.
cluster_resources Ranura de tareas y métricas de escalado automático para el proceso de respaldo.
planning_information Información de la fase de planificación de la actualización.
gc_pressure Presión de recolección de elementos no utilizados en controladores o ejecutores.
abnormal_termination La actualización finalizó anómalamente.
disk_space Presión de espacio en disco en el clúster.
hook_progress Progreso del ciclo de vida de un enlace de canalización.
dataset_life_cycle Evento de ciclo de vida del conjunto de datos.
background_operation Una operación en segundo plano pasa a través de un estado .
remote_api_usage La canalización realizó una llamada API saliente.
operation_progress Progreso de una operación genérica.
stream_progress Progreso de una consulta de streaming que respalda un flujo.
rewind_summary Resumen de una operación de rebobinado de canalización.
advisory Mensaje de aviso del motor.
runtime_details Configuración detallada del entorno de ejecución.
resource_info Información de recursos (clúster, tipo de instancia, etc.).
file_notification_set_up Estado de configuración de notificaciones de archivos (para orígenes de archivos en la nube).
behavior_change_in_spark_connect Notificación de cambio de comportamiento en Spark Connect.
user_action Una acción iniciada por el usuario en la canalización.
user_code_context Contexto sobre el código de usuario asociado al evento.