Podstawy platformy Spark

Podstawowe pojęcia stanowiące podstawę określania rozmiaru, optymalizacji i rozwiązywania problemów. Przeczytaj to najpierw, jeśli dopiero zaczynasz korzystać z platformy Spark w Fabric.

Ogólne zasady i zakazy

Scenariusz: dopiero zaczynasz korzystać z platformy Spark. Co robić, a czego unikać
Przypadek użycia Najlepsze rozwiązania
Używanie zoptymalizowanych formatów serializowanych Preferuj formaty, takie jak Avro, Parquet lub Optimized Row Columnar (ORC), ponieważ zawierają schemat, są kompaktowe i optymalizują przechowywanie i przetwarzanie. W Fabric użyj formatu Delta w celu zapewnienia niepodzielności, spójności, izolacji, trwałości (ACID) i zalet związanych z wydajnością.
Zachowaj ostrożność przy użyciu formatu XML/JSON Nie polegaj na wnioskowaniu schematu dla dużych plików JSON (JavaScript Object Notation) ani Extensible Markup Language (XML), ponieważ platforma Spark odczytuje cały zestaw danych w celu wnioskowania schematu, co spowalnia przetwarzanie i intensywnie zużywa pamięć.

Podaj statyczny schemat podstawowy podczas odczytywania kodu JSON/XML lub używania .option("samplingRatio", 0.1) go do przyspieszenia operacji odczytu, ale należy pamiętać, że jeśli przykład nie reprezentuje pełnego zestawu danych, operacje odczytu mogą zakończyć się niepowodzeniem. Bezpieczne podejście polega na wywnioskowaniu schematu z reprezentatywnej próbki i jego utrwaleniu dla wszystkich operacji odczytu.

Unikaj analizowania dużych plików XML. Analizowanie XML działa z natury wolniej z powodu przetwarzania tagów i rzutowania typów.
Optymalizowanie sprzężeń i filtrowania Do: Zastosuj przycinanie kolumn i filtrowanie na poziomie wiersza przed połączeniami, aby zmniejszyć przetwarzanie danych i użycie pamięci.

Optimizator Catalyst automatycznie obsługuje przenoszenie predykatów podczas korzystania z API DataFrame. Unikaj używania interfejsów API odpornego rozproszonego zestawu danych (RDD), gdyż omijają optymalizacje prowadzone przez Catalyst.
Preferuj ramki danych zamiast RDD Do: Używaj DataFrame'y zamiast RDD do większości operacji. Ramki danych korzystają z optymalizatora Catalyst oraz silnika wykonawczego Tungsten dla efektywnego wykonywania.
Włączanie adaptacyjnego wykonywania zapytań (AQE) Włącz AQE dla dynamicznej optymalizacji partycji shuffle i automatycznego obsługiwania niesymetrycznych danych.

Zarządzanie pamięcią funkcji wykonawczej

Scenariusz: chcesz zrozumieć zarządzanie pamięcią funkcji wykonawczej na potrzeby dostrajania wydajności.

Nawet jeśli funkcja wykonawcza jest skonfigurowana z pamięcią o rozmiarze 56 GB, platforma Spark nie zezwala na bezpośrednie używanie wszystkich tych funkcji na potrzeby danych użytkownika. Platforma Spark Core dzieli pamięć funkcji wykonawczej i zarządza nią:

  • Pamięć zarezerwowana: Stała część przeznaczona na wewnętrzne obciążenia systemowe i Spark (na przykład maszyna wirtualna Java (JVM), komponenty wewnętrzne).

  • Pamięć użytkownika: Przechowuje funkcje zdefiniowane przez użytkownika (UDF), zmienne lokalne, struktury danych (listy, mapy, słowniki) i obiekty utworzone podczas obliczeń.

  • Pamięć przechowywania: Przechowuje buforowane/utrwalane dane, zmienne rozgłaszania i dane przetasowania, które można buforować.

  • Pamięć wykonywania: Służy do obliczeń pośrednich (mieszania, sprzężeń, sortowania, agregacji).

  • Dynamiczne udostępnianie pamięci: Granica między magazynem a pamięcią wykonywania jest wymienna. Platforma Spark może pożyczyć pamięć z jednego regionu do drugiego, co pozwala na elastyczne użycie pamięci.

  • Przepełnienie: Występuje, gdy zapotrzebowanie na pamięć dyskową lub pamięć operacyjną przekracza dostępną pamięć po zapożyczeniu. Wymusza to na dysku dane, co może mieć wpływ na wydajność.

    Diagram zarządzania pamięcią i rozlewania w Sparku.

