Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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_timeochcommit_timeär tillgängliga automatiskt. På Databricks Runtime 16.4–18.1 är dessa fält endast tillgängliga närcloudFiles.cleanSourceär aktiverade. -
Databricks Runtime 16.4 och senare med
cloudFiles.cleanSourceaktiverat:archive_time,archive_modeochmove_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
_metadatakolumnen. Samla in minstfile_pathochfile_modification_time. Se Filmetadatakolumnen. - Aktivera
_rescued_dataoch_corrupt_recordkolumner.
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}_sourcevy (källdefinitionen för automatisk inläsning), en{table}_bronzeströmmande tabell (rådatainmatning med_rescued_dataoch_corrupt_recordkolumner), encorrupt_records_sinksom 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 NULLidentifierar oväntade schemaändringar och_corrupt_record IS NULLidentifierar 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
StreamingQueryListenerför att hämta Auto Loader-specifika mätvärden från varje batch genom att läsa frånsource.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,inputRowsPerSecondochprocessedRowsPerSecondfrå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_timeochcommit_timefråncloud_files_state()för latens från slutpunkt till slutpunkt. För att bearbeta svarstider använder du uppdelningendurationMs(till exempellatestOffset,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 iStreamingQueryListenerförloppshändelser underobservedMetrics.
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()ochevent_type = 'flow_progress'händelser. - Utveckla tabellstatistik med hjälp av radantal och datavolym per tabell, som härleds från
num_output_rowsi 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ärdata_qualityhar 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_dataför antalet förväntelseöverträdelser indikerar schemaförskjutning. Sök i händelseloggen efterfailed_records > 0utifrån förväntningenno rescued data. - Ändringar i
_schemaskatalogen i den konfigureradecloudFiles.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
onQueryTerminatedhändelse följt avonQueryStartedfö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 –_schemaskatalogändringar eller_rescued_dataförväntningsöverträdelser – innan du drar slutsatsen att schemautvecklingen inträffade. - Använd
_metadata.file_pathför att identifiera vilka filer som införde schemaändringar. Koppla detta tillcloud_files_state()på fältetpathfö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.