Referência da tabela do sistema de eventos de pipeline

Important

Esta tabela do sistema está em Beta.

Este artigo é uma referência para a pipeline_events tabela do sistema, que regista entradas de registo de eventos dos pipelines Lakeflow para pipelines na sua conta. Cada linha é um evento imutável do registo de eventos do pipeline, capturando transições do ciclo de vida, progresso do fluxo, métricas de qualidade dos dados, erros, recursos do cluster e outros dados operacionais em todos os pipelines e espaços de trabalho dentro de uma região.

Requirements

  • Para aceder a esta tabela do sistema, os utilizadores devem ou:

Tabelas de eventos do pipeline disponíveis

A tabela do sistema de eventos do pipeline permanece no lakeflow_pipeline_events_preview esquema durante a Beta e passa para o lakeflow esquema na disponibilidade geral:

Table Description Suporta streaming Período de retenção gratuito Inclui dados globais ou regionais
pipeline_events (Beta) Regista entradas do registo de eventos do pipeline emitidas por execuções do pipeline Yes 13 meses Regional

Note

O esquema está lakeflow_pipeline_events_preview durante o Beta. Na disponibilidade geral, a tabela move-se para o lakeflow esquema (o caminho final da tabela será system.lakeflow.pipeline_events). As consultas escritas contra o esquema Beta devem ser atualizadas quando a tabela se move.

Referência detalhada do esquema

Esquema da tabela de eventos do pipeline

A tabela de eventos do pipeline é apenas para acréscimos. Cada linha regista um único evento emitido por uma atualização do pipeline no momento em que foi emitida, e as linhas nunca são modificadas ou eliminadas no local.

Quais campos são preenchidos numa linha depende do tipo de evento. error, update_id, e muitos origin.* subcorpos são definidos apenas em eventos onde se aplicam, e a estrutura do details corpo também varia por event_type.

Use esta tabela para consultar a atividade histórica do oleoduto, criar alertas sobre falhas do oleoduto e correlacionar o comportamento do oleoduto com outras tabelas do sistema Lakeflow.

Caminho da tabela: system.lakeflow_pipeline_events_preview.pipeline_events

Tonalidade primária: (account_id, pipeline_event_id)

Nome da coluna Tipo de dados Description Notes
account_id cadeia (de caracteres) O ID da conta a que este evento do pipeline pertence
workspace_id cadeia (de caracteres) O ID do espaço de trabalho a que este evento do pipeline pertence
pipeline_id cadeia (de caracteres) O ID do pipeline que emitiu o evento
update_id cadeia (de caracteres) O ID da atualização do pipeline que emitiu o evento
pipeline_event_id cadeia (de caracteres) Identificador globalmente único para o evento
event_type cadeia (de caracteres) O tipo de evento (por exemplo, flow_progress, update_progress, create_update) Veja valores de tipo de evento para o conjunto completo de valores.
origin estrutura Metadados contextuais sobre a origem do evento, como fornecedor de cloud, região, tipo de pipeline, nomes de tabelas ou fluxos, e outros identificadores Ver Campos de estruturas de origem.
message cadeia (de caracteres) Descrição legível para humanos do evento Pode estar vazio para alguns eventos.
level cadeia (de caracteres) Nível de gravidade do evento Um de INFO, WARN, ERROR, METRICS. Ver valores de Nível.
maturity_level cadeia (de caracteres) Estabilidade do esquema de eventos Um dos STABLE, EVOLVING, DEPRECATED. Ver Valores do nível de maturidade.
error estrutura Detalhes do erro. Preenchido apenas para eventos que transportam informação de erro Ver campos de estruturas de erro.
details variante Carga útil específica de evento. Os campos que contém dependem do event_type Ver campo Detalhes.
event_time carimbo de data/hora A hora em que o evento foi emitido pelo oleoduto O fuso horário registado como +00:00 (UTC).

Campos de estruturas de origem

