Systemtabellreferens för pipelinehändelser

Viktigt!

Den här systemtabellen finns i Beta.

Den här artikeln är en referens för systemtabellen pipeline_events som registrerar händelseloggposter för Lakeflow-pipelines för pipelines i ditt konto. Varje rad är en oföränderlig händelse från pipelinehändelseloggen, som samlar in livscykelövergångar, flödesförlopp, datakvalitetsmått, fel, klusterresurser och andra driftdata i alla pipelines och arbetsytor i en region.

Requirements

  • För att få åtkomst till den här systemtabellen måste användarna antingen:

Tillgängliga pipelinehändelsetabeller

Pipeline-händelsesystemtabellen finns i schemat lakeflow_pipeline_events_preview under Beta och flyttar till schemat lakeflow vid allmän tillgänglighet:

Bord Description Stöder direktuppspelning Fri kvarhållningsperiod Innehåller globala eller regionala data
pipeline_events (Beta) Registrerar pipelinehändelseloggposter som genereras av pipelinekörningar Yes 13 månader Regionala

Note

Schemat är lakeflow_pipeline_events_preview under Beta. Vid allmän tillgänglighet flyttar tabellen till schemat lakeflow (den slutliga tabellens väg blir system.lakeflow.pipeline_events). Frågor skrivna mot Beta-schemat måste uppdateras när bordet flyttas.

Detaljerad schemareferens

Tabellschema för pipelinehändelser

Tabellen pipelinehändelser är endast tillägg. Varje rad registrerar en enskild händelse som genereras av en pipelineuppdatering när den genereras och rader ändras eller tas aldrig bort på plats.

Vilka fält som fylls i på en rad beror på händelsetypen. error, update_id, och många origin.* underfält anges endast för händelser där de gäller, och strukturen för details fältet varierar också med event_type.

Använd den här tabellen om du vill köra frågor mot historisk pipelineaktivitet, skapa aviseringar om pipelinefel och korrelera pipelinebeteendet med andra Lakeflow-systemtabeller.

Tabellsökväg: system.lakeflow_pipeline_events_preview.pipeline_events

Primärnyckel: (account_id, pipeline_event_id)

Kolumnnamn Datatyp Description Notes
account_id string ID:t för det konto som pipelinehändelsen tillhör
workspace_id string ID:t för arbetsytan som den här pipelinehändelsen tillhör
pipeline_id string ID:t för pipelinen som genererade händelsen
update_id string ID:t för pipelineuppdateringen som genererade händelsen
pipeline_event_id string Globalt unik identifierare för händelsen
event_type string Typ av händelse (till exempel flow_progress, update_progress, create_update) Se Händelsetypvärden för hela uppsättningen av värden.
origin struct Kontextuella metadata om händelsens ursprung, till exempel molnleverantör, region, pipelinetyp, tabell- eller flödesnamn och andra identifierare Se Origin-strukturfält.
message string Beskrivning av händelsen som kan läsas av människor Kan vara tomt för vissa händelser.
level string Allvarlighetsgrad för händelsen En av INFO, WARN, ERROR, METRICS. Se nivåvärden.
maturity_level string Stabilitet i händelseschemat En av STABLE, EVOLVING, DEPRECATED. Se Värden på mognadsnivå.
error struct Felinformation. Endast ifylld för händelser som har felinformation Se Felstruktureringsfält.
details variant Händelsespecifik nyttolast. Vilka fält som den innehåller beror på event_type Se detaljfältet.
event_time timestamp Den tid då händelsen genererades av pipelinen Tidszon som registrerats som +00:00 (UTC).

Ursprungsstrukturfält

