Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
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_timeacommit_timejsou k dispozici automaticky. V Databricks Runtime 16.4–18.1 jsou tato pole dostupná pouze v případě, žecloudFiles.cleanSourceje povolená. -
Databricks Runtime 16.4 a vyšší s povolenou funkcí
cloudFiles.cleanSource: k dispozici jsouarchive_time,archive_modeamove_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
_metadataZachyťte alespoňfile_pathafile_modification_time. Podívejte se na sloupec metadat souboru . - Povolte sloupce
_rescued_dataa_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_dataa_corrupt_record), prvekcorrupt_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 NULLdetekuje neočekávané změny schématu a_corrupt_record IS NULLdetekuje 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 zsource.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,inputRowsPerSecondaprocessedRowsPerSecondz 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_timeacommit_timezcloud_files_state()pro latenci typu end-to-end. Ke zjištění latence zpracování použijtedurationMsčlenění (napříkladlatestOffset,addBatcha 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í vStreamingQueryListenerprobíhajících událostech v částiobservedMetrics.
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í aevent_type = 'flow_progress'událostí. - Vyvíjejte statistiky tabulky pomocí počtu řádků a objemu dat na tabulku odvozenou z
num_output_rowsprotokolu 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ýmdata_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_dataindikují odchylku schématu. Vyhledejte v protokolu událostí položkufailed_records > 0pro očekáváníno rescued data. -
_schemasZměny adresáře uvnitř nakonfigurovanéhocloudFiles.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í
onQueryTerminatednásledovanouonQueryStartedstejný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 —_schemaszměnami v adresáři nebo_rescued_dataporušeními očekávání — než dospějete k závěru, že došlo k evoluci schématu. - Slouží
_metadata.file_pathk identifikaci souborů, které zavedly změny schématu. Spojte ho scloud_files_state()polem a proveďte korelaci změn schématupaths 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.