Monitorujte a sledujte Auto Loader

Pipeline Auto Loaderu vyžadují aktivní monitorování, aby bylo možné odhalit problémy, jako jsou narůstající fronty, změny schématu, poškozená data a zaseknuté streamy, dříve než ovlivní navazující systémy. Tato stránka popisuje, jak monitorovat klíčové metriky, stav na úrovni souboru dotazu, vytvářet řídicí panely pozorovatelnosti a řešit běžné problémy.

Podrobnosti o konfiguraci pro produkční prostředí najdete v článku Konfigurace Auto Loaderu pro produkční úlohy. Osvědčené postupy konfigurace najdete v části Auto Loader – osvědčené postupy.

Předpoklady

Několik pracovních postupů monitorování na této stránce spoléhá na cloud_files_state() k monitorování stavu ingestu jednotlivých souborů, včetně dotazů na backlog, výpočtů latence a detekce odchylek ve schématu. cloud_files_state() je tabulková funkce, která vrací stav ingestce na úrovni souboru pro kontrolní bod Auto Loaderu. Ve výchozím nastavení nejsou k dispozici všechna jeho pole. Dostupnost závisí na vaší verzi a konfiguraci Databricks Runtime:

  • Databricks Runtime 18.2 a novější: discovery_time, processed_time a commit_time jsou k dispozici automaticky. V Databricks Runtime 16.4–18.1 jsou tato pole dostupná pouze v případě, že cloudFiles.cleanSource je povolená.
  • Databricks Runtime 16.4 a vyšší s povolenou funkcí cloudFiles.cleanSource: k dispozici jsou archive_time, archive_mode a move_location.

Povolení cloudFiles.cleanSource má určitou výkonnostní režii. Než to povolíte v produkci, otestujte to na svých úlohách v předprodukčním prostředí.

Additionally:

  • Přidávání poznámek k přijatým datům ve sloupci _metadata Zachyťte alespoň file_path a file_modification_time. Podívejte se na sloupec metadat souboru .
  • Povolte sloupce _rescued_data a _corrupt_record

Klíčové metriky Auto Loaderu

Následující tabulka shrnuje nejdůležitější metriky ke sledování pro pipeline Auto Loaderu. Tyto metriky jsou k dispozici v událostech průběhu StreamingQueryListener, přičemž hodnoty specifické pro Auto Loader jsou uvedené v mapě metrics každého zdroje.

Metric Co vám to řekne
numFilesOutstanding Počet souborů v backlogu čekajících na zpracování
numBytesOutstanding Velikost backlogu souboru v bajtech
approximateQueueSize Hloubka cloudové fronty (pouze režim oznámení souborů)
numInputRows Počet řádků zpracovaných v dávce
inputRowsPerSecond Rychlost doručení dat
processedRowsPerSecond Propustnost zpracování
durationMs Členění Na co se v každé dávce spotřebuje čas

Co sledovat

Následující příznaky naznačují, že vaše pipeline může vyžadovat vaši pozornost.

  • Narůstající numFilesOutstanding: Backlog narůstá. Vaše pipeline nestíhá příchozí data.
  • processedRowsPerSecond < inputRowsPerSecond: Kanál zpracovává data pomaleji, než dorazí.
  • Velké durationMs.latestOffset: Zjišťování souborů je pomalé. Zvažte přepnutí na události souborů.
  • Velké durationMs.addBatch: Zpracování dat je pomalé. Zvažte škálování výpočetních prostředků nebo optimalizaci transformací.

Kompletní referenci metrik najdete v tématu Zdrojové metriky Auto Loaderu.

Dotazování na stav na úrovni souboru pomocí cloud_files_state

Tabulková funkce cloud_files_state() poskytuje podrobné informace o každém souboru zjištěném nástrojem Auto Loader. K dispozici jsou následující pole. Pole označená jako vyžadující Databricks Runtime 16.4 a vyšší nebo 18.2 a vyšší jsou vyplněná pouze za podmínek popsaných v požadavcích.