Subcampo Tipo de dados Description
cloud cadeia (de caracteres) Fornecedor de cloud (por exemplo, AWS, AZURE, GCP)
region cadeia (de caracteres) Região do fornecedor de cloud
org_id bigint ID de organização do espaço de trabalho
pipeline_type cadeia (de caracteres) O tipo de gasoduto
pipeline_name cadeia (de caracteres) O nome do pipeline fornecido pelo usuário
cluster_id cadeia (de caracteres) O ID do cluster de computação que suporta a atualização do pipeline
maintenance_id cadeia (de caracteres) O ID da atualização de manutenção, se o evento for de uma execução de manutenção
dataset_name cadeia (de caracteres) O nome do conjunto de dados (tabela ou vista) a que o evento se refere
sink_name cadeia (de caracteres) O nome da pia a que o evento se refere
catalog_name cadeia (de caracteres) O nome do catálogo Unity Catalog
schema_name cadeia (de caracteres) O nome do esquema do Catálogo Unity
flow_id cadeia (de caracteres) O ID do fluxo a que o evento se refere
flow_name cadeia (de caracteres) O nome do fluxo a que o evento se refere
batch_id bigint O micro-batch ID para fluxos em streaming. bigint Para compatibilidade
request_id cadeia (de caracteres) O ID do pedido que iniciou a ação
materialization_name cadeia (de caracteres) O nome de materialização
operation_id cadeia (de caracteres) O ID da operação
source_name cadeia (de caracteres) O nome da fonte de dados
uc_table_id cadeia (de caracteres) O ID da tabela do Catálogo Unity
ingestion_source_type cadeia (de caracteres) O tipo de fonte de ingestão (por exemplo, SQL_SERVER, SALESFORCE)
ingestion_source_connection_name cadeia (de caracteres) O nome da ligação para a fonte de ingestão
ingestion_source_catalog_name cadeia (de caracteres) O nome do catálogo de origem no sistema a montante
ingestion_source_schema_name cadeia (de caracteres) O nome do esquema de origem no sistema a montante
ingestion_source_table_name cadeia (de caracteres) O nome da tabela de origem no sistema a montante
ingestion_source_table_version cadeia (de caracteres) A versão da tabela de origem, quando aplicável

Campos de estruturas de erro

Subcampo Tipo de dados Description
fatal boolean Se o erro causou a terminação da atualização
exceptions Estrutura de array<> Cadeia de exceções associada ao erro (causa raiz por último)
exceptions[].sql_state cadeia (de caracteres) Código SQLSTATE, se disponível
exceptions[].error_class cadeia (de caracteres) Classe de erro Databricks, se disponível

Campo de detalhes

A details coluna é um VARIANT, e os campos que contém dependem do event_type. Para os campos disponíveis em cada tipo de evento, consulte Esquema do registo de eventos do pipeline. Use a variant_get sintaxe da função ou do ponto para ler valores aninhados. Veja as consultas de exemplo abaixo para padrões típicos de acesso.

A event_type chave envolve a carga útil. Por exemplo, as métricas de um flow_progress evento estão em $.flow_progress.metrics, não $.metrics. Inclua a chave do tipo evento em cada caminho.

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

Exemplos de consultas

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

Padrões comuns de junção

Junte-se à pipelines tabela para filtrar pelo nome do pipeline

A pipelines tabela é uma dimensão que muda lentamente (SCD2). Tome a versão mais recente de cada pipeline antes de aderir.

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

Juntar com pipeline_update_timeline em 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

Configuração de alertas

Podes criar alertas usando pipeline_eventsalertas SQL do Databricks. Escreve uma consulta SQL contra pipeline_events (opcionalmente ligada a outras tabelas do sistema Lakeflow), agenda-a num SQL warehouse e configura um destino de notificação (email, Slack, webhook, PagerDuty).

Alguns pontos úteis para começar:

Alerta quando não surgiram eventos para um oleoduto nos últimos N minutos

Use isto para detetar pipelines bloqueados ou a falhar silenciosamente.

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

Avise quando o atraso para um fluxo específico for demasiado elevado

O backlog é reportado em flow_progress eventos como backlog_bytes, e para fontes de ficheiros também como backlog_files. Ativa-se quando a leitura mais recente ultrapassa um limiar (por exemplo, 100 MB de trabalho por processar). Nem todas as fontes reportam todas as métricas, por isso filtra pela que a tua fonte preenche.

-- 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 qual fonte está por trás, $.flow_progress.metrics.source_metrics é um conjunto de leituras por fonte, cada uma com source_name o de backlog_bytes, backlog_records ou backlog_files.

Alerta sobre quedas na qualidade dos dados num pipeline

Cada flow_progress evento reporta o número de carreiras descartadas segundo EXPECT … DROP as expectativas. Soma-os por conjunto de dados numa janela de atualização e alerta quando o total ultrapassar um limiar.

-- 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 quando um fluxo processa poucas linhas

