Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
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_timeicommit_timesą dostępne automatycznie. W środowisku Databricks Runtime 16.4–18.1 te pola są dostępne tylko wtedy, gdycloudFiles.cleanSourcejest włączona. -
Databricks Runtime 16.4 i nowsze z włączonym
cloudFiles.cleanSource:archive_time,archive_modeimove_locationsą 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 najmniejfile_pathifile_modification_time. Zobacz Kolumna metadanych pliku. - Włącz
_rescued_datai_corrupt_recordkolumny.
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_datai_corrupt_record), elementcorrupt_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 NULLwykrywa nieoczekiwane zmiany schematu i_corrupt_record IS NULLwykrywa 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 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())
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,inputRowsPerSecondiprocessedRowsPerSecondz 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_timeicommit_timewcloud_files_state(), aby obliczyć opóźnienie end-to-end. W przypadku opóźnień przetwarzania użyj podziałudurationMs(na przykładlatestOffset,addBatchi 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ępuStreamingQueryListenerw sekcjiobservedMetrics.
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()ievent_type = 'flow_progress'. - Generuj statystyki tabel na podstawie liczby wierszy i ilości danych dla każdej tabeli, uzyskanych z
num_output_rowsw 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 polemdata_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_dataliczbie naruszeń oczekiwań wskazują dryf schematu. Przeszukaj dziennik zdarzeń dlafailed_records > 0pod kątem oczekiwaniano rescued data. - Zmiany w
_schemaskatalogu wewnątrz skonfigurowanegocloudFiles.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ępujeonQueryStarteddla 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_schemaslub naruszeniami oczekiwań_rescued_data— zanim uznasz, że doszło do ewolucji schematu. - Służy
_metadata.file_pathdo identyfikowania plików, które wprowadziły zmiany schematu. Połącz to zcloud_files_state()na polupath, 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.