Referência da tabela do sistema de eventos do pipeline
Beta
Esta tabela do sistema está em Beta.
Este artigo é uma referência para a tabela de sistema pipeline_events, que registra entradas de log de eventos dos LakeFlow Pipelines para pipelines em sua account. Cada linha é um evento imutável do logs de eventos do pipeline, capturando transições de ciclo de vida, progresso de fluxo, métricas de qualidade de dados, erros, recurso de cluster e outros dados operacionais em todos os pipelines e Workspace dentro de uma região.
Requisitos
- Para acessar esta tabela do sistema, os usuários devem:
- Ser ao mesmo tempo um administrador do metastore e um administrador da account, ou
- Ter as permissões
USEeSELECTnos esquemas do sistema. Consulte Conceder acesso às tabelas do sistema.
Tabelas de eventos de pipeline disponíveis
A tabela do sistema de eventos de pipeline reside no esquema lakeflow_pipeline_events_preview durante a versão Beta e é movida para o esquema lakeflow na disponibilidade geral:
Tabela | Descrição | Suporta transmissão | Período de retenção gratuita | Inclui dados globais ou regionais |
|---|---|---|---|---|
pipeline_events (Beta) | Registra entradas de log de eventos do pipeline emitidas por execuções de pipeline. | Sim | 13 meses | Regional |
O esquema é lakeflow_pipeline_events_preview durante a versão Beta. Na disponibilidade geral, a tabela é movida para o esquema lakeflow (o caminho final da tabela será system.lakeflow.pipeline_events). As queries escritas para o esquema Beta devem ser atualizadas quando a tabela for movida.
Referência detalhada do esquema
Esquema da tabela de eventos do pipeline
A tabela de eventos de pipeline é somente de anexo. 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.
Os campos que são preenchidos em uma linha dependem do tipo de evento. error, update_id e muitos subcampos origin.* são definidos somente em eventos onde se aplicam, e a estrutura do campo details também varia por event_type.
Use esta tabela para consultar a atividade histórica do pipeline, criar alertas sobre falhas do 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 | Descrição | Notas |
|---|---|---|---|
| string | O ID da account à qual este evento de pipeline pertence. | |
| string | O ID do workspace ao qual esse evento do pipeline pertence. | |
| string | O ID do pipeline que emitiu o evento | |
| string | O ID da atualização do pipeline que emitiu o evento | |
| string | Identificador globalmente exclusivo para o evento | |
| string | O tipo de evento (por exemplo, | Consulte Valores de tipo de evento para o conjunto completo de valores. |
| struct | Metadados contextuais sobre a origem do evento, como provedor de cloud, região, tipo de pipeline, nomes de tabela ou fluxo e outros identificadores. | Consulte Campos de struct de origem. |
| string | Descrição legível por humanos do evento | Pode estar vazio para alguns eventos. |
| string | Nível de gravidade do evento | Um de |
| string | Estabilidade do esquema de eventos | Um de |
| struct | Detalhes do erro. Preenchido apenas para eventos que contêm informações de erro | Consulte Campos de struct de erro. |
| Variante | Payload específico do evento. Os campos que contém dependem d | Consulte Campo de detalhes. |
| carimbo de data/hora | A hora em que o evento foi emitido pelo pipeline | Fuso horário registrado como |
Campos de estrutura de origem
Subcampo | Tipo de dados | Descrição |
|---|---|---|
| string | Provedor de nuvem (por exemplo, |
| string | Região do provedor de cloud |
| BigInt | ID de organização do workspace |
| string | O tipo de pipeline |
| string | O nome do pipeline fornecido pelo usuário |
| string | O ID do cluster de compute da atualização do pipeline |
| string | O ID da atualização de manutenção, se o evento for de uma execução de manutenção |
| string | O nome do dataset (tabela ou view) ao qual o evento se refere |
| string | O nome do destino ao qual o evento se refere |
| string | O nome do catálogo do Unity Catalog |
| string | O nome do esquema do Unity Catalog |
| string | O ID do fluxo a que o evento se refere |
| string | O nome do fluxo ao qual o evento se refere |
| BigInt | O ID de micro-lotes para fluxos de transmissão. |
| string | O ID da solicitação que iniciou a ação |
| string | O nome da materialização |
| string | ID da operação |
| string | O nome da fonte de dados |
| string | O ID da tabela do Unity Catalog |
| string | O tipo de origem de ingestão (por exemplo, |
| string | O nome da conexão para a fonte de ingestão |
| string | O nome do catálogo de origem no sistema a montante |
| string | O nome do esquema de origem no sistema upstream |
| string | O nome da tabela de origem no sistema upstream |
| string | A versão da tabela de origem, quando aplicável |
Campos de estrutura de erro
Subcampo | Tipo de dados | Descrição |
|---|---|---|
| boolean | Se o erro causou o encerramento da atualização |
| array<struct> | Cadeia de exceções associada ao erro (causa raiz por último) |
| string | Código SQLSTATE, se disponível |
| string | Classe de erro do Databricks, se disponível |
Campo de detalhes
A coluna details é 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 log de eventos do pipeline. Utilize a função variant_get ou a sintaxe de ponto para ler valores aninhados. Veja os exemplos de consultas abaixo para padrões de acesso típicos.
A key event_type envolve o payload. Por exemplo, as métricas de um evento flow_progress estão em $.flow_progress.metrics, não em $.metrics. Inclua a chave key em todos os caminhos.
-- 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 join
Faça a join com a tabela pipelines para filtrar por nome do pipeline
A tabela pipelines é uma dimensão que muda lentamente (SCD). Utilize a última versão de cada pipeline antes de unir.
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
Join 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
Configurar alertas
Você pode criar alertas em pipeline_events usando alertas do Databricks SQL. Escreva uma consulta SQL contra pipeline_events (opcionalmente unida a outras tabelas do sistema Lakeflow), programe-a em um SQL warehouse e configure um destino de notificação (e-mail, 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 parados ou falhando 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
Alertar quando o backlog para um fluxo específico estiver muito alto
O backlog é relatado em flow_progress eventos como backlog_bytes e, para origens de arquivo, também como backlog_files. Acionar trigger quando a leitura mais recente ultrapassar um limite (por exemplo, 100 MB de trabalho não processado). Nem toda origem relata todas as métricas, portanto, filtre pela que sua origem 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á atrasada, $.flow_progress.metrics.source_metrics é uma matriz de leituras por fonte, cada uma com source_name ao lado de backlog_bytes, backlog_records ou backlog_files dessa fonte.
Alerta sobre quedas na qualidade dos dados em um pipeline
Cada evento flow_progress relata o número de linhas descartadas por EXPECT … DROP expectativas. Some esses valores por dataset durante uma janela de atualização e emita um alerta quando o total exceder 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
Alerta quando um fluxo processa poucas linhas
flow_progress os eventos relatam metrics.num_output_rows como uma contagem por micro-lote, portanto, somar os eventos em uma janela fornece as linhas gravadas durante essa janela. Crie um alerta para quando o throughput cair abaixo de um limite mínimo esperado. Por exemplo, um fluxo que normalmente grava milhares de linhas por hora, mas produz quase zero, pode indicar uma origem configurada incorretamente.
Esta query relata apenas fluxos que emitiram um evento flow_progress com uma contagem de linhas na janela. Um fluxo totalmente travado não emite eventos, portanto, combine este alerta com o alerta de eventos ausentes 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
Alerta quando a latência de novos dados for muito alta
Para fluxos de transmissão, flow_progress eventos relatam a latência em streaming_metrics. stream_latency_ms é o tempo desde quando os dados chegaram à origem até quando o microlote foi confirmado na tabela Delta. Você pode definir um trigger para quando a leitura mais recente ultrapassar um limite (por exemplo, 5 minutos).
Apenas fluxos de transmissão com um tempo de evento marcado relatam stream_latency_ms, e somente quando as métricas de tempo SDP estão habilitadas. Outros fluxos retornam NULL em cada evento, e este 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
- Filtrar por
pipeline_id(e porworkspace_id, se os alertas forem mantidos por workspace) para que cada alerta tenha como alvo um escopo específico em vez de toda a account. - Selecione 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, a cada hora) para pendências e tendências de qualidade de dados. A condição "gatilho quando a consulta retorna mais de 0 linhas" funciona para a maioria dos casos.
Referenciar valores
Valores de nível
Valor | Descrição |
|---|---|
| Atividade normal do pipeline (progresso do fluxo, transições do ciclo de vida de atualizações, alterações de configuração). |
| Problemas não fatais dos quais o pipeline se recuperou, ou que podem exigir atenção. |
| Falhas que impediram o pipeline de progredir em um fluxo ou atualização. |
| Medidas quantitativas emitidas durante a execução (contagens de linha, taxa de transferência, latência). |
Valores do nível de maturidade
Valor | Descrição |
|---|---|
| O esquema de eventos é estável. Mudanças radicais não são esperadas. Seguro para construir queries e alertas de produção. |
| O esquema de eventos poderá mudar em versões futuras. Usar com cuidado. |
| O tipo de evento ou esquema foi descontinuado e será removido em versões futuras. Migrar disso. |
Valores do tipo de evento
O campo event_type é uma enumeração. O conjunto completo de valores:
Valor | Descrição |
|---|---|
| Uma nova atualização de pipeline foi solicitada. |
| Uma atualização de pipeline passou por um estado do ciclo de vida. |
| Um fluxo (dataset) em uma atualização passou por um estado. |
| Metadados estáticos sobre um fluxo. |
| Metadados estáticos sobre um dataset. |
| Metadados estáticos sobre um destino de saída. |
| Um recurso descontinuado foi usado pelo pipeline. |
| Decisão de dimensionamento automático de cluster |
| Uma operação que não é compatível na configuração atual. |
| Métricas de slot de tarefa e de autoscale para o compute de apoio. |
| Informações da fase de planejamento para a atualização. |
| Pressão de coleta de lixo no driver ou nos executores. |
| A atualização foi encerrada anormalmente. |
| Pressão de espaço em disco no cluster. |
| Progresso do ciclo de vida de um gancho de pipeline. |
| Um evento de ciclo de vida de dataset. |
| Uma operação em segundo plano passou por um estado. |
| O pipeline fez uma chamada de API de saída. |
| Progresso para uma operação genérica. |
| Progresso para uma consulta de transmissão que sustenta um fluxo |
| Resumo de uma operação de rebobinamento de pipeline. |
| Uma mensagem consultiva do motor. |
| Configuração detalhada de runtime. |
| Informações de recurso (cluster, tipo de instância, etc.). |
| Status de configuração de notificação de arquivo (para fontes de arquivos cloud). |
| Notificação de mudança de comportamento no Spark Connect. |
| Uma ação iniciada pelo usuário contra o pipeline. |
| Contexto sobre o código do usuário associado ao evento. |