flow_progress os eventos reportam metrics.num_output_rows como uma contagem por micro-lote, por isso, somando os eventos numa janela, as linhas escritas sobre essa janela. Crie um alerta para quando o débito descer abaixo de um mínimo esperado. Por exemplo, um fluxo que normalmente escreve milhares de linhas por hora mas produz quase zero pode indicar uma fonte mal configurada.

Esta consulta apenas reporta fluxos que emitiram um flow_progress evento com contagem de linhas na janela. Um fluxo totalmente parado não emite eventos, por isso emparelhe este alerta com o alerta de eventos em falta acima.

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

Alerte quando a latência de novos dados for demasiado elevada

Para fluxos de streaming, flow_progress os eventos reportam latência em streaming_metrics. stream_latency_ms é o tempo desde que os dados chegaram a montante até ao momento em que o micro-batch se comprometeu na tabela Delta. Pode definir um gatilho para quando a leitura mais recente ultrapassa um limiar (por exemplo, 5 minutos).

Apenas o streaming flui com um relatório stream_latency_msde tempo de evento marcado, e apenas quando as métricas de tempo SDP estão ativadas. Outros fluxos regressam NULL a cada evento, e este alerta nunca é ativado para eles.

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

Dicas para alertas de produção

  • Filtra por pipeline_id (e workspace_id se mantiveres alertas por espaço de trabalho) para que cada alerta tenha como alvo um âmbito específico em vez de toda a conta.
  • Escolha uma cadência de avaliação que corresponda à sensibilidade do alerta. Use um intervalo curto (por exemplo, a cada 5 minutos) para sinais de falha rápida e um intervalo mais longo (por exemplo, horário) para atrasos e tendências de qualidade dos dados. Uma condição de "disparo quando a consulta retorna mais de 0 linhas" funciona na maioria dos casos.

Valores de referência

Valores dos níveis

Value Description
INFO Atividade normal do pipeline (progresso do fluxo, transições do ciclo de vida da atualização, alterações de configuração).
WARN Problemas não fatais dos quais o gasoduto recuperou, ou que possam exigir atenção.
ERROR Falhas que impediam o pipeline de avançar num fluxo ou atualização.
METRICS Medições quantitativas emitidas durante a execução (contagem de linhas, throughput, latência).

Valores do nível de maturidade

Value Description
STABLE O esquema de eventos é estável. Não se esperam mudanças repentinas. É seguro construir consultas e alertas de produção.
EVOLVING O esquema de eventos pode mudar em futuras versões. Use com cuidado.
DEPRECATED O tipo de evento ou esquema está obsoleto e será removido em futuras versões. Migra para trás.

Valores do tipo de evento

O event_type campo é uma enumeração. O conjunto completo de valores:

Value Description
create_update Foi solicitada uma nova atualização do pipeline.
update_progress Uma atualização do pipeline passou por um estado de ciclo de vida.
flow_progress Um fluxo (conjunto de dados) dentro de uma atualização transitava por um estado.
flow_definition Metadados estáticos sobre um fluxo.
dataset_definition Metadados estáticos sobre um conjunto de dados.
sink_definition Metadados estáticos sobre um sumidouro de saída.
deprecation Uma funcionalidade obsoleta foi utilizada pelo pipeline.
autoscale Decisão de escalonamento automático do cluster.
unsupported_operation Uma operação que não é suportada na configuração atual.
cluster_resources Métricas de slot de tarefa e autoescala para o cálculo de apoio.
planning_information Informação da fase de planeamento para a atualização.
gc_pressure Pressão de recolha de lixo sobre o condutor ou executores.
abnormal_termination A atualização terminou de forma anormal.
disk_space Pressão no espaço do disco sobre o cluster.
hook_progress Progresso do ciclo de vida de um gancho de pipeline.
dataset_life_cycle Um evento do ciclo de vida do conjunto de dados.
background_operation Uma operação de fundo transitava por um estado.
remote_api_usage O pipeline fazia uma chamada de API de saída.
operation_progress Progresso para uma operação genérica.
stream_progress Progresso para uma consulta de streaming que apoia um fluxo.
rewind_summary Resumo de uma operação de rebobinamento de oleoduto.
advisory Uma mensagem de aviso do motor.
runtime_details Configuração detalhada em tempo de execução.
resource_info Informação de recursos (cluster, tipo de instância, etc.).
file_notification_set_up Estado de configuração de notificações de ficheiro (para fontes de ficheiros na cloud).
behavior_change_in_spark_connect Notificação de mudança de comportamento no Spark Connect.
user_action Uma ação iniciada pelo utilizador contra o pipeline.
user_code_context Contexto sobre o código do utilizador associado ao evento.