Monitorowanie i obserwowanie automatycznego modułu ładującego

Potoki Auto Loadera wymagają aktywnego monitorowania, aby wykrywać problemy, takie jak narastające zaległości, zmiany schematu, uszkodzone dane i zatrzymane strumienie, zanim wpłyną one na odbiorców danych znajdujących się dalej w łańcuchu przetwarzania. Na tej stronie opisano, jak monitorować kluczowe metryki, sprawdzać stan na poziomie pliku, tworzyć pulpity monitorowania i rozwiązywać typowe problemy.

Aby uzyskać szczegółowe informacje o konfiguracji produkcyjnej, zobacz Konfigurowanie automatycznego modułu ładującego dla obciążeń produkcyjnych. Aby uzyskać najlepsze rozwiązania dotyczące konfiguracji, zobacz Najlepsze rozwiązania dotyczące automatycznego modułu ładującego.

Wymagania wstępne

Kilka procesów monitorowania na tej stronie opiera się na cloud_files_state(), aby obserwować stan przetwarzania poszczególnych plików — w tym zapytania dotyczące zaległości, obliczenia opóźnień i wykrywanie dryfu schematu. cloud_files_state() to funkcja tabelaryczna, która zwraca stan pozyskiwania danych na poziomie pliku dla punktu kontrolnego Auto Loader. Nie wszystkie pola są domyślnie dostępne. Dostępność zależy od wersji i konfiguracji środowiska Databricks Runtime:

  • Databricks Runtime 18.2 i nowsze: discovery_time, processed_time i commit_time są dostępne automatycznie. W środowisku Databricks Runtime 16.4–18.1 te pola są dostępne tylko wtedy, gdy cloudFiles.cleanSource jest włączona.
  • Databricks Runtime 16.4 i nowsze z włączonym cloudFiles.cleanSource: archive_time, archive_mode i move_location są dostępne.

Włączenie cloudFiles.cleanSource ma pewne obciążenie związane z wydajnością. Przed włączeniem tego w środowisku produkcyjnym przetestuj to na swoich obciążeniach w środowisku przedprodukcyjnym.

Additionally:

  • Dodawanie adnotacji do pozyskanych danych za pomocą kolumny _metadata . Przechwyć co najmniej file_path i file_modification_time. Zobacz Kolumna metadanych pliku.
  • Włącz _rescued_data i _corrupt_record kolumny.

Kluczowe metryki modułu ładującego automatycznego

Poniższa tabela podsumowuje najważniejsze metryki, które należy monitorować w potokach Auto Loader. Te metryki są dostępne w zdarzeniach postępu StreamingQueryListener, a wartości specyficzne dla Auto Loader są udostępniane w mapie metrics każdego źródła.

Metric Co to ci mówi
numFilesOutstanding Liczba plików na liście prac oczekujących na przetworzenie
numBytesOutstanding Rozmiar kolejki zaległych plików w bajtach
approximateQueueSize Głębokość kolejki chmurowej (tylko tryb powiadomień o plikach)
numInputRows Wiersze przetwarzane w partii
inputRowsPerSecond Szybkość przylotu danych
processedRowsPerSecond Przepustowość przetwarzania
durationMs Zestawienie Na co przeznaczany jest czas w każdej partii

Co należy obejrzeć

Poniższe wzorce wskazują, że potok może wymagać uwagi.

  • Rośnie numFilesOutstanding: Zaległości się tworzą. Potok przetwarzania nie nadąża za napływającymi danymi.
  • processedRowsPerSecond < inputRowsPerSecond: potok przetwarza dane wolniej niż docierają.
  • Wysoki durationMs.latestOffset: Wyszukiwanie plików jest powolne. Rozważ przejście na zdarzenia plikowe.
  • Duży durationMs.addBatch: przetwarzanie danych działa wolno. Rozważ skalowanie zasobów obliczeniowych lub optymalizowanie przekształceń.

Pełny opis metryk znajdziesz w sekcji Metryki źródła Auto Loader.

Wykonywanie zapytań o stan na poziomie pliku za pomocą polecenia cloud_files_state

Funkcja cloud_files_state() z wartością tabeli zawiera szczegółowe informacje o każdym pliku odnalezionym przez moduł automatycznego ładowania. Dostępne są następujące pola. Pola oznaczone jako wymagające środowiska Databricks Runtime 16.4 lub nowszego lub 18.2 lub nowszego są wypełniane tylko zgodnie z warunkami opisanymi w sekcji Wymagania wstępne.

