Referência da tabela do sistema de eventos de pipeline

Importante

Esta tabela do sistema está em Beta.

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

Requirements

  • Para acessar essa tabela do sistema, os usuários devem:

Tabelas de eventos de pipeline disponíveis

A tabela de eventos do pipeline permanece no lakeflow_pipeline_events_preview esquema durante o Beta, e se move para o lakeflow esquema na disponibilidade geral:

Tabela Description Dá suporte ao streaming Período de retenção gratuito Inclui dados globais ou regionais
pipeline_events (Beta) Registra entradas de log de eventos de pipeline emitidas por execuções de pipeline Yes 13 meses Regional

Observação

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

Referência de esquema detalhada

Esquema da tabela de eventos de pipeline

A tabela de eventos de pipeline é somente acréscimo. Cada linha registra um único evento emitido por uma atualização de pipeline no momento em que foi emitido e as linhas nunca são modificadas ou excluídas no local.

Quais campos são preenchidos em uma linha depende do tipo de evento. error, update_ide muitos origin.* sub-campos são definidos apenas em eventos em que se aplicam e a estrutura do details campo também varia de acordo event_typecom .

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

Caminho da tabela: system.lakeflow_pipeline_events_preview.pipeline_events

Chave primária: (account_id, pipeline_event_id)

Nome da coluna Tipo de dados Description Notes
account_id cadeia A ID da conta à qual este evento de pipeline pertence
workspace_id cadeia A ID do workspace ao qual este evento de pipeline pertence
pipeline_id cadeia A ID do pipeline que emitiu o evento
update_id cadeia A ID da atualização de pipeline que emitiu o evento
pipeline_event_id cadeia Identificador global exclusivo para o evento
event_type cadeia O tipo de evento (por exemplo, flow_progress, , update_progress) create_update Veja valores do tipo de evento para o conjunto completo de valores.
origin struct Metadados contextuais sobre a origem do evento, como provedor de nuvem, região, tipo de pipeline, nomes de tabela ou fluxo e outros identificadores Veja Campos de estruturas de origem.
message cadeia Descrição legível por humanos do evento Pode estar vazio para alguns eventos.
level cadeia Nível de gravidade do evento Uma opção entre INFO, WARN, ERROR, METRICS. Veja valores de Nível.
maturity_level cadeia Estabilidade do esquema de eventos Um dos STABLE, EVOLVING, DEPRECATED. Veja Valores de nível de maturidade.
error struct Detalhes do erro. Populado somente para eventos que carregam informações de erro Veja Campos de estrutura de erro.
details variante Conteúdo específico do evento. Os campos que ele contém dependem do event_type Veja o campo Detalhes.
event_time carimbo de data/hora A hora em que o evento foi emitido pelo pipeline Fuso horário registrado como +00:00 (UTC).

Campos de estruturas de origem

Sub-campo Tipo de dados Description
cloud cadeia Provedor de nuvem (por exemplo, AWS, AZURE, GCP)
region cadeia Região do provedor de nuvem
org_id bigint ID da organização do workspace
pipeline_type cadeia O tipo de pipeline
pipeline_name cadeia O nome do pipeline fornecido pelo usuário
cluster_id cadeia A ID do cluster de computação que dá suporte à atualização do pipeline
maintenance_id cadeia A ID da atualização de manutenção, se o evento for de uma execução de manutenção
dataset_name cadeia O nome do conjunto de dados (tabela ou exibição) ao qual o evento se refere
sink_name cadeia O nome do coletor ao qual o evento se refere
catalog_name cadeia O nome do catálogo do Catálogo do Unity
schema_name cadeia O nome do esquema do Catálogo do Unity
flow_id cadeia A ID do fluxo ao qual o evento se refere
flow_name cadeia O nome do fluxo ao qual o evento se refere
batch_id bigint A ID do microlote para fluxos de streaming. bigint para compatibilidade
request_id cadeia A ID da solicitação que iniciou a ação
materialization_name cadeia O nome da materialização
operation_id cadeia A ID da operação
source_name cadeia O nome da fonte de dados
uc_table_id cadeia A ID da tabela catálogo do Unity
ingestion_source_type cadeia O tipo de fonte de ingestão (por exemplo, SQL_SERVER, ) SALESFORCE
ingestion_source_connection_name cadeia O nome da conexão para a origem da ingestão
ingestion_source_catalog_name cadeia O nome do catálogo de origem no sistema upstream
ingestion_source_schema_name cadeia O nome do esquema de origem no sistema upstream
ingestion_source_table_name cadeia O nome da tabela de origem no sistema upstream
ingestion_source_table_version cadeia A versão da tabela de origem, quando aplicável

Campos de estruturas de erro

Sub-campo Tipo de dados Description
fatal boolean Se o erro causou o término da atualização
exceptions struct de matriz<> Cadeia de exceções associadas ao erro (causa raiz por último)
exceptions[].sql_state cadeia Código SQLSTATE, se disponível
exceptions[].error_class cadeia Classe de erro do Databricks, se disponível

Campo de detalhes

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

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

Consultas de exemplo

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

Ingressar com a pipelines tabela para filtrar pelo nome do pipeline