Underfält Datatyp Description
cloud string Molnleverantör (till exempel AWS, AZURE, GCP)
region string Molnleverantörsregion
org_id bigint Organisations-ID för arbetsyta
pipeline_type string Typ av pipeline
pipeline_name string Det användaruppgivna namnet på pipelinen
cluster_id string Beräkningskluster-ID:t som stöder pipelineuppdateringen
maintenance_id string ID:t för underhållsuppdateringen, om händelsen kommer från en underhållskörning
dataset_name string Namnet på datamängden (tabell eller vy) som händelsen refererar till
sink_name string Namnet på mottagaren som händelsen refererar till
catalog_name string Katalognamnet för Unity-katalogen
schema_name string Schemanamnet för Unity-katalogen
flow_id string ID för flödet som händelsen refererar till
flow_name string Namnet på flödet som händelsen refererar till
batch_id bigint Mikrobatch-ID:t för strömningsflöden. bigint för kompatibilitet
request_id string Begärande-ID:t som initierade åtgärden
materialization_name string Materialiseringsnamnet
operation_id string Åtgärds-ID
source_name string Namnet på datakällan
uc_table_id string Tabell-ID för Unity Catalog
ingestion_source_type string Typ av inmatningskälla (till exempel SQL_SERVER, SALESFORCE)
ingestion_source_connection_name string Anslutningsnamnet för inmatningskällan
ingestion_source_catalog_name string Källkatalognamnet i det överordnade systemet
ingestion_source_schema_name string Källschemanamnet i det överordnade systemet
ingestion_source_table_name string Källtabellens namn i det överordnade systemet
ingestion_source_table_version string Källtabellversionen, om tillämpligt

Felstrukturfält

Underfält Datatyp Description
fatal boolean Om felet gjorde att uppdateringen avslutades
exceptions matris<struct> Kedja av undantag som är associerade med felet (rotorsaken sist)
exceptions[].sql_state string SQLSTATE-kod, om den är tillgänglig
exceptions[].error_class string Databricks-felklass, om det är tillgängligt

Detaljfältet

Kolumnen details är en VARIANT, och de fält som den innehåller beror på event_type. De fält som är tillgängliga under varje händelsetyp finns i Schema för pipelinehändelselogg. variant_get Använd funktionen eller punktsyntaxen för att läsa kapslade värden. Se exempelfrågorna nedan för vanliga åtkomstmönster.

Nyckeln event_type omsluter nyttolasten. Till exempel är en flow_progress händelses mått vid $.flow_progress.metrics, inte $.metrics. Inkludera händelsetypnyckeln i varje väg.

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

Exempel på förfrågningar

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

Vanliga kopplingsmönster

Koppla med tabellen pipelines för att filtrera efter pipelinenamn

Tabellen pipelines är en långsamt föränderlig dimension (SCD2). Ta den senaste versionen av varje pipeline innan du ansluter.

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

Anslut med pipeline_update_timelineupdate_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

Att ställa in varningar

Du kan skapa aviseringar om pipeline_events hur du använder Databricks SQL-aviseringar. Skriv en SQL-fråga mot pipeline_events (eventuellt ansluten till andra Lakeflow-systemtabeller), schemalägg den på ett SQL-lager och konfigurera ett meddelandemål (e-post, Slack, webhook, PagerDuty).

Några användbara utgångspunkter:

Avisering när inga händelser har anlänt för en pipeline under de senaste N minuterna

Använd det här alternativet för att identifiera fastnade eller tyst misslyckade pipelines.

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

Avisering när kvarvarande uppgifter för ett visst flöde är för högt

Backlogg rapporteras på flow_progress händelser som backlog_bytes, och för filkällor även som backlog_files. Trigga när den senaste avläsningen passerar en tröskel (till exempel 100 MB obehandlat arbete). Inte varje källa rapporterar varje mått, så filtrera på den som din källa fyller på.

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

För att se vilken källa som ligger bakom, $.flow_progress.metrics.source_metrics finns en uppsättning avläsningar per källa, var och en med source_name bredvid källans backlog_bytes, backlog_records eller backlog_files.

Avisering om datakvalitetsfall i en pipeline

Varje flow_progress händelse rapporterar antalet rader som släppts av EXPECT … DROP förväntningar. Summera dessa per dataset under ett uppdateringsfönster och varning när totalen överskrider en tröskel.

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

Varning när ett flöde bearbetar för få rader

flow_progress Händelser rapporteras metrics.num_output_rows som en räkning per mikrobatch, så att summera händelserna i ett fönster ger raderna skrivna över det fönstret. Skapa en varning för när genomströmningen sjunker under en förväntad golvnivå. Till exempel kan ett flöde som normalt skriver tusentals rader per timme men producerar nära noll indikera en felkonfigurerad källa.