Błędy braku pamięci (OOM)

Scenariusz: Zadania Spark kończą się niepowodzeniem z powodu błędów przekroczenia dostępnej pamięci (OOM).

Sterownik OOM:

Błędy OOM sterownika Spark występują, gdy przekracza on przydzieloną pamięć.

Typowa przyczyna: operacje intensywne pod względem użycia sterowników, takie jak collect(), countByKey(), lub duże toPandas() wywołania, które pobierają zbyt dużo danych do pamięci sterownika.

Środki zaradcze: unikaj operacji intensywnie korzystających ze sterowników, jeśli jest to możliwe. Jeśli jest to nieuniknione, zwiększ rozmiar sterownika i test porównawczy, aby znaleźć optymalną konfigurację.

Wykonawca bez pamięci (OOM):

Błędy funkcji wykonawczej OOM występują, gdy funkcja wykonawcza platformy Spark przekracza przydzieloną pamięć.

Typowa przyczyna: przekształcenia intensywnie korzystające z pamięci i obliczeń w dużych zestawach danych (na przykład szerokie sprzężenia, agregacje, przetasowania) lub buforowane/utrwalane zestawy danych, które przekraczają dostępną pamięć wykonawcy (wykonywanie i regiony magazynu).

Środki zaradcze: w razie potrzeby zwiększ pamięć wykonawcy, dostosuj wartości konfiguracyjne pamięci Spark (spark.memory.fraction, spark.memory.storageFraction) i dokonaj selektywnego utrwalenia. Upewnij się, że buforowane dane mieszczą się w dostępnej pamięci.

Zniekształcenia danych

Objawy niesymetryczności:

  • Kilka zadań trwa dłużej niż inne w interfejsie użytkownika platformy Spark (zadania etapowe pokazują ciężki ogon).
  • Duża różnica między medianą a maksymalnym czasem zadań w metrykach etapu.
  • Fazy o dużych rozmiarach odczytu lub zapisu związanych z procesem shuffle dla kilku partycji.

Typowe przyczyny:

  • Nierównomierna dystrybucja danych dla kluczy sprzężenia i grupowania (gorące klucze).
  • Niepoprawne partycjonowanie lub zbyt mało partycji dla woluminu danych.
  • Anomalie danych nadrzędnych, które generują duże rekordy lub wiele kluczy o wartości null/pustych.

Łagodzenia:

  • Ponowne partycjonowanie lub łączenie w celu zwiększenia równoległości partycji i rozmiarów równowagi.
  • Zastosuj solenie kluczy lub partycjonowanie niestandardowe, aby rozłożyć intensywnie używane klucze między partycjami.
  • Użyj AQE (adaptacyjnego wykonywania zapytań), aby połączyć partycje po shuffle i włączyć optymalizacje sprzężenia niesymetrycznego.
  • Używaj sprzężeń rozgłaszających dla małych tabel wyszukiwania, aby uniknąć całkowitego przetasowywania.
  • Utrwalanie zrównoważonych zestawów danych pośrednich przed kosztownymi etapami i ponowne uruchomienie zadania.

Najlepsze rozwiązania dotyczące funkcji UDF

Scenariusz: należy zastosować logikę niestandardową, której nie można wyrazić za pomocą wbudowanych funkcji ramki danych.

Używaj API Spark DataFrame, kiedy to możliwe. Optymalizator Catalyst optymalizuje wbudowane funkcje i uruchamia je natywnie na maszynie JVM, dzięki czemu zapewniają najlepszą wydajność.