Pole Typ Description
path STRING Cesta k souboru
size BIGINT Velikost souboru v bajtech
create_time TIMESTAMP Po vytvoření souboru
discovery_time TIMESTAMP Když Auto Loader zjistil soubor (Databricks Runtime 16.4 a vyšší)
processed_time TIMESTAMP Když Auto Loader zpracoval soubor (Databricks Runtime 16.4 a vyšší)
commit_time TIMESTAMP Když byl soubor potvrzen do kontrolního bodu (Databricks Runtime 16.4 a vyšší)
archive_time TIMESTAMP Při archivaci souboru (vyžaduje cloudFiles.cleanSource)
archive_mode STRING MOVE, DELETEnebo NULL (vyžaduje cloudFiles.cleanSource)
move_location STRING Cílová cesta, pokud cloudFiles.cleanSource je MOVE
ingestion_state STRING Aktuální stav příjmu souborů

Prozkoumání stavu příjmu souborů

Následující dotazy se týkají běžných diagnostických scénářů.

Vyhledejte všechny nezpracované soubory (aktuální backlog):

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

Vypočítat průměrnou latenci příjmu dat (doba od vytvoření souboru do potvrzení):

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;

Vyhledání poškozených nebo přeskočených souborů:

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

Sledování průběhu archivace (vyžaduje cloudFiles.cleanSource):

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

Vyhledejte soubory s vysokou latencí při zjišťování a potvrzením a identifikujte kritické body:

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;

Úplnou referenci k SQL naleznete v části cloud_files_statefunkce vracející tabulku.

Monitorování Auto Loaderu v pipelinech Lakeflow

Databricks doporučuje používat kanály Lakeflow pro produkční kanály automatického zavaděče. Využití integrovaných možností monitorování:

  • Uložte protokol událostí pipelinů Lakeflow do tabulky Delta, aby z něj bylo možné pomocí dotazů získávat data pro monitorování. Nakonfigurujte to v pokročilých nastaveních pipeline nebo prostřednictvím rozhraní API. Podrobnosti naleznete v protokolu událostí pipeline.

  • Strukturujte kanál pro pozorovatelnost. Dobře strukturovaný kanál Auto Loaderu v kanálech Lakeflow zahrnuje zobrazení {table}_source (definici zdroje Auto Loaderu), streamovací tabulku {table}_bronze (pro ingesti nezpracovaných dat se sloupci _rescued_data a _corrupt_record), prvek corrupt_records_sink, který přesouvá řádky s nenaparsovatelnými daty do karantény, a čisté zobrazení {table} pro následné využití.

  • Nastavte očekávání u vašich bronzových streamovacích tabulek pro monitorování odchylek schématu a poškození dat. _rescued_data IS NULL detekuje neočekávané změny schématu a _corrupt_record IS NULL detekuje neparsovatelná data. Kanály Lakeflow tyto očekávání vyhodnocují při příchodu dat a generují stopu pozorovatelnosti. Můžete nakonfigurovat očekávání tak, aby vyvolala upozornění, vyřadila řádky nebo způsobila selhání kanálu.

Po vytvoření zobrazení event_log_raw pro váš kanál použijte následující dotazy pro metriky specifické pro Auto Loader.

Monitorujte propustnost příjmu dat pro každý tok:

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;

Monitorujte objem nevyřízených dat pro každý tok:

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;

Shrňte porušení očekávání, abyste odhalili odchylky schématu a poškozená 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;

Obecné pokyny k monitorování kanálů Lakeflow najdete v tématu Monitorování kanálů a protokolu událostí kanálu.

Monitorujte Auto Loader pomocí Structured Streaming

Při spouštění nástroje Auto Loader mimo pipeline Lakeflow použijte následující přístupy ke monitorování nástroje Structured Streaming.

  • Implementujte StreamingQueryListener, chcete-li zachytávat metriky specifické pro Auto Loader z každé dávky čtením z 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())

Note

