Naslaginformatie over systeemtabel voor pijplijn gebeurtenissen

Important

Deze systeemtabel bevindt zich in de bètaversie.

Dit artikel is een verwijzing naar de pipeline_events systeemtabel, waarin gebeurtenislogboekvermeldingen van Lakeflow-pijplijnen voor pijplijnen in uw account worden vastgelegd. Elke rij is een onveranderbare gebeurtenis uit het gebeurtenislogboek van de pijplijn, het vastleggen van levenscyclusovergangen, stroomvoortgang, metrische gegevens over gegevenskwaliteit, fouten, clusterresources en andere operationele gegevens in alle pijplijnen en werkruimten binnen een regio.

Requirements

  • Gebruikers moeten het volgende doen om toegang te krijgen tot deze systeemtabel:
    • Zowel een metastore-beheerder als een accountbeheerder zijn, of
    • Beschikken over USE- en SELECT machtigingen voor de systeemschema's. Zie Verleen toegang tot systeemtabellen.

Beschikbare pijplijn gebeurtenistabellen

De systeemtabel voor pijplijngebeurtenissen leeft in het lakeflow_pipeline_events_preview schema tijdens de bèta, en schakelt over naar het lakeflow schema bij algemene beschikbaarheid:

Table Description Ondersteunt streaming Gratis bewaarperiode Bevat globale of regionale gegevens
pipeline_events (Beta) Registreert vermeldingen in pijplijnlogboeken die worden verzonden door pijplijnuitvoeringen Yes 13 maanden Regional

Note

Het schema is lakeflow_pipeline_events_preview tijdens de bèta. Bij algemene beschikbaarheid verplaatst de tabel naar het lakeflow schema (het uiteindelijke tabelpad zal zijn system.lakeflow.pipeline_events). Queries die op het Beta-schema zijn geschreven, moeten worden bijgewerkt wanneer de tabel verplaatst.

Gedetailleerde schemareferentie

Tabelschema pijplijn gebeurtenissen

De tabel met pijplijngebeurtenissen is alleen toevoegen. Elke rij registreert één gebeurtenis die wordt verzonden door een pijplijnupdate op het moment dat deze werd verzonden, en rijen worden nooit gewijzigd of verwijderd.

Welke velden in een rij worden ingevuld, is afhankelijk van het gebeurtenistype. error, update_iden veel origin.* subvelden zijn alleen ingesteld op gebeurtenissen waarop ze van toepassing zijn, en de structuur van het details veld varieert ook per event_type.

Gebruik deze tabel om een query uit te voeren op historische pijplijnactiviteit, om waarschuwingen te bouwen voor pijplijnfouten en om pijplijngedrag te correleren met andere Lakeflow-systeemtabellen.

Tabelpad: system.lakeflow_pipeline_events_preview.pipeline_events

Primaire sleutel: (account_id, pipeline_event_id)

Kolomnaam Gegevenstype Description Notes
account_id string De id van het account waartoe deze pijplijn gebeurtenis behoort
workspace_id string De id van de werkruimte waartoe deze pijplijn gebeurtenis behoort
pipeline_id string De id van de pijplijn die de gebeurtenis heeft verzonden
update_id string De id van de pijplijnupdate die de gebeurtenis heeft verzonden
pipeline_event_id string Globaal unieke id voor de gebeurtenis
event_type string Het type gebeurtenis (bijvoorbeeld flow_progress, update_progress, create_update) Zie Gebeurtenistypewaarden voor de volledige set waarden.
origin struct Contextuele metagegevens over de oorsprong van de gebeurtenis, zoals cloudprovider, regio, pijplijntype, tabel- of stroomnamen en andere id's Zie Origin structvelden.
message string Leesbare beschrijving van de gebeurtenis Kan leeg zijn voor sommige gebeurtenissen.
level string Ernstniveau van de gebeurtenis Een vanINFO, WARN, ERROR, . METRICS Zie Niveauwaarden.
maturity_level string Stabiliteit van het gebeurtenisschema Eén van STABLE, EVOLVING, . DEPRECATED Zie Volwassenheidsniveauwaarden.
error struct Foutdetails. Alleen ingevuld voor gebeurtenissen die foutinformatie bevatten Zie Foutstructvelden.
details variant Gebeurtenisspecifieke nettolading. De velden die deze bevat, zijn afhankelijk van de event_type Zie het veld Details.
event_time timestamp De tijd waarop de gebeurtenis is verzonden door de pijplijn Tijdzone vastgelegd als +00:00 (UTC).

Oorsprongsstructvelden

Subveld Gegevenstype Description
cloud string Cloudprovider (bijvoorbeeldAWS, , AZUREGCP)
region string Regio van cloudprovider
org_id Bigint Organisatie-id van werkruimte
pipeline_type string Het type pijplijn
pipeline_name string De door de gebruiker opgegeven naam van de pijplijn
cluster_id string De id van het rekencluster die de pijplijnupdate back-upt
maintenance_id string De id van de onderhoudsupdate als de gebeurtenis afkomstig is van een onderhoudsuitvoering
dataset_name string De naam van de gegevensset (tabel of weergave) waarnaar de gebeurtenis verwijst
sink_name string De naam van de sink waarnaar de gebeurtenis verwijst
catalog_name string De naam van de Catalogus van Unity
schema_name string De naam van het Unity Catalog-schema
flow_id string De id van de stroom waarnaar de gebeurtenis verwijst
flow_name string De naam van de stroom waarnaar de gebeurtenis verwijst
batch_id Bigint De microbatch-id voor streamingstromen. bigint voor compatibiliteit
request_id string De aanvraag-id waarmee de actie is gestart
materialization_name string De materialisatienaam
operation_id string De bewerkings-id
source_name string De naam van de gegevensbron
uc_table_id string De tabel-id van Unity Catalog
ingestion_source_type string Het type opnamebron (bijvoorbeeld SQL_SERVER, SALESFORCE)
ingestion_source_connection_name string De verbindingsnaam voor de opnamebron
ingestion_source_catalog_name string De naam van de broncatalogus in het upstream-systeem
ingestion_source_schema_name string De naam van het bronschema in het upstream-systeem
ingestion_source_table_name string De naam van de brontabel in het upstream-systeem
ingestion_source_table_version string De versie van de brontabel, indien van toepassing

Foutstructvelden

Subveld Gegevenstype Description
fatal boolean Of de fout heeft veroorzaakt dat de update werd beëindigd
exceptions matrixstruct<> Keten van uitzonderingen die zijn gekoppeld aan de fout (laatste hoofdoorzaak)
exceptions[].sql_state string SQLSTATE-code, indien beschikbaar
exceptions[].error_class string Databricks-foutklasse, indien beschikbaar

Detailsveld

De details kolom is een VARIANT, en de velden die deze bevat, zijn afhankelijk van de event_type. Zie het schema van het gebeurtenislogboek voor pijplijn voor de velden die beschikbaar zijn onder elk gebeurtenistype. Gebruik de syntaxis van de variant_get functie of punt om geneste waarden te lezen. Zie de onderstaande voorbeeldquery's voor typische toegangspatronen.

De event_type sleutel omhult de lading. Bijvoorbeeld, de metrics van een flow_progress gebeurtenis zijn bij $.flow_progress.metrics, niet $.metrics. Neem de gebeurtenistype-sleutel op in elk pad.

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

Voorbeeldvragen

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

Algemene joinpatronen

Samenvoegen met de pipelines tabel om te filteren op pijplijnnaam

De pipelines tabel is een langzaam veranderende dimensie (SCD2). Neem de nieuwste versie van elke pijplijn voordat u deelneemt.

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

Deelnemen met pipeline_update_timeline aan 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

Waarschuwingen instellen

U kunt waarschuwingen bouwen voor pipeline_events het gebruik van Databricks SQL-waarschuwingen. Schrijf een SQL-query op pipeline_events (eventueel gekoppeld aan andere Lakeflow-systeemtabellen), plan deze in een SQL-warehouse en configureer een meldingsdoel (e-mail, Slack, webhook, PagerDuty).

Enkele nuttige uitgangspunten:

Waarschuwing wanneer er in de afgelopen N minuten geen gebeurtenissen voor een pijplijn zijn aangekomen

Gebruik deze optie om vastgelopen of mislukte pijplijnen op de achtergrond te detecteren.

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

Waarschuwing wanneer de achterstand voor een specifieke stroom te hoog is

Backlog wordt gerapporteerd op flow_progress gebeurtenissen als backlog_bytes, en voor bestandsbronnen ook als backlog_files. Trigger wanneer de meest recente meting een drempel overschrijdt (bijvoorbeeld 100 MB onverwerkt werk). Niet elke bron rapporteert elke metriek, dus filter op degene die je bron vult.

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

Om te zien welke bron achterloopt, $.flow_progress.metrics.source_metrics is een reeks per-bron-lezingen, elk met source_name naast die bron , backlog_bytesbacklog_records of backlog_files.

Waarschuwing over dalingen van gegevenskwaliteit in een pijplijn

Elke flow_progress gebeurtenis rapporteert het aantal rijen dat is verwijderd door EXPECT … DROP verwachtingen. Tel deze op per dataset over een updatevenster en waarschuw wanneer het totaal een drempel overschrijdt.

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

Waarschuw wanneer een flow te weinig rijen verwerkt

flow_progress Gebeurtenissen rapporteren metrics.num_output_rows als een aantal per microbatch, dus het optellen van de gebeurtenissen in een venster geeft de rijen die over dat venster zijn geschreven. Maak een waarschuwing wanneer de doorvoer onder een verwachte vloer daalt. Bijvoorbeeld, een flow die normaal gesproken duizenden rijen per uur schrijft maar bijna nul produceert, kan wijzen op een verkeerd geconfigureerde bron.

Deze query rapporteert alleen stromen die een flow_progress gebeurtenis met een rijaantal in het venster hebben veroorzaakt. Een volledig gestagneerde stroom zendt geen gebeurtenissen uit, dus combineer deze alert met de melding van ontbrekende gebeurtenissen hierboven.

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

Waarschuw wanneer de nieuwe data-latentie te hoog is

Voor streaming flows rapporteren flow_progress events latentie in streaming_metrics. stream_latency_ms is de tijd vanaf het moment dat data upstream landde tot het moment dat de microbatch zich aanmeldde voor de Delta-tabel. Je kunt een trigger instellen wanneer de meest recente meting een drempel overschrijdt (bijvoorbeeld 5 minuten).

Alleen streamen met een getagd event time-rapport stream_latency_ms, en alleen wanneer SDP-tijdmetrics zijn ingeschakeld. Bij elk evenement keren er andere stromen terug NULL , en deze waarschuwing gaat nooit voor hen af.

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

  • Filter op pipeline_id (en workspace_id als u waarschuwingen per werkruimte onderhoudt), zodat elke waarschuwing een specifiek bereik heeft in plaats van het hele account.
  • Kies een evaluatiefrequentie die overeenkomt met de gevoeligheid van de waarschuwing. Gebruik een kort interval (bijvoorbeeld om de 5 minuten) voor snel-failsignalen en een langer interval (bijvoorbeeld elk uur) voor achterstands- en gegevenskwaliteitstrends. Een voorwaarde 'trigger wanneer query meer dan 0 rijen retourneert' werkt voor de meeste gevallen.

Referentiewaarden

Niveauwaarden

Value Description
INFO Normale pijplijnactiviteit (stroomvoortgang, overgangen van levenscyclus bijwerken, configuratiewijzigingen).
WARN Niet-fatale problemen waaruit de pijplijn is hersteld of waarvoor mogelijk aandacht nodig is.
ERROR Fouten waardoor de pijplijn geen voortgang kan boeken in een stroom of update.
METRICS Kwantitatieve metingen die worden verzonden tijdens de uitvoering (aantal rijen, doorvoer, latentie).

Waarden van volwassenheidsniveaus

Value Description
STABLE Het gebeurtenisschema is stabiel. Belangrijke wijzigingen worden niet verwacht. Veilig voor het bouwen van productiequery's en waarschuwingen.
EVOLVING Het gebeurtenisschema kan in toekomstige releases worden gewijzigd. Gebruik met zorg.
DEPRECATED Het gebeurtenistype of schema is afgeschaft en wordt verwijderd in toekomstige releases. Migreer het uit.

Gebeurtenistypewaarden

Het event_type veld is een opsomming. De volledige set waarden:

Value Description
create_update Er is een nieuwe pijplijnupdate aangevraagd.
update_progress Een pijplijnupdate is overgezet naar een levenscyclusstatus.
flow_progress Een stroom (gegevensset) binnen een update die is overgezet via een status.
flow_definition Statische metagegevens over een stroom.
dataset_definition Statische metagegevens over een gegevensset.
sink_definition Statische metagegevens over een uitvoersink.
deprecation Er is een afgeschafte functie gebruikt door de pijplijn.
autoscale Beslissing over automatische schaalaanpassing van clusters.
unsupported_operation Een bewerking die niet wordt ondersteund in de huidige configuratie.
cluster_resources Taaksite en metrische gegevens voor automatische schaalaanpassing voor de back-up berekenen.
planning_information Informatie over de planningsfase voor de update.
gc_pressure Garbagecollectiondruk op stuurprogramma of uitvoerders.
abnormal_termination De update is abnormaal beëindigd.
disk_space Druk op schijfruimte op het cluster.
hook_progress Voortgang van de levenscyclus van een pijplijnhook.
dataset_life_cycle Een gebeurtenis voor de levenscyclus van een gegevensset.
background_operation Een achtergrondbewerking is overgezet via een status.
remote_api_usage De pijplijn heeft een uitgaande API-aanroep gemaakt.
operation_progress Voortgang voor een algemene bewerking.
stream_progress Voortgang van een streamingquery voor het maken van een stroom.
rewind_summary Samenvatting van een pijplijnherstelbewerking.
advisory Een adviesbericht van de engine.
runtime_details Gedetailleerde runtimeconfiguratie.
resource_info Resourcegegevens (cluster, exemplaartype, enzovoort).
file_notification_set_up Status van installatie van bestandsmeldingen (voor cloudbestandenbronnen).
behavior_change_in_spark_connect Melding over gedrag wijzigen onder Spark Connect.
user_action Een door de gebruiker geïnitieerde actie voor de pijplijn.
user_code_context Context over de gebruikerscode die bij het evenement hoort.