Jeśli musisz użyć funkcji zdefiniowanej przez użytkownika (UDF), unikaj zwykłych funkcji UDF języka Python PySpark. Zamiast tego należy wziąć pod uwagę następujące alternatywy:

  • UDF w Pandas (znane również jako zwektoryzowane UDF): wykorzystują Apache Arrow do efektywnego przesyłania danych między środowiskami JVM i Python. Funkcje Pandas UDF umożliwiają wektoryzowane operacje, co znacznie poprawia wydajność w porównaniu z funkcjami UDF w Pythonie przetwarzającymi dane pojedynczo dla każdego wiersza.

  • Scala/Java UDF-y: Uruchamiane bezpośrednio na maszynie JVM, unikając obciążeń związanych z serializacją w Pythonie. Funkcje zdefiniowane przez użytkownika (UDF) w językach Scala/Java zazwyczaj działają szybciej niż UDFy w języku Python.

Zachowaj ostrożność przy użyciu funkcji użytkownika w Pythonie. Każdy wykonawca uruchamia oddzielny proces języka Python, wymagający serializacji i deserializacji danych między maszyną wirtualną JVM a językiem Python. Powoduje to wąskie gardło wydajności, szczególnie na dużą skalę. 

Rejestrowanie błędów

Scenariusz: Najlepsze rozwiązania dotyczące logowania błędów w Fabric Spark
  1. Użyj log4j zamiast print(), co mocno obciąża kierowcę. Za pomocą log4j można mieć dostęp do dzienników sterowników i je przeszukiwać (na przykład za pomocą nazwy rejestratora: PySparkLogger).

    Diagram dzienników platformy Spark.

  2. Zawijaj odczyty, zapisy i przekształcenia w blokach try i except. Użyj logger.error dla wyjątków i logger.info dla komunikatów postępu.

    • Rejestrowanie języka Python: Idealne rozwiązanie do rejestrowania operacji, aktualizacji stanu lub debugowania informacji z kodu wykonywanego tylko na sterowniku Spark. Moduł logowania w Pythonie nie przekazuje zapisów do dzienników wykonawczych. Zapoznaj się z dokumentacją tworzenia, wykonywania notesów i zarządzania nimi.

    • Dziennik Spark4j: Standard niezawodnego rejestrowania aplikacji na poziomie produkcyjnym na platformie Spark, który integruje się natywnie z dziennikami sterowników/funkcji wykonawczej platformy Spark.

    Przykładowe użycie log4j w narzędziu PySpark:

    import traceback
    # Get log4j logger
    log4jLogger = spark._jvm.org.apache.log4j
    logger = log4jLogger.LogManager.getLogger("PySparkLogger")
    logger.info("Application started.")
    try:
        # Create DataFrame with 20 records
        data = [(f"Name{i}", i) for i in range(1, 21)]  # 20 records
        df = spark.createDataFrame(data, ["name", "age"])
        logger.info("DataFrame created successfully with 20 records.")
        df.show(s)  # 's' is not defined -> will throw error but the application will not fail
    except Exception as e:
        logger.error(f"Error while creating or showing DataFrame: {str(e)}\n{traceback.format_exc()}")
    
  3. Scentralizowane monitorowanie błędów:

    • Użyj rozszerzenia diagnostycznego emitera (Monitorowanie aplikacji Apache Spark za pomocą usługi Azure Log Analytics) w środowisku i dołącz do notesów z uruchomionymi aplikacjami Spark. Emiter może wysyłać dzienniki zdarzeń, dzienniki niestandardowe (takie jak log4j) i metryki do usługi Azure Log Analytics/Azure Storage/Azure Event Hubs. Przekaż nazwę log4j do właściwości : spark.synapse.diagnostic.emitter.\<destination\>.filter.loggerName.match.

    • Ponadto na potrzeby debugowania można również zbierać nieudane wiersze/rekordy w tabelach Lakehouse (LH) w celu przechwytywania nieprawidłowych danych na poziomie rekordu.