Denna fråga rapporterar endast flöden som sände en flow_progress händelse med radräkning i fönstret. Ett helt stillastående flöde sänder inga händelser, så kombinera denna varning med varningen om saknade händelser ovan.

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

Varning när fördröjningen för nya data är för hög

För strömmande flöden rapporterar händelser flow_progress latens i streaming_metrics. stream_latency_ms är tiden från när data landade uppströms till när mikrobatchen satte in sig i Delta-tabellen. Du kan ställa in en trigger när den senaste avläsningen passerar en tröskel (till exempel 5 minuter).

Endast strömmande flöden med en taggad händelsetidsrapport stream_latency_ms, och endast när SDP-tidsmetrik är aktiverade. Andra flöden återkommer NULL vid varje händelse, och denna varning utlöses aldrig för dem.

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

Tips för produktionsaviseringar

  • pipeline_id Filtrera efter (och workspace_id om du underhåller aviseringar per arbetsyta) så att varje avisering riktar sig mot ett specifikt omfång snarare än hela kontot.
  • Välj en utvärderingstakt som matchar aviseringens känslighet. Använd ett kort intervall (till exempel var 5:e minut) för snabba felsignaler och ett längre intervall (till exempel varje timme) för kvarvarande uppgifter och datakvalitetstrender. Ett "utlösare när frågan returnerar fler än 0 rader" fungerar i de flesta fall.

Referensvärden

Nivåvärden

Value Description
INFO Normal pipelineaktivitet (flödesstatus, uppdatering av livscykelövergångar, konfigurationsändringar).
WARN Icke-dödliga problem som pipelinen återställs från, eller som kan kräva uppmärksamhet.
ERROR Fel som hindrade pipelinen från att göra framsteg i ett flöde eller en uppdatering.
METRICS Kvantitativa mått som genereras under körningen (radantal, dataflöde, svarstid).

Värden på mognadsnivå

Value Description
STABLE Händelseschemat är stabilt. Icke-bakåtkompatibla ändringar förväntas inte. Säkert att skapa produktionsfrågor och aviseringar på.
EVOLVING Händelseschemat kan ändras i framtida versioner. Använd med försiktighet.
DEPRECATED Händelsetypen eller schemat är inaktuellt och tas bort i framtida versioner. Migrera från den.

Händelsetypvärden

Fältet event_type är en uppräkning. Den fullständiga uppsättningen värden:

Value Description
create_update En ny pipelineuppdatering begärdes.
update_progress En pipelineuppdatering övergick genom ett livscykeltillstånd.
flow_progress Ett flöde (datauppsättning) i en uppdatering övergick via ett tillstånd.
flow_definition Statiska metadata om ett flöde.
dataset_definition Statiska metadata om en datauppsättning.
sink_definition Statiska metadata om en utdatamottagare.
deprecation En inaktuell funktion användes av pipelinen.
autoscale Beslut om automatisk skalning av kluster.
unsupported_operation En åtgärd som inte stöds i den aktuella konfigurationen.
cluster_resources Mått för aktivitetsfack och autoskalning för beräkningen.
planning_information Planeringsfasinformation för uppdateringen.
gc_pressure Skräpinsamlingstryck på drivrutins- eller exekutorer.
abnormal_termination Uppdateringen avslutades onormalt.
disk_space Diskutrymmestryck på klustret.
hook_progress Livscykel för en pipeline-hook.
dataset_life_cycle En livscykelhändelse för datamängden.
background_operation En bakgrundsåtgärd övergick genom ett tillstånd.
remote_api_usage Pipelinen gjorde ett utgående API-anrop.
operation_progress Förlopp för en allmän åtgärd.
stream_progress Förlopp för en strömmande fråga som stöder ett flöde.
rewind_summary Sammanfattning av en pipelineåterspolningsåtgärd.
advisory Ett rådgivande meddelande från motorn.
runtime_details Detaljerad körningskonfiguration.
resource_info Resursinformation (kluster, instanstyp osv.).
file_notification_set_up Installationsstatus för filmeddelanden (för molnfilkällor).
behavior_change_in_spark_connect Meddelande om beteendeändring under Spark Connect.
user_action En användarinitierad åtgärd mot pipelinen.
user_code_context Kontexten om användarkoden kopplad till händelsen.