Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
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:
- Ser um administrador do metastore e um administrador de conta ou
- Tenha permissões de
USEeSELECTnos esquemas do sistema. Veja Conceder acesso às tabelas do sistema.
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(eworkspace_idse 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. |