A pipelines tabela é uma SCD2 (dimensão de alteração lenta). Pegue a versão mais recente de cada pipeline antes de ingressar.

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

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

Configuração de alertas

Você pode criar alertas sobre pipeline_events como usar alertas SQL do Databricks. Escreva uma consulta pipeline_events SQL (opcionalmente unida a outras tabelas do sistema Lakeflow), agende-a em um sql warehouse e configure um destino de notificação (email, Slack, webhook, PagerDuty).

Alguns pontos de partida úteis:

Alerta quando nenhum evento chegou para um pipeline nos últimos N minutos

Use isso para detectar pipelines travados ou silenciosamente com falha.

-- 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 quando a lista de pendências para um fluxo específico é muito alta

O backlog é reportado em flow_progress eventos como backlog_bytes, e para fontes de arquivos também como backlog_files. Aciona quando a leitura mais recente ultrapassa um limite (por exemplo, 100 MB de trabalho não processado). Nem toda fonte reporta todas as métricas, então filtre pela que sua 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á atrás, $.flow_progress.metrics.source_metrics é uma matriz de leituras por fonte, cada uma com source_name o de , backlog_bytes ou backlog_records.backlog_files

Alerta sobre quedas de qualidade de dados em um pipeline

Cada flow_progress evento relata o número de linhas descartadas pelas EXPECT … DROP expectativas. Somare-os por conjunto de dados durante uma janela de atualização e alerte quando o total ultrapassar um limite.

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

flow_progress os eventos reportam metrics.num_output_rows como uma contagem por micro-lote, então somando os eventos em uma janela são as linhas escritas sobre essa janela. Crie um alerta para quando o throughput cair abaixo de um piso esperado. Por exemplo, um fluxo que normalmente escreve milhares de linhas por hora, mas produz quase zero, pode indicar uma fonte mal configurada.

Essa consulta apenas reporta fluxos que emitiram um flow_progress evento com contagem de linhas na janela. Um fluxo totalmente parado não emite eventos, então pareie esse alerta com o alerta de eventos perdidos 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 estiver muito alta

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

Apenas fluxos de streaming com um relatório stream_latency_msde tempo de evento marcado, e somente quando métricas de tempo SDP estiverem ativadas. Outros fluxos retornam NULL a cada evento, e esse alerta nunca é acionado 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

  • Filtre por pipeline_id (e workspace_id se você mantiver alertas por workspace) para que cada alerta tenha como destino um escopo 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, por hora) para tendências de backlog e qualidade de dados. Uma condição "gatilho quando a consulta retorna mais de 0 linhas" funciona para a maioria dos casos.

Valores de referência

Valores de níveis

Value Description
INFO Atividade de pipeline normal (progresso do fluxo, transições de ciclo de vida de atualização, alterações de configuração).
WARN Problemas não fatais dos quais o pipeline se recuperou ou que podem exigir atenção.
ERROR Falhas que impediram que o pipeline progredisse em um fluxo ou atualização.
METRICS Medidas quantitativas emitidas durante a execução (contagens de linhas, taxa de transferência, latência).

Valores de nível de maturidade

Value Description
STABLE O esquema de eventos é estável. Alterações interruptivas não são esperadas. Seguro para criar consultas de produção e alertas.
EVOLVING O esquema de eventos pode ser alterado em versões futuras. Use com cuidado.
DEPRECATED O tipo de evento ou esquema foi preterido e será removido em versões futuras. Migre-o.

Valores do tipo de evento

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

Value Description
create_update Uma nova atualização de pipeline foi solicitada.
update_progress Uma atualização de pipeline passou por um estado de ciclo de vida.
flow_progress Um fluxo (conjunto de dados) em uma atualização passou 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 coletor de saída.
deprecation Um recurso preterido foi usado pelo pipeline.
autoscale Decisão de dimensionamento automático do cluster.
unsupported_operation Uma operação que não tem suporte na configuração atual.
cluster_resources Slot de tarefa e métricas de dimensionamento automático para a computação de backup.
planning_information Informações de fase de planejamento para a atualização.
gc_pressure Pressão de coleta de lixo sobre driver ou executores.
abnormal_termination A atualização foi encerrada anormalmente.
disk_space Pressão de espaço em disco no cluster.
hook_progress Progresso do ciclo de vida de um gancho de pipeline.
dataset_life_cycle Um evento de ciclo de vida do conjunto de dados.
background_operation Uma operação em segundo plano passou por um estado.
remote_api_usage O pipeline fez uma chamada à API de saída.
operation_progress Progresso para uma operação genérica.
stream_progress Progresso para uma consulta de streaming que faz backup de um fluxo.
rewind_summary Resumo de uma operação de retrocesso de pipeline.
advisory Uma mensagem de consultoria do mecanismo.
runtime_details Configuração detalhada do runtime.
resource_info Informações de recurso (cluster, tipo de instância etc.).
file_notification_set_up Status de configuração de notificação de arquivo (para fontes de arquivos de nuvem).
behavior_change_in_spark_connect Notificação de alteração de comportamento no Spark Connect.
user_action Uma ação iniciada pelo usuário em relação ao pipeline.
user_code_context Contexto sobre o código do usuário associado ao evento.