Logika zpracování v naslouchacích procesech může zpomalit zpracování dotazů. Omezte výpočty v callbackách posluchače a vyhněte se v nich synchronním externím zápisům; místo toho asynchronně odesílejte nenáročná telemetrická data nebo předejte metriky samostatné úloze k jejich uložení.

  • Použijte numInputRows, inputRowsPerSecond a processedRowsPerSecond z průběhu zdroje k výpočtu propustnosti — souborů za sekundu a řádků za sekundu pro každou dávku.

  • Pro výpočet latence ingestování dat porovnejte create_time a commit_time z cloud_files_state() pro latenci typu end-to-end. Ke zjištění latence zpracování použijte durationMs členění (například latestOffset, addBatch a další vykazované fáze dávkového zpracování), abyste určili, která fáze je úzkým místem.

  • Pomocí df.observe() definujte inline metriky kvality dat přímo ve streamovaném DataFrame. Metriky se zobrazují v StreamingQueryListener probíhajících událostech v části 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")
)
  • Pomocí .queryName() přiřaďte každému streamu jedinečný název, aby bylo snazší rozlišit streamy Auto Loaderu na kartě Streaming v uživatelském rozhraní Spark a v monitorovacích dashboardech.

Úplné referenční informace k monitorování strukturovaného streamování najdete v tématu Monitorování dotazů strukturovaného streamování na Azure Databricks.

Vytvoření řídicího panelu pozorovatelnosti

Kombinujte data z více zdrojů a vytvořte komplexní dashboard pro monitorování vašich Auto Loader pipeline. Tato tabulka zobrazuje některé navrhované zdroje, které můžete použít ke strukturování řídicího panelu pozorovatelnosti.

Zdroj dat Data pozorovatelnosti
cloud_files_state() Stav ingesce dat na úrovni souboru: časová razítka detekce, zpracování, zápisu a archivace pro každý soubor
Protokol událostí kanálů Lakeflow Historie spuštění kanálu, metriky toku pro jednotlivé dávky a očekávané výsledky kvality dat
Výstupní tabulky kanálu Počty řádků a objem dat zapsané pro každou ingestovanou tabulku

Pak můžete agregovat data pozorovatelnosti do vyhrazených tabulek, které slouží jako základ pro řídicí panely a výstrahy:

  • Shrňte stavy běhů pipeline (úspěch nebo selhání) v čase, odvozené z událostí event_type = 'update_progress'.
  • Agregované metriky příjmu souborů (velikost backlogu, propustnost, latence na dávku) odvozené od cloud_files_state() událostí a event_type = 'flow_progress' událostí.
  • Vyvíjejte statistiky tabulky pomocí počtu řádků a objemu dat na tabulku odvozenou z num_output_rows protokolu událostí.
  • Shromážděte informace pro ladění z podrobných protokolů chyb a porušení očekávání pro každou aktualizaci, odvozených z událostí event_type = 'flow_progress' s vyplněným data_quality.

Tyto agregované tabulky mohou sloužit jako podklad pro dashboard AI/BI a SQL upozornění. Mezi doporučené panely řídicích panelů patří časová osa spuštění kanálu, trend backlogu příjmu dat, trend propustnosti, distribuce latence příjmu dat, metriky kvality dat, události vývoje schématu a stav archivace souborů.

Monitorování událostí vývoje schématu

K detekci změn schématu použijte následující přístupy.

  • Nenulové hodnoty v počtech porušení očekávání v _rescued_data indikují odchylku schématu. Vyhledejte v protokolu událostí položku failed_records > 0 pro očekávání no rescued data.
  • _schemas Změny adresáře uvnitř nakonfigurovaného cloudFiles.schemaLocation (nebo uvnitř kontrolního bodu pouze v případech, kdy umístění schématu není nastaveno samostatně) značí, že došlo k vývoji schématu. Tento adresář můžete pravidelně kontrolovat pomocí samostatné monitorovací úlohy.
  • Nezacházejte s událostí onQueryTerminated následovanou onQueryStarted stejným názvem datového proudu jako dostatečný důkaz vývoje schématu samostatně. Streamy se restartují z mnoha důvodů (restartování clusteru, nasazení kódu, přechodné chyby úložiště). Dávejte restarty do souvislosti s nezávislými signály — _schemas změnami v adresáři nebo _rescued_data porušeními očekávání — než dospějete k závěru, že došlo k evoluci schématu.
  • Slouží _metadata.file_path k identifikaci souborů, které zavedly změny schématu. Spojte ho s cloud_files_state() polem a proveďte korelaci změn schématu path s konkrétními soubory a dávkami.

Pomocí tohoto ukázkového dotazu odhalíte nedávné změny schématu na základě porušení očekávání:

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;

Nastavení upozornění pro běžné problémy

Pomocí výstrah Databricks SQL nebo oznámení pipeline můžete odhalit problémy dříve, než ovlivní navazující příjemce.

Následující příkaz SQL slouží k detekci narůstající fronty nevyřízených položek a lze jej použít jako základ pro vytvoření výstrahy v Databricks SQL. Naplánujte pravidelné spuštění (například každých 5 minut) a upozorňování, když je výsledek neprázdný.

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

Následující tabulka shrnuje doporučené podmínky upozornění:

Co zjistit Jak ho zjistit Kdy upozornit
Rostoucí backlog numFilesOutstanding trend směrem nahoru Setrvalé zvýšení během několika šarží
Zaseknutý datový tok Žádné události o průběhu Žádné události pro N minut (na základě očekávaného intervalu triggeru)
Vysoká latence příjmu dat commit_time - create_time Překračuje váš práh SLA
Snížení kvality dat Očekávaná míra selhání Rostoucí procento řádků nesplňujících očekávání
Událost vývoje schématu _rescued_data IS NOT NULL Všechny hodnoty, které nejsou null v počtu očekávaných porušení
Pomalé zjišťování souborů durationMs.latestOffset Výrazně vyšší než výchozí hodnota

Řešení běžných potíží

Následující tabulka popisuje běžné problémy s pipeline Auto Loaderu, jejich pravděpodobné příčiny a doporučené kroky k jejich odstranění.

Issue Možná příčina Doporučená akce
Objem nevyřízených položek roste rychleji než rychlost jejich zpracování Nedostatečná výpočetní kapacita, nerovnoměrné rozložení dat nebo omezené limity rychlosti Škálujte výpočetní prostředky, zkontrolujte nerovnoměrné rozložení dat pomocí uživatelského rozhraní Spark a zkontrolujte nastavení maxFilesPerTrigger pro řízení velikosti dávek
Nezjišťované soubory Nesprávná konfigurace událostí souboru, problém s oprávněními nebo datový proud nebyl spuštěn během 7 dnů Ověřte oprávnění k externímu umístění, zkontrolujte nastavení událostí souborů v uživatelském rozhraní katalogu Unity a ujistěte se, že stream běží alespoň každých 7 dnů, aby se zabránilo vypršení platnosti stavu RocksDB.
Spuštění streamu trvá příliš dlouho Stažení stavu velkého kontrolního bodu (RocksDB) Upgradujte na Databricks Runtime 15.3 nebo novější kvůli asynchronnímu načítání stavu, které zkracuje dobu spouštění přibližně o 90 %
Duplicitní zpracování souborů Příliš agresivní cloudFiles.maxFileAge nastavení nebo poškození kontrolních bodů Použijte konzervativní nastavení maxFileAge (alespoň 90 dní), ověřte integritu checkpointů a vyhněte se zásadám životního cyklu pro úložiště checkpointů.
Vývoj schématu způsobující restartování kanálu Časté nebo nekompatibilní změny schématu Zkontrolujte schemaEvolutionMode, kvůli povyšování typů přejděte na addNewColumnsWithTypeWidening nebo pro vysoce dynamická schémata použijte typ Variant.
Poškozená data nahromaděná v jímce Problémy s kvalitou zdrojových dat Zkontrolujte vzory v _corrupt_record karanténním úložišti, prověřte generování zdrojových dat a zvažte přidání ověřování na vstupu.
discovery_time a commit_time nejsou vyplněny Spuštění na platformě Databricks Runtime ve verzi nižší než 18.2 bez cleanSource Upgradujte na Databricks Runtime 18.2 a novější nebo povolte cloudFiles.cleanSource v Databricks Runtime 16.4–18.1

Další informace k řešení potíží najdete v Nejčastějších dotazech k Auto Loaderu.