Pole Typ Opis
path STRING Ścieżka pliku
size BIGINT Rozmiar pliku w bajtach
create_time TIMESTAMP Kiedy plik został utworzony
discovery_time TIMESTAMP Gdy Auto Loader wykrył plik (Databricks Runtime 16.4 i nowsze)
processed_time TIMESTAMP Gdy Auto Loader przetworzy plik (Databricks Runtime 16.4 i nowsze)
commit_time TIMESTAMP Gdy plik został zapisany w punkcie kontrolnym (Databricks Runtime 16.4 i nowsze)
archive_time TIMESTAMP Gdy plik został zarchiwizowany (wymaga cloudFiles.cleanSource)
archive_mode STRING MOVE, DELETElub NULL (wymaga cloudFiles.cleanSource)
move_location STRING Ścieżka docelowa, gdy cloudFiles.cleanSource jest MOVE
ingestion_state STRING Bieżący stan importowania plików

Sprawdź stan ingestii plików

Poniższe zapytania obejmują typowe scenariusze diagnostyczne.

Znajdź wszystkie nieprzetworzone pliki (bieżąca lista prac):

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

Średnie opóźnienie pozyskiwania zasobów obliczeniowych (czas od utworzenia pliku do zatwierdzenia):

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;

Znajdź uszkodzone lub pominięte pliki:

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

Śledź postęp archiwizacji (wymaga cloudFiles.cleanSource):

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

Znajdowanie plików z dużym opóźnieniem odnajdywania do zatwierdzenia w celu zidentyfikowania wąskich gardeł:

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;

Aby uzyskać pełną dokumentację języka SQL, zobacz cloud_files_state funkcja z wartościami tabeli.

Monitoruj Auto Loader w potokach Lakeflow

Databricks zaleca używanie potoków Lakeflow dla potoków Auto Loader używanych w środowisku produkcyjnym. Aby skorzystać z wbudowanych funkcji monitorowania:

  • Zapisz dziennik zdarzeń potoków Lakeflow w tabeli Delta, aby można było wysyłać do niego zapytania o dane monitorowania. Skonfiguruj to za pomocą ustawień zaawansowanych potoku lub interfejsu API. Aby uzyskać szczegółowe informacje, zobacz Dziennik zdarzeń potoku.

  • Zaprojektuj potok przetwarzania pod kątem obserwowalności. Dobrze ustrukturyzowany potok Auto Loader w potokach Lakeflow obejmuje widok {table}_source (definicję źródła Auto Loader), tabelę strumieniową {table}_bronze (surowe pozyskiwanie danych z kolumnami _rescued_data i _corrupt_record), element corrupt_records_sink, który poddaje kwarantannie wiersze z danymi, których nie można przeanalizować, oraz czysty widok {table} do wykorzystania w dalszych etapach przetwarzania.

  • Skonfiguruj oczekiwania dla tabel strumieniowych warstwy brązowej, aby monitorować dryf schematu i uszkodzenia danych. _rescued_data IS NULL wykrywa nieoczekiwane zmiany schematu i _corrupt_record IS NULL wykrywa nieparzysalne dane. Potoki Lakeflow weryfikują te oczekiwania w miarę napływu danych i generują ślad monitorowania. Możesz skonfigurować oczekiwania dotyczące ostrzegania, porzucania wierszy lub niepowodzenia potoku.

Po utworzeniu widoku event_log_raw dla potoku użyj następujących zapytań dotyczących metryk specyficznych dla funkcji Auto Loader.

Monitoruj przepływność pozyskiwania danych dla każdego przepływu:

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;

Monitoruj backlog danych dla każdego przepływu:

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;

Podsumowanie naruszeń oczekiwań w celu wykrycia dryfu schematu i uszkodzonych danych:

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;

Ogólne wskazówki dotyczące monitorowania potoków usługi Lakeflow można znaleźć w sekcjach Monitorowanie potoków i Dziennik zdarzeń potoku.

Monitoruj Auto Loader za pomocą Structured Streaming

Podczas uruchamiania Auto Loader poza potokami Lakeflow stosuj następujące metody monitorowania Structured Streaming.

  • Zaimplementuj StreamingQueryListener, aby rejestrować metryki specyficzne dla Auto Loader z każdej partii poprzez odczyt 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())

Uwaga / Notatka

Logika przetwarzania w odbiornikach może spowolnić przetwarzanie zapytań. Ogranicz obliczenia w funkcjach zwrotnych odbiornika i unikaj synchronicznych zapisów do systemów zewnętrznych w tym miejscu; zamiast tego emituj lekką telemetrię asynchronicznie lub przekazuj metryki do oddzielnego zadania w celu trwałego zapisania.

  • Użyj numInputRows, inputRowsPerSecond i processedRowsPerSecond z postępu źródła, aby obliczyć przepustowość — pliki na sekundę i wiersze na sekundę dla każdej partii.

  • Aby obliczyć opóźnienie przyjmowania danych, porównaj create_time i commit_time w cloud_files_state(), aby obliczyć opóźnienie end-to-end. W przypadku opóźnień przetwarzania użyj podziału durationMs (na przykład latestOffset, addBatch i innych raportowanych etapów przetwarzania wsadowego), aby zidentyfikować, który etap jest wąskim gardłem.

  • Służy df.observe() do definiowania wbudowanych metryk jakości danych bezpośrednio w ramce danych przesyłania strumieniowego. Metryki są widoczne w zdarzeniach postępu StreamingQueryListener w sekcji 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")
)
  • Użyj .queryName(), aby przypisać unikalną nazwę każdemu strumieniowi, co ułatwi odróżnianie strumieni Auto Loader na karcie Streaming w interfejsie Spark UI oraz na panelach monitorowania.

Pełne informacje referencyjne dotyczące monitorowania Structured Streaming można znaleźć w artykule Monitorowanie zapytań Structured Streaming w usłudze Azure Databricks.

Utwórz panel obserwowalności

Łączenie danych z wielu źródeł w celu utworzenia kompleksowego pulpitu nawigacyjnego z obserwacją potoków modułu automatycznego ładowania. W tej tabeli przedstawiono sugerowane źródła, których możesz użyć do utworzenia pulpitu obserwowalności.

Źródło danych Dane dotyczące obserwacji
cloud_files_state() Stan importu dla pliku: znaczniki czasu wykrywania, przetwarzania, zatwierdzenia i archiwizacji dla każdego pliku
Dziennik zdarzeń potoków Lakeflow Historia przebiegów potoku, metryki przepływu dla partii i wyniki oczekiwań dotyczących jakości danych
Tabele wyjściowe pipeline’u Liczba wierszy i wolumen danych zapisywane dla każdej zaimportowanej tabeli

Następnie można agregować dane dotyczące obserwacji w dedykowanych tabelach, które służą jako podstawa dla pulpitów nawigacyjnych i alertów:

  • Podsumuj stany przebiegu potoku (powodzenie lub niepowodzenie) w czasie pochodzące z event_type = 'update_progress' zdarzeń.
  • Zagregowane metryki importu plików (rozmiar zaległości, przepustowość, opóźnienie dla każdej partii), na podstawie zdarzeń cloud_files_state() i event_type = 'flow_progress'.
  • Generuj statystyki tabel na podstawie liczby wierszy i ilości danych dla każdej tabeli, uzyskanych z num_output_rows w dzienniku zdarzeń.
  • Zbierz informacje diagnostyczne ze szczegółowych dzienników błędów i naruszeń oczekiwań dla każdej aktualizacji, na podstawie zdarzeń event_type = 'flow_progress' z wypełnionym polem data_quality.

Te zagregowane tabele mogą zasilać pulpit nawigacyjny sztucznej inteligencji/analizy biznesowej i alerty SQL. Zalecane panele na pulpicie nawigacyjnym obejmują oś czasu stanu uruchomień potoku, trend zaległości przetwarzania danych, trend przepływności, rozkład opóźnień przetwarzania danych, metryki jakości danych, zdarzenia ewolucji schematu oraz status archiwizacji plików.

Monitorowanie zdarzeń ewolucji schematu

Użyj poniższych metod wykrywania zmian schematu w miarę ich występowania.

  • Wartości inne niż NULL w _rescued_data liczbie naruszeń oczekiwań wskazują dryf schematu. Przeszukaj dziennik zdarzeń dla failed_records > 0 pod kątem oczekiwania no rescued data.
  • Zmiany w _schemas katalogu wewnątrz skonfigurowanego cloudFiles.schemaLocation (lub wewnątrz punktu kontrolnego tylko wtedy, gdy lokalizacja schematu nie jest ustawiona oddzielnie) wskazują, że nastąpiła ewolucja schematu. Ten katalog można odpytywać z poziomu oddzielnego zadania monitorującego.
  • Nie traktuj zdarzenia onQueryTerminated, po którym następuje onQueryStarted dla tej samej nazwy strumienia, jako wystarczającego dowodu ewolucji schematu samego w sobie. Strumienie są restartowane z wielu powodów (ponowne uruchomienia klastra, wdrożenia kodu, tymczasowe błędy pamięci masowej). Powiąż ponowne uruchomienia z niezależnymi sygnałami — zmianami katalogu _schemas lub naruszeniami oczekiwań _rescued_data — zanim uznasz, że doszło do ewolucji schematu.
  • Służy _metadata.file_path do identyfikowania plików, które wprowadziły zmiany schematu. Połącz to z cloud_files_state() na polu path, aby skorelować zmiany schematu z konkretnymi plikami i partiami.

Użyj tego przykładowego zapytania, aby wykryć niedawny dryf schematu za pośrednictwem naruszeń oczekiwań:

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;

Konfigurowanie alertów dla typowych problemów

Użyj alertów SQL usługi Databricks lub powiadomień potoków, aby wykrywać problemy, zanim wpłyną one na odbiorców zależnych.

Poniższy kod SQL wykrywa rosnącą listę prac i może służyć jako podstawa alertu SQL usługi Databricks. Zaplanuj jego cykliczne uruchamianie (na przykład co 5 minut) i generuj alert, gdy wynik nie jest pusty.

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

Poniższa tabela zawiera podsumowanie zalecanych warunków alertów:

Co należy wykryć Jak go wykryć Kiedy wysyłać alerty
Rosnące zaległości numFilesOutstanding trend wzrostowy Utrzymujący się wzrost w kolejnych partiach
Utknięty strumień Brak zdarzeń postępu Brak zdarzeń przez N minut (na podstawie oczekiwanego interwału wyzwalacza)
Wysokie opóźnienie przyjmowania danych commit_time - create_time Przekracza próg umowy SLA
Obniżenie jakości danych Wskaźnik niespełnienia oczekiwań Rosnący odsetek wierszy niespełniających oczekiwań
Zdarzenie ewolucji schematu _rescued_data IS NOT NULL Wszystkie wartości inne niż NULL w liczbie naruszeń oczekiwań
Powolne odnajdywanie plików durationMs.latestOffset Znacznie wyższy niż punkt odniesienia

Rozwiązywanie typowych problemów

W poniższej tabeli opisano typowe problemy z potokiem Auto Loader, ich prawdopodobne przyczyny oraz zalecane działania służące do ich rozwiązania.

Issue Możliwa przyczyna Zalecana akcja
Zaległości rosną szybciej niż przetwarzanie Zbyt mała moc obliczeniowa, nierównomierny rozkład danych lub ograniczone limity szybkości Skalowanie zasobów obliczeniowych, sprawdzanie niesymetryczności za pomocą interfejsu użytkownika platformy Spark i przeglądanie maxFilesPerTrigger ustawień w celu kontrolowania rozmiaru partii
Nie odnaleziono plików Błędnie skonfigurowane zdarzenia związane z plikami, problem z uprawnieniami lub strumień nie został uruchomiony w ciągu 7 dni Sprawdź uprawnienia do lokalizacji zewnętrznej, sprawdź konfigurację zdarzeń plików w interfejsie użytkownika Unity Catalog i upewnij się, że strumień jest uruchamiany co najmniej raz na 7 dni, aby uniknąć wygaśnięcia stanu RocksDB
Uruchamianie usługi Stream trwa zbyt długo Pobieranie stanu dużego punktu kontrolnego (RocksDB) Uaktualnienie do środowiska Databricks Runtime w wersji 15.3 lub nowszej w celu załadowania stanu asynchronicznego, co skraca czas uruchamiania o ok. 90%
Zduplikowane przetwarzanie plików Agresywne cloudFiles.maxFileAge ustawienia lub uszkodzenie punktu kontrolnego Użyj konserwatywnego maxFileAge (co najmniej 90 dni), zweryfikuj integralność punktu kontrolnego i unikaj zasad zarządzania cyklem życia w magazynie punktów kontrolnych
Ewolucja schematu powodująca ponowne uruchamianie potoku Częste lub niezgodne zmiany schematu Przejrzyj schemaEvolutionMode, przejdź na addNewColumnsWithTypeWidening w przypadku promowania typów lub użyj typu Variant dla wysoce dynamicznych schematów
Uszkodzone dane gromadzące się w ujściu Problemy z jakością danych źródłowych Sprawdź repozytorium kwarantanny _corrupt_record pod kątem wzorców, przejrzyj generowanie danych źródłowych i rozważ dodanie walidacji na wcześniejszym etapie
discovery_time i commit_time nie są wypełnione Uruchamianie w środowisku Databricks Runtime poniżej wersji 18.2 bez cleanSource Uaktualnij do Databricks Runtime 18.2 lub nowszego albo włącz cloudFiles.cleanSource w Databricks Runtime 16.4–18.1

Aby uzyskać dodatkowe informacje dotyczące rozwiązywania problemów, zobacz Automatyczne ładowanie — często zadawane pytania.