Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
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- enSELECTmachtigingen 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(enworkspace_idals 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. |