Strona główna▸Artykuły▸

Wnętrze Rozproszonego Modelu Wykonywania Apache Spark

Jak podział sterownika-wykonawcy w Spark’u, ocena opóźniona i wewnątrzwymiarowe przemieszanie danych pozwalają mu przewyższać MapReduce w przypadku obciążeń iteracyjnych, a także co naprawdę się dzieje, gdy uruchamia się zadanie Spark.

mysimulator teamZaktualizowano — czerwiec 2026≈ 5 min czytania▶ Otwórz symulację

Dlaczego MapReduce nie wystarczało

Dyscyplina MapReduce polegająca na zapisywaniu każdego pośredniego rezultatu na dysku pomiędzy etapami była bezpieczna, ale wolna, a koszty były najbardziej odczuwalne w algorytmach iteracyjnych – np. trenowanie modelu regresji logistycznej metodą gradientową, które mogły wymagać przeszukiwania tego samego zbioru danych kilkudziesięciokrotnie lub stu kilkadziesiąt razy. W systemie MapReduce każdy z tych pięćdziesięciu przebiegów był nowym zadaniem: odczytywanie z HDFS, obliczanie, zapisywanie z powrotem do HDFS i powtarzanie. Zakładem Sparka, opracowanym w AMPLab Uniwersytetu Kalifornijskiego w Berkeley, było to, że jeśli klastr ma wystarczającą ilość pamięci agregującej (RAM) do przechowywania danych roboczych – co jest prawdą dla wielu rzeczywistych obciążeń – to utrzymywanie danych w pamięci przez iteracje zamiast wielokrotnego przesyłania ich z dysku może skrócić czas trwania zadania z godzin na minuty.

Abstrakcją, która umożliwiała bezpieczne wykonywanie tego na dużą skalę, był Zbiór Rozproszonych i Odporny na Błędy (RDD): kolekcja podzielona na fragmenty rozłożona na klastrze, którą Spark może odtworzyć po awarii węzła nie poprzez replikację danych, ale poprzez odtwarzanie sekwencji transformacji, które je wygenerowały, korzystając z grafu pochodzenia. Jest to znacząca różnica w stosunku do replikacji bloków HDFS – zamiast ponosić koszty przechowywania na początku, Spark ponosi koszt ponownych obliczeń tylko wtedy, gdy awaria faktycznie wystąpiła, co jest wystarczająco rzadkie, aby amortyzowany koszt był niższy.

Ocena bezczynności i plan zapytania

Kod Spark czyta się tak, jakby wykonywał operacje natychmiast — wywołujesz `.filter()`, następnie `.groupBy()`, a w końcu `.agg()` — ale żaden z tych wywołań nie dotyka danych. Każdy z nich jedynie dodaje węzeł do logicznego planu, grafu skierowanego transformacji, i nic nie uruchamia się, dopóki nie zostanie wykonane działanie, takie jak `.collect()`, `.write()` lub `.count()`. Ta leniwość nie jest jedynie szczegółem implementacyjnym; to właśnie umożliwia optymalizatorowi Catalyst zobaczyć całą łańcuch operacji naraz i zrewidować ją przed wykonaniem — przesuwając filtr w przód, aby wykonywał się przed połączeniem, usuwając kolumny, których nikt nie potrzebował, lub przekształcając kolejność połączeń, aby umieścić najmniejszą tabelę na pierwszym miejscu.

Nieuwzględne podejście — wykonywanie każdej transformacji natychmiast, jedna po drugiej — marnowałoby ogromne ilości pracy: filtrowanie wierszy tylko po to, aby odrzucić większość kolumn złączonej tabeli dwa kroki później, gdy filtr mógłby zostać uruchomiony wcześniej i zmniejszyć wszystko, co znajduje się poniżej. Ponieważ Catalyst widzi cały plan na wgląd, może automatycznie zastosować taki rodzaj pushdown predykatów i projekcji bez konieczności ręcznego dostrajania kolejności operacji przez programistę — przekształcając to, co byłoby sekwencją naiwnych przejść przez cały zbiór danych, w plan, który dotyka minimalnej ilości danych na każdym etapie.

Sterownik, wykonawcy i przemieszanie

Uruchomione aplikacje Spark mają jeden sterownik, który przechowuje SparkContext, buduje plany logiczne i fizyczne oraz koordynuje wszystko, a także zestaw procesów wykonawczych rozproszonych na węzłach roboczych klastra, każdy działający w oddzielnym JVM z kawałkiem rdzeni CPU i pamięci. Sterownik dzieli fizyczny plan na etapy, dzieli każdy etap na wiele małych zadań – jedno zadanie na partycję danych – i wysyła te zadania do wykonawców do uruchomienia równolegle. Wykonawcy zgłaszają wyniki i sygnały serca sterownikowi, który jest również jedynym punktem awarii dla całej aplikacji: jeśli sterownik umrze, praca zakończy się, nawet jeśli wszystkie wykonawcy są zdrowe, dlatego w środowiskach produkcyjnych sterownik uruchamiany jest na odpornym węźle, a w przypadku długotrwałych lub krytycznych zadań może być wykonywany checkpoint stanu, aby restart nie musiał odtwarzać wszystkiego od nowa.

Granice etapów są określane przez przemieszania – każda operacja, taka jak `groupBy` lub `join` na danych niezgodnych z partycją, która wymaga przenoszenia rekordów między partycjami w sieci. W ramach etapu zadania działają niezależnie i nigdy nie muszą się ze sobą komunikować, dlatego Spark może potokowe wiele wąskich transformacji (map, filter) do jednego etapu bez tworzenia pośrednich wyników. Przemieszanie, z drugiej strony, wymusza twardy punkt synchronizacji: każde zadanie w górnym etapie musi zakończyć pisanie wyjściowych danych przemieszania przed rozpoczęciem przez jakiekolwiek zadanie w dolnym etapie odczytu, ponieważ partycja dolna może potrzebować danych przyczynionych przez wszystkie zadania górne. Jest to ten sam konceptualny wątek blokujący jak przemieszanie w MapReduce, ale Spark łagodzi jego koszt, utrzymując wyjściowe dane przemieszania w pamięci lub na lokalnym dysku w klastrze zamiast transportować je przez rozproszony system plików i pozwalając sterownikowi zminimalizować liczbę przemieszczeń, których potrzebuje zapytanie.

DataFrames, partycje i nierównomierne rozłożenie danych

Współczesny kod Spark w większości wykorzystuje API DataFrames zamiast surowych RDDs, co dodaje schemat — nazwy kolumn i typy — na szczycie tego samego modelu wykonania z partycjonowaniem i opóźnionym wykonywaniem, a ten schemat jest dokładnie tym, dzięki czemu Catalyst może rozumieć kolumny i typy wystarczająco dobrze, aby przeprowadzać optymalizacje. Pod spodem DataFrame nadal jest kolekcją partycji rozproszonych na wykonawcach, a liczba i rozmiar tych partycji wpływa materialnie na wydajność: zbyt mało partycji powoduje niewykorzystanie rdzeni klastra; zbyt wiele prowadzi do nadmiernych kosztów planowania zarządzania tysiącami małych zadań, które zaczynają dominować w rzeczywistej pracy.

Bardziej niepokojącym problemem jest nierównomierne rozłożenie danych: jeśli jeden klucz połączenia — na przykład bardzo popularny identyfikator produktu — odpowiada za niezrównane znaczenie liczby wierszy, partycja zawierająca ten klucz staje się ogromna, a jej sąsiedzi prawie puste, a czas wykonywania całego etapu skraca się do czasu trwania tej jednej przeciążonej zadania, niezależnie od tego, ile wykonawców stoi nieużywanych. Wykonywanie adaptacyjne zapytań Spark może wykryć to w czasie rzeczywistym i automatycznie podzielić dużą partycję na kilka mniejszych, ale inżynierowie regularnie ręcznie dostosowują się do znanych kluczy z obciążeniem, dodając do nich sól — dodając losowe sufiksy, aby rozłożyć gorący klucz na wiele partycji — ponieważ żaden automatyczny optymalizator nie wychwyci wszystkich rzeczywistych wzorców obciążenia.

Często zadawane pytania

Czym jest RDD i czy ludzie nadal go używają bezpośrednio?

RDD (Resilient Distributed Dataset) to oryginalne, partycjonowane, odporne na błędy zbiory danych Spark; większość nowoczesnego kodu wykorzystuje wyższego poziomu, schematowo świadomą API DataFrame, ale DataFrames wciąż kompilują operacje RDD pod spodem.

Dlaczego ocena opóźniona jest szybsza niż wykonywanie każdego kroku natychmiast?

Ponieważ nic nie wykonuje, dopóki nie zostanie wywołana akcja. Optymalizator Spark może zobaczyć całą łańcuch transformacji naraz i przekształcić ją – przesuwając filtry wcześniej, pomijając nieużywane kolumny, przekształcając kolejność połączeń – zamiast wykonywać każdy krok naiwnie w kolejności, w jakiej został napisany.

Co wywołuje shuffle w Spark i dlaczego jest to drogie?

Operacje takie jak groupBy lub łączenia na danych, które nie są już partycjonowane według klucza połączenia, zmuszają rekordy do przemieszczania się przez sieć między wykonawcami, co wymaga, aby wszystkie upstream zadania zakończyły się przed rozpoczęciem zadań downstream, tworząc twardy punkt synchronizacji, który blokuje równoległość.

Co to jest data skew i dlaczego szkodzi on zadaniom Spark?

Data skew występuje, gdy określone wartości kluczy są znacznie częstsze niż inne, powodując, że niektóre partycje zawierają znacznie więcej danych niż inne; ponieważ etap nie może zakończyć się dopóki jego najwolniejsze zadanie nie skończy, jedna duża partia może dominować całkowity czas wykonywania nawet wtedy, gdy gdzie indziej znajdują się nieaktywne wykonawcy.

Wypróbuj na żywo

Wszystko powyżej działa bezpośrednio w Twojej przeglądarce — otwórz the simulation i zmieniaj parametry podczas działania. Nic nie jest instalowane ani przesyłane na serwer, cały model działa w jednej karcie.

▶ Otwórz symulację the simulation

Co znalazłeś?

Dodaj kroki odtworzenia (opcjonalnie)