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.
Uruchamiaj produkcyjne obciążenia Structured Streaming jako zaplanowane zadania Lakeflow w usłudze Azure Databricks. Zobacz Zadania lakeflow.
Usługa Databricks zaleca, aby zawsze konfigurować następujące elementy:
- Usuń niepotrzebny kod z notatników, który może zwrócić wyniki, takie jak
displayicount. - Nie uruchamiaj obciążeń Structured Streaming w środowisku obliczeniowym ogólnego przeznaczenia. Zawsze planuj strumienie jako zadania lakeflow przy użyciu obliczeń zadań.
- Planowanie zadań Lakeflow w trybie
Continuousmode. Dotyczy to funkcji planowania zadań w Azure Databricks, a nie funkcji Structured Streaming dotyczącej interwału wyzwalania . - Nie włączaj automatycznego skalowania mocy obliczeniowej dla zadań Structured Streaming.
Niektóre obciążenia korzystają z następujących elementów:
- Konfigurowanie magazynu stanów bazy danych RocksDB w Azure Databricks
- Asynchroniczne sprawdzanie stanu dla zapytań stanowych
- Śledzenie postępu asynchronicznego
Usługa Databricks wprowadziła potoki Lakeflow w celu zmniejszenia złożoności zarządzania infrastrukturą produkcyjną dla obciążeń przesyłania strumieniowego ze strukturą. Databricks zaleca używanie potoków Lakeflow do tworzenia nowych potoków Structured Streaming. Zobacz Potoki deklaratywne platformy Spark.
Uwaga
Automatyczne skalowanie zasobów obliczeniowych ma ograniczenia dotyczące zmniejszania rozmiaru klastra dla obciążeń przetwarzania strumieniowego ze zdefiniowaną strukturą. Usługa Databricks zaleca używanie potoków deklaratywnych platformy Spark w usłudze Lakeflow z rozszerzonym skalowaniem automatycznym na potrzeby obciążeń przesyłania strumieniowego. Zobacz Optymalizowanie wykorzystania klastra pipeline’u Lakeflow za pomocą automatycznego skalowania.
:::uwaga Bezserwerowe obliczenia
W przypadku obliczeń bezserwerowych obsługiwane są tylko funkcje Trigger.AvailableNow() i Trigger.Once() . Databricks zaleca Trigger.AvailableNow().
W przypadku ciągłego przesyłania strumieniowego na obliczeniach bezserwerowych użyj trybu wyzwalanego i ciągłego potoku w trybie ciągłym.
Zobacz Ograniczenia przesyłania strumieniowego.
:::
Zmniejszenie opóźnień podczas operacyjnego streamingu
Operacyjne obciążenia strumieniowe pobierają, transformują i przetwarzają dane niemal w czasie rzeczywistym. Typowe przykłady to wykrywanie oszustw, anomalii, personalizacja oraz monitorowanie i powiadamianie w czasie rzeczywistym, gdzie opóźnienie w przetwarzaniu bezpośrednio wpływa na wyniki biznesowe. Niskie opóźnienia dla tych obciążeń zazwyczaj oznaczają dziesiątki do setek milisekund, choć wiele zespołów ustala umowy o poziomie usług (SLA) w zakresie sekund, aby uwzględnić zmienność w wyższych percentylach.
Aby uzyskać najniższe całkowite opóźnienie, należy użyć trybu czasu rzeczywistego, który zapewnia całkowite opóźnienie poniżej jednej sekundy w skrajnym przypadku i około 300 milisekund w typowych przypadkach. Zobacz koncepcje trybu czasu rzeczywistego.
Gdy tryb czasu rzeczywistego nie pasuje do Twojego obciążenia, poniższe najlepsze praktyki zmniejszają opóźnienia dla mikro-batchowego strumieniowania strukturalnego:
- Tryb wyjściowy: Użyj trybu aktualizacji, gdzie operatory zapytań i sink go wspierają. Tryb aktualizacji generuje zaktualizowane wiersze po każdym wyzwalaczu i aktualizuje je aż do wygaśnięcia znaku wodnego, więc spraw, by Twój downstream sink był idempotent do obsługi zaktualizowanych wyników. Używaj trybu dopisywania dla obciążeń, których nie obsługuje tryb aktualizacji, takich jak złączenia strumień–strumień, lub gdy możesz odrzucać dane napływające z opóźnieniem. Nie używaj trybu kompletnego dla niskich opóźnień. Zobacz Wybieranie trybu danych wyjściowych dla przesyłania strumieniowego ze strukturą.
-
Wyzwalacz: Użyj wyzwalacza
processingTimez interwałem0, który uruchamia następną mikropartię zaraz po zakończeniu poprzedniej i pojawieniu się nowych danych. Zapewnia to najniższe opóźnienie dla mikropartii, ale zwiększa koszty interfejsu API usługi przechowywania w chmurze. Nie używajAvailableNow,Once, aniContinuousdo zadań operacyjnych. Zobacz Konfigurowanie interwałów wyzwalacza strukturalnego przesyłania strumieniowego. - Znak wodny: Ustaw znak wodny na tyle długi, aby uwzględnić dane o późnym nadejściu, dzięki czemu Twoje obciążenie nie może spadnąć. Znacznik wodny określa, jak długo zapytanie przyjmuje dostarczane poza kolejnością dane czasu zdarzenia, zanim je odrzuci i usunie stan, więc zbyt krótki znacznik wodny po cichu odrzuca prawidłowe spóźnione rekordy. W ramach tego ograniczenia krótszy watermark obniża opóźnienie i wymaga utrzymywania mniejszego stanu, a dłuższy watermark toleruje więcej spóźnionych danych kosztem większego opóźnienia i większego stanu. Niewielka wielokrotność limitu opóźnienia określonego w SLA, na przykład 2x, to rozsądny punkt wyjścia do dostrajania. Zobacz Stosowanie wodnych znaków do kontroli progów przetwarzania danych.
-
Źródła i ujścia: Odczyt z źródeł o niskich opóźnieniach, takich jak szyny komunikatów (Apache Kafka, Amazon Kinesis, Apache Pulsar lub Google Cloud Pub/Sub) lub strumienie danych o zmianach z tabel Delta Lake i Apache Iceberg. Zapisuj do miejsc docelowych o niskich opóźnieniach i wysokiej przepustowości, takich jak szyny komunikatów, operacyjne bazy danych lub miejsca docelowe
foreach. Projektuj operacje zapisu tak, aby były idempotentne, dzięki czemu odbiorcy downstream będą mogli obsługiwać duplikaty i dane docierające z opóźnieniem. - Stan i tworzenie punktów kontrolnych: W przypadku zapytań stanowych używaj magazynu stanu RocksDB, który jest wymagany zarówno do tworzenia punktów kontrolnych dziennika zmian, jak i asynchronicznego tworzenia punktów kontrolnych stanu. Włącz tworzenie punktów kontrolnych dziennika zmian, aby utrwalać tylko przyrostowe zmiany stanu. Gdy zapisywanie punktów kontrolnych stanu stanowi wąskie gardło w czasie przetwarzania partii, włącz asynchroniczne zapisywanie punktów kontrolnych stanu, aby zapisy punktów kontrolnych nakładały się na przetwarzanie kolejnej mikropartii, po wcześniejszym przeanalizowaniu zastrzeżeń dotyczących odzyskiwania po awarii i skalowania klastra. Przypisz każdemu zapytaniu własny katalog punktów kontrolnych w trwałej pamięci masowej w chmurze. Zobacz Konfigurowanie magazynu stanu RocksDB w usłudze Azure Databricks, Asynchroniczne tworzenie punktów kontrolnych stanu dla zapytań stanowych oraz Punkty kontrolne w Structured Streaming.
-
Zarządzanie offsetami: Aby zmniejszyć opóźnienia wynikające z zapisywania punktów kontrolnych offsetów w strumieniach ciągłych, włącz asynchroniczne śledzenie postępu, które aktualizuje offsety i dzienniki zatwierdzeń bez blokowania przetwarzania danych. Nie jest kompatybilne z wyzwalaczami
AvailableNowlubOnce. Zobacz Śledzenie postępu asynchronicznego. - Skoki pamięci: Utrzymuj obliczenia w jednym strumieniu strumieniowym, jeśli to możliwe. Rozdzielanie logiki między wiele zadań lub potoków przetwarzania dodaje dodatkowe etapy zapisu i odczytu danych, co zwiększa opóźnienia.
Projektuj obciążenia strumieniowe, aby uwzględniały możliwość awarii
Databricks zaleca, aby zawsze konfigurować zadania strumieniowe tak, aby były automatycznie uruchamiane ponownie po awarii. Niektóre funkcje, w tym ewolucja schematu, wymagają, aby obciążenia Structured Streaming były automatycznie ponawiane. Zobacz Konfigurowanie zadań strumieniowania strukturalnego do ponownego uruchamiania zapytań strumieniowych w przypadku niepowodzenia.
Niektóre operacje, takie jak foreachBatch, zapewniają gwarancje co najmniej jednokrotne zamiast gwarancji dokładnie jednokrotnych. W przypadku tych operacji upewnij się, że potok przetwarzania jest idempotentny. Zobacz Używanie polecenia foreachBatch do zapisu w dowolnych odbiornikach danych.
Uwaga
Po ponownym uruchomieniu zapytania mikropartia przetwarzana podczas poprzedniego uruchomienia. Jeśli twoje zadanie nie powiodło się z powodu błędu braku pamięci lub ręcznie anulowałeś je z powodu zbyt dużego mikrosadowego przetwarzania, może być konieczne zwiększenie mocy obliczeniowej, aby pomyślnie przetworzyć mikrosadę.
Jeśli zmienisz konfiguracje między przebiegami, te konfiguracje zostaną zastosowane do pierwszej nowej partii zaplanowanej. Zobacz Odzyskiwanie po zmianach w zapytaniu Strukturowanego przesyłania strumieniowego.
Gdy zadanie jest ponawiane
W ramach zadania Azure Databricks można zaplanować wiele zadań. Podczas konfigurowania zadania przy użyciu wyzwalacza ciągłego nie można ustawić zależności między zadaniami.
Możesz zdecydować się na zaplanowanie wielu strumieni w jednym zadaniu, korzystając z jednego z poniższych podejść.
- Wiele zadań: Zdefiniuj zadanie obejmujące wiele zadań, które przetwarzają obciążenia strumieniowe przy użyciu wyzwalacza ciągłego.
- Wiele zapytań: zdefiniuj wiele zapytań przesyłanych strumieniowo w kodzie źródłowym dla jednego zadania.
Można również połączyć te strategie. W poniższej tabeli porównaliśmy te podejścia.
| Strategia | Wiele zadań | Wiele zapytań |
|---|---|---|
| Jak współużytkowane są zasoby obliczeniowe? | Databricks zaleca wdrożenie zasobów obliczeniowych o odpowiednim rozmiarze do każdego zadania przesyłania strumieniowego. Opcjonalnie możesz udostępniać zasoby obliczeniowe między zadaniami. | Wszystkie zapytania współdzielą te same obliczenia. Opcjonalnie można przypisywać zapytania do pul harmonogramu. |
| Jak są obsługiwane ponawianie prób? | Wszystkie zadania muszą zakończyć się niepowodzeniem, zanim praca zostanie ponownie podjęta. | Zadanie ponawia próbę, jeśli jakiekolwiek zapytanie zakończy się niepowodzeniem. |
Aby uzyskać więcej informacji na temat pracy z wieloma zadaniami lub zapytaniami, zobacz Uruchamianie wielu zapytań przesyłania strumieniowego ze strukturą w tym samym klastrze.
Konfigurowanie zadań przesyłania strumieniowego ze strukturą w celu ponownego uruchamiania zapytań przesyłanych strumieniowo w przypadku niepowodzenia
Usługa Databricks zaleca skonfigurowanie wszystkich obciążeń przesyłania strumieniowego przy użyciu wyzwalacza ciągłego. Zobacz Uruchamianie zadań w sposób ciągły.
Wyzwalacz ciągły domyślnie ma następujące zachowanie:
- Powstrzymuje więcej niż jednoczesne uruchomienie zadania.
- Uruchamia nowy przebieg, gdy poprzedni przebieg zakończy się niepowodzeniem.
- Używa wykładniczego odstępu dla ponownych prób.
Databricks zaleca zawsze używanie zasobów obliczeniowych dla zadań zamiast ogólnych zasobów obliczeniowych podczas planowania przepływów pracy. W przypadku niepowodzenia zadania i ponawiania próby nowe zasoby obliczeniowe są wdrażane.
Uwaga
Usługa Databricks zaleca, aby nie używać streamingQuery.awaitTermination() ani spark.streams.awaitAnyTermination(). Zobacz Kiedy używać awaitTermination().
Kiedy należy używać awaitTermination()
streamingQuery.awaitTermination() i spark.streams.awaitAnyTermination() blokują bieżący wątek do momentu zakończenia przesyłanego strumieniowo zapytania. To, czy używać tych funkcji, zależy od środowiska wykonawczego.
W przypadku zadań Lakeflow nie należy używać streamingQuery.awaitTermination() ani spark.streams.awaitAnyTermination(). Te funkcje nie są niezbędne, ponieważ usługa Zadań automatycznie uniemożliwia zakończenie przebiegu, gdy zapytanie strumieniowe jest aktywne. Obie funkcje blokują wykonanie komórek notatnika i uniemożliwiają usłudze Jobs śledzenie zapytania przesyłania strumieniowego, co zakłóca metryki zadań w kolejce i powiadomienia o zadaniach.
Użyj awaitTermination() w następujących przypadkach:
| Przypadek użycia | Zachowanie |
|---|---|
| Interaktywne notebooki w uniwersalnych obliczeniach |
awaitTermination() utrzymuje działanie komórki, pozwala monitorować stan zapytania i zapewnia, że błędy są widoczne w danych wyjściowych notesu. |
| Środowiska lokalne i programistyczne | Podczas lokalnego uruchamiania programu Spark proces kończy się po zakończeniu głównego wątku. Wywołaj awaitTermination() aby utrzymać działanie programu aż do zakończenia lub niepowodzenia zapytania strumieniowego. |
| Propagacja błędu do sterownika | Bez awaitTermination() parametru niepowodzenie zapytania strumieniowego w kontekście niezwiązanym bezpośrednio z zadaniem może nie być propagowane do wątku wywołującego. Zapytanie może zakończyć się niepowodzeniem w trybie dyskretnym, co utrudnia wykrywanie i diagnozowanie błędów. Wywołanie awaitTermination() ponownie powoduje zgłoszenie wyjątku zapytania w sterowniku. |