Övervaka och bevaka Auto Loader

Auto Loader-pipelines kräver aktiv övervakning för att identifiera problem som växande kvarvarande uppgifter, schemaavvikelser, skadade data och blockerade strömmar innan de påverkar nedströmskonsumenter. Den här sidan beskriver hur du övervakar nyckelmått, frågar efter status på filnivå, skapar instrumentpaneler för observabilitet och felsöker vanliga problem.

Mer information om konfigurationsinformation för produktion finns i Konfigurera Auto Loader för produktionsarbetslaster. Metodtips för konfiguration finns i Metodtips för automatisk inläsning.

Förutsättningar

Flera övervakningsarbetsflöden på den här sidan förlitar sig på cloud_files_state() att observera inmatningstillstånd per fil – inklusive frågor om kvarvarande uppgifter, svarstidsberäkningar och identifiering av schemaavvikelser. cloud_files_state() är en tabellvärdesfunktion som returnerar inläsningstillstånd på filnivå för en Auto Loader-kontrollpunkt. Alla fält är inte tillgängliga som standard. Tillgängligheten beror på din Databricks Runtime-version och konfiguration:

  • Databricks Runtime 18.2 och senare versioner: discovery_time, processed_time och commit_time är tillgängliga automatiskt. På Databricks Runtime 16.4–18.1 är dessa fält endast tillgängliga när cloudFiles.cleanSource är aktiverade.
  • Databricks Runtime 16.4 och senare med cloudFiles.cleanSource aktiverat: archive_time, archive_modeoch move_location är tillgängliga.

Att aktivera cloudFiles.cleanSource medför viss prestandapåverkan. Jämför med dina arbetsbelastningar i en förproduktionsmiljö innan du aktiverar det i produktionsmiljö.

Additionally:

  • Kommentera inmatade data med _metadata kolumnen. Samla in minst file_path och file_modification_time. Se Filmetadatakolumnen.
  • Aktivera _rescued_data och _corrupt_record kolumner.

Viktiga mått för Auto Loader

Tabellen nedan sammanfattar de viktigaste mätvärdena att övervaka för Auto Loader-pipelines. Dessa mätvärden är tillgängliga i StreamingQueryListener-förloppshändelser, och Auto Loader-specifika värden visas under varje källas metrics-mappning.

Mått Vad det säger dig
numFilesOutstanding Antal filer i kvarvarande uppgifter som väntar på att bearbetas
numBytesOutstanding Filloggens storlek i byte
approximateQueueSize Molnködjup (endast filmeddelandeläge)
numInputRows Antal rader som bearbetas per omgång
inputRowsPerSecond Data ankomstfrekvens
processedRowsPerSecond Bearbetningskapacitet
durationMs uppdelning Var tiden spenderas i varje batch

Vad du ska titta efter

Följande mönster indikerar att din pipeline kan behöva åtgärdas.

  • Växande numFilesOutstanding: Eftersläpningen håller på att byggas upp. Din datapipeline släpar efter inflödet av data.
  • processedRowsPerSecond < inputRowsPerSecond: Pipelinen bearbetar data långsammare än den anländer.
  • Stor durationMs.latestOffset: Filidentifieringen är långsam. Överväg att byta till filhändelser.
  • Stor durationMs.addBatch: Databearbetningen är långsam. Överväg att skala beräkning eller optimera transformeringar.

För den fullständiga referensen för mätvärden, se Källmätvärden för Auto Loader.

Fråga efter filnivåtillstånd med cloud_files_state

cloud_files_state()-tabellvärdesfunktionen ger detaljerad information om varje fil som identifierats av Auto Loader. Följande fält är tillgängliga. Fält som markerats som att de kräver Databricks Runtime 16.4 och senare eller 18.2 och senare fylls endast i under de villkor som beskrivs i Förutsättningar.

Fält Type Description
path STRING Sökvägen till filen
size BIGINT Storleken på filen i byte
create_time TIMESTAMP När filen skapades
discovery_time TIMESTAMP När Auto Loader upptäckte filen (Databricks Runtime 16.4 och senare)
processed_time TIMESTAMP När Auto Loader bearbetade filen (Databricks Runtime 16.4 och senare)
commit_time TIMESTAMP När filen skrevs till kontrollpunkten (Databricks Runtime 16.4 och senare)
archive_time TIMESTAMP När filen arkiverades (kräver cloudFiles.cleanSource)
archive_mode STRING MOVE, DELETE, eller NULL (kräver cloudFiles.cleanSource)
move_location STRING Målsökväg när cloudFiles.cleanSource är MOVE
ingestion_state STRING Aktuellt tillstånd för filinläsning

Undersöka filinmatningstillstånd

Följande frågor omfattar vanliga diagnostikscenarier.

Hitta alla obearbetade filer (den aktuella kvarvarande informationen):

SELECT * FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state != 'COMMITTED';

Beräkna genomsnittlig svarstid för inmatning (tid från filskapande till incheckning):

SELECT avg(unix_timestamp(commit_time) - unix_timestamp(create_time)) AS avg_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL AND create_time IS NOT NULL;

Hitta skadade eller överhoppade filer:

SELECT path, ingestion_state, size, create_time
FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state LIKE 'SKIPPED%';

Spåra arkiveringsförloppet (kräver cloudFiles.cleanSource):

SELECT archive_mode, count(*) AS file_count
FROM cloud_files_state('path/to/checkpoint')
GROUP BY archive_mode;

Hitta filer med hög latens från identifiering till incheckning för att hitta flaskhalsar:

SELECT
  path,
  size,
  unix_timestamp(commit_time) - unix_timestamp(discovery_time) AS processing_latency_seconds,
  unix_timestamp(commit_time) - unix_timestamp(create_time) AS end_to_end_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL
ORDER BY end_to_end_latency_seconds DESC
LIMIT 20;

Den fullständiga SQL-referensen finns i cloud_files_state tabellvärdesfunktionen.

Övervaka Auto Loader i Lakeflow-pipelines

Databricks rekommenderar att du använder Lakeflow-pipelines för produktionspipelines i Auto Loader. Så här drar du nytta av de inbyggda övervakningsfunktionerna:

  • Lagra händelseloggen för Lakeflow Pipelines i en Delta-tabell så att den kan användas för att fråga efter observerbarhetsdata. Konfigurera detta via pipelinens avancerade inställningar eller API:et. Mer information finns i Händelselogg för pipeline.

  • Strukturera din pipeline för observerbarhet. En välstrukturerad pipeline för automatisk inläsning i Lakeflow-pipelines innehåller en {table}_source vy (källdefinitionen för automatisk inläsning), en {table}_bronze strömmande tabell (rådatainmatning med _rescued_data och _corrupt_record kolumner), en corrupt_records_sink som placerar rader i karantän med oparserbara data och en {table} ren vy för nedströmsförbrukning.

  • Definiera förväntningar för dina strömmande bronstabeller för att övervaka schemaavvikelser och datakorruption. _rescued_data IS NULL identifierar oväntade schemaändringar och _corrupt_record IS NULL identifierar oförsebara data. Lakeflow pipelines utvärderar dessa förväntningar när data anländer och genererar ett spår för observerbarhet. Du kan konfigurera förväntningar för att varna, släppa rader eller stoppa pipelinen.

När du har skapat vyn event_log_raw för din pipeline använder du följande frågor för Auto Loader-specifika mätvärden.

Övervaka datainmatningens genomströmning per flöde:

SELECT
  origin.flow_name,
  origin.update_id,
  timestamp,
  TRY_CAST(details:flow_progress.metrics.num_output_rows AS BIGINT) AS rows_written
FROM event_log_raw
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC;

Övervaka dataeftersläpning per flöde:

SELECT
  origin.flow_name,
  timestamp,
  DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes
FROM event_log_raw
WHERE event_type = 'flow_progress'
  AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
ORDER BY timestamp DESC;

Sammanfatta förväntningsöverträdelser för att identifiera schemaavvikelser och skadade data:

SELECT
  origin.flow_name,
  explode(from_json(
    details:flow_progress.data_quality.expectations,
    'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
  )) AS expectation
FROM event_log_raw
WHERE event_type = 'flow_progress'
  AND details:flow_progress.data_quality.expectations IS NOT NULL;

Allmän vägledning om övervakning av Lakeflow-pipelines finns i Övervaka pipelines och Pipelinehändelseloggen.

Övervaka Auto Loader med Structured Streaming

När du kör Auto Loader utanför Lakeflow-pipelines använder du följande metoder för övervakning av strukturerad direktuppspelning.

  • Implementera en StreamingQueryListener för att hämta Auto Loader-specifika mätvärden från varje batch genom att läsa från source.metrics.
from pyspark.sql.streaming import StreamingQueryListener

class AutoLoaderMonitor(StreamingQueryListener):
    def onQueryStarted(self, event):
        pass

    def onQueryProgress(self, event):
        for source in event.progress.sources:
            if "CloudFilesSource" in source.description:
                metrics = source.metrics
                files_outstanding = metrics.get("numFilesOutstanding", "0")
                bytes_outstanding = metrics.get("numBytesOutstanding", "0")
                rows_per_sec = source.processedRowsPerSecond
                # Push metrics to your monitoring system (for example, write to a Delta table)

    def onQueryIdle(self, event):
        pass

    def onQueryTerminated(self, event):
        pass

spark.streams.addListener(AutoLoaderMonitor())

Anmärkning

Bearbetning av logik i lyssnare kan göra frågebearbetningen långsammare. Begränsa beräkningar i lyssnaråteranrop och undvik synkrona externa skrivningar där. Skicka i stället lättviktig telemetri asynkront eller lämna över mätvärden till ett separat jobb för lagring.

  • Använd numInputRows, inputRowsPerSecondoch processedRowsPerSecond från källans förlopp för att beräkna dataflödet – filer per sekund och rader per sekund för varje batch.

  • För att beräkna inmatningslatens, jämför create_time och commit_time från cloud_files_state() för latens från slutpunkt till slutpunkt. För att bearbeta svarstider använder du uppdelningen durationMs (till exempel latestOffset, addBatchoch andra rapporterade batchfaser) för att identifiera vilket stadium som är flaskhalsen.

  • Använd df.observe() för att definiera interna datakvalitetsmått direkt på strömmande DataFrame. Mätvärden visas i StreamingQueryListener förloppshändelser under observedMetrics.

from pyspark.sql.functions import count, lit, col

observed_df = df.observe(
    "auto_loader_quality",
    count(lit(1)).alias("total_rows"),
    count(col("_rescued_data")).alias("rescued_rows"),
    count(col("_corrupt_record")).alias("corrupt_rows")
)
  • Använd .queryName() för att ge varje dataström ett unikt namn, så att det blir enklare att skilja mellan Auto Loader-strömmar på fliken Streaming i Spark-gränssnittet och i övervakningspaneler.

Fullständig referens för övervakning av strukturerad direktuppspelning finns i Övervaka strukturerade strömningsfrågor på Azure Databricks.

Skapa en instrumentpanel för observerbarhet

Kombinera data från flera källor för att bygga en heltäckande dashboard för övervakning av dina Auto Loader-pipelines. Den här tabellen visar några föreslagna källor som du kan använda för att strukturera instrumentpanelen för observerbarhet.

Datakälla Observerbarhetsdata
cloud_files_state() Inmatningsstatus på filnivå: tidsstämplar för upptäckt, bearbetning, slutförande och arkivering per fil
Händelselogg för Lakeflow-pipelines Historik för pipelinekörningar, flödesmätvärden per batch och resultat för datakvalitetsförväntningar
Utdatatabeller för pipeline Antal rader och datavolym som skrivits per inmatad tabell

Du kan sedan aggregera observerbarhetsdata i dedikerade tabeller som utgör grunden för instrumentpaneler och aviseringar:

  • Sammanfatta status för pipelinekörningar (lyckade eller misslyckade) över tid, härledda från event_type = 'update_progress' händelser.
  • Sammanställa mått för filinmatning (storlek på kvarvarande uppgifter, dataflöde, svarstid per batch), härledda från cloud_files_state() och event_type = 'flow_progress' händelser.
  • Utveckla tabellstatistik med hjälp av radantal och datavolym per tabell, som härleds från num_output_rows i händelseloggen.
  • Samla in felsökningsinformation från detaljerade felloggar och avvikelser från förväntat beteende per uppdatering, härledda från event_type = 'flow_progress'-händelser där data_quality har fyllts i.

Dessa aggregerade tabeller kan driva en AI/BI-instrumentpanel och SQL-aviseringar. Rekommenderade instrumentpaneler inkluderar tidslinje för pipelinekörningsstatus, inmatningstrend, dataflödestrend, datainmatningssvarstidsdistribution, datakvalitetsmått, schemautvecklingshändelser och filarkiveringsstatus.

Övervaka schemautvecklingshändelser

Använd följande metoder för att identifiera schemaändringar när de inträffar.

  • Icke-NULL-värden i _rescued_data för antalet förväntelseöverträdelser indikerar schemaförskjutning. Sök i händelseloggen efter failed_records > 0 utifrån förväntningen no rescued data.
  • Ändringar i _schemas katalogen i den konfigurerade cloudFiles.schemaLocation (eller inuti kontrollpunkten endast när schemaplatsen inte anges separat) indikerar att schemautvecklingen har skett. Du kan avsöka den här katalogen från ett separat övervakningsjobb.
  • Behandla inte en onQueryTerminated händelse följt av onQueryStarted för samma strömnamn som tillräckliga bevis för schemautvecklingen på egen hand. Strömmar startas om av många orsaker (omstarter av kluster, koddistributioner, tillfälliga lagringsfel). Korrelera omstarter med oberoende signaler – _schemas katalogändringar eller _rescued_data förväntningsöverträdelser – innan du drar slutsatsen att schemautvecklingen inträffade.
  • Använd _metadata.file_path för att identifiera vilka filer som införde schemaändringar. Koppla detta till cloud_files_state() på fältet path för att korrelera schemaändringar med specifika filer och batcher.

Använd den här exempelfrågan för att identifiera den senaste schemaavvikelsen via förväntansöverträdelser:

SELECT
  timestamp,
  origin.flow_name,
  exp.name AS expectation_name,
  exp.failed_records
FROM (
  SELECT
    timestamp,
    origin,
    explode(from_json(
      details:flow_progress.data_quality.expectations,
      'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
    )) AS exp
  FROM event_log_raw
  WHERE event_type = 'flow_progress'
    AND details:flow_progress.data_quality.expectations IS NOT NULL
)
WHERE exp.name = '<rescued-data expectation name>'
  AND exp.failed_records > 0
ORDER BY timestamp DESC;

Konfigurera aviseringar för vanliga problem

Använd Databricks SQL-aviseringar eller pipelinemeddelanden för att identifiera problem innan de påverkar nedströmskonsumenter.

Följande SQL identifierar en växande kvarvarande information och kan användas som grund för en Databricks SQL-avisering. Schemalägg den så att den körs regelbundet (till exempel var 5:e minut) och avisera när resultatet inte är tomt.

-- Alert when backlog exceeds threshold or trends upward across recent batches
WITH recent_backlog AS (
  SELECT
    origin.flow_name,
    timestamp,
    DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes,
    ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
  FROM event_log_raw
  WHERE event_type = 'flow_progress'
    AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
)
SELECT flow_name, backlog_bytes, timestamp
FROM recent_backlog
WHERE rn = 1
  AND backlog_bytes > 1073741824  -- alert when backlog exceeds 1 GB

I följande tabell sammanfattas rekommenderade aviseringsvillkor:

Vad du ska identifiera Så här identifierar du det När du ska avisera
Växande kvarvarande uppgifter numFilesOutstanding pekar uppåt Ihållande ökning jämfört med flera batchar
Stoppad ström Inga förloppshändelser Inga händelser för N-minuter (baserat på förväntat utlösarintervall)
Långa svarstider för inmatning commit_time - create_time Överskrider tröskelvärdet för serviceavtal
Datakvalitetsförsämring Förväntad felfrekvens Ökande andel rader som inte uppfyller förväntningarna
Händelse för schemautveckling _rescued_data IS NOT NULL Alla icke-NULL-värden i antalet förväntade överträdelser
Långsam filsökning durationMs.latestOffset Betydligt högre än baslinjen

Felsökning av vanliga problem

I tabellen nedan beskrivs vanliga problem med Auto Loader-pipelines, deras troliga orsaker och rekommenderade åtgärder för att åtgärda dem.

Issue Möjlig orsak Rekommenderad åtgärd
Kvarvarande uppgifter växer snabbare än bearbetning Otillräcklig beräkningskapacitet, snedfördelade data eller strypta frekvensgränser Skala beräkning, sök efter skevhet med Spark-användargränssnittet och granska maxFilesPerTrigger inställningarna för att kontrollera batchstorleken
Filer som inte identifieras Felkonfigurerade filhändelser, behörighetsproblem eller strömmen har inte körts inom de senaste 7 dagarna Verifiera behörigheter för extern plats, kontrollera konfigurationen av filhändelser i Unity Catalog-användargränssnittet och se till att strömmen körs minst var sjunde dag för att undvika att RocksDB-tillståndet upphör att gälla
Stream-start tar för lång tid Ladda ned stort kontrollpunktstillstånd (RocksDB) Uppgradera till Databricks Runtime 15.3 och senare för asynkron tillståndsinläsning, vilket minskar starttiden med ~90%
Duplicerad filbearbetning Aggressiva cloudFiles.maxFileAge inställningar eller skadade kontrollpunkter Använd en konservativ maxFileAge (minst 90 dagar), verifiera kontrollpunktsintegriteten och undvik livscykelprinciper för kontrollpunktslagring
Schemaändringar som orsakar omstarter i pipelinen Frekventa eller inkompatibla schemaändringar Granska schemaEvolutionMode, växla till addNewColumnsWithTypeWidening för typuppflyttning eller använd typen Variant för mycket dynamiska scheman
Skadade data som ackumuleras i mottagare Problem med källdatakvalitet Kontrollera _corrupt_record karantänsinken efter mönster, granska genereringen av källdata och överväg att lägga till validering uppströms
discovery_time och commit_time inte ifyllda Körs på Databricks Runtime under 18.2 utan cleanSource Uppgradera till Databricks Runtime 18.2 och senare eller aktivera cloudFiles.cleanSource på Databricks Runtime 16.4–18.1

För ytterligare felsökning, se Auto Loader FAQ.