Usługa zarządzana dla Apache Spark

1. Omówienie usługi Serverless for Apache Spark

Usługa zarządzana dla Apache Spark to w pełni zarządzana i wysoce skalowalna usługa do uruchamiania Apache Spark, Apache Flink, Presto oraz wielu innych narzędzi i platform open source. Usługa zarządzana Apache Spark jest zintegrowana z Google Cloud i możesz z niej korzystać do modernizacji jezior danych, realizacji procesów ETL / ELT i bezpiecznego badania danych na globalną skalę. Zarządzana usługa Apache Spark jest też w pełni zintegrowana z kilkoma usługami Google Cloud, w tym BigQuery, Cloud Storage, Gemini Enterprise Agent Engine i Knowledge Catalog.

Zarządzana usługa Apache Spark jest dostępna w 2 wariantach:

  • Zarządzana usługa Apache Spark bez serwera umożliwia uruchamianie zadań PySpark bez konieczności konfigurowania infrastruktury i autoskalowania. Usługa zarządzana Apache Spark obsługuje zadania wsadowe i sesje / notatniki PySpark.
  • Zarządzane klastry Apache Spark umożliwiają zarządzanie klastrem Hadoop YARN dla zadań Spark opartych na YARN, a także narzędziami open source, takimi jak Flink i Presto. Możesz dostosować klastry w chmurze za pomocą dowolnego skalowania w pionie lub w poziomie, w tym autoskalowania.

W tym ćwiczeniu dowiesz się, jak korzystać z Dataproc Serverless na kilka różnych sposobów.

Platforma Apache Spark została pierwotnie opracowana z myślą o jej uruchamianiu w klastrach Hadoop, a jako menedżera zasobów używała rozwiązania YARN. Obsługa klastrów Hadoop wymaga specjalistycznej wiedzy i zapewnienia prawidłowej konfiguracji wielu różnych ustawień klastrów. Oprócz tego Spark wymaga od użytkownika ustawienia osobnego zestawu parametrów. W wielu przypadkach programiści poświęcają więcej czasu na konfigurowanie infrastruktury niż na pracę nad samym kodem Spark.

Dataproc Serverless eliminuje konieczność ręcznego konfigurowania klastrów Hadoop lub Spark. Dataproc Serverless nie działa na platformie Hadoop i wykorzystuje własne dynamiczne przydzielanie zasobów do określania wymagań dotyczących zasobów, w tym autoskalowania. Za pomocą Dataproc Serverless można dostosować niewielki podzbiór usług Spark, ale w większości przypadków nie trzeba ich zmieniać.

2. Skonfiguruj

Zaczniesz od skonfigurowania środowiska i zasobów używanych w tych ćwiczeniach z programowania.

Utwórz projekt Google Cloud. Możesz użyć istniejącego.

Otwórz Cloud Shell, klikając go na pasku narzędzi Cloud Console.

ba0bb17945a73543.png

Cloud Shell udostępnia gotowe do użycia środowisko powłoki, z którego możesz korzystać w tym ćwiczeniu.

68c4ebd2a8539764.png

Cloud Shell domyślnie ustawi nazwę projektu. Sprawdź to, uruchamiając echo $GOOGLE_CLOUD_PROJECT. Jeśli w danych wyjściowych nie widzisz identyfikatora projektu, ustaw go.

export GOOGLE_CLOUD_PROJECT=<your-project-id>

Ustaw dla swoich zasobów region Compute Engine, na przykład us-central1 lub europe-west2.

export REGION=<your-region>

Włącz interfejsy API

W tym ćwiczeniu wykorzystujemy te interfejsy API:

  • BigQuery
  • Dataproc

Włącz niezbędne interfejsy API. Zajmie to około minuty. Po zakończeniu pojawi się komunikat o powodzeniu.

gcloud services enable bigquery.googleapis.com
gcloud services enable dataproc.googleapis.com

Konfigurowanie dostępu do sieci

Sterowniki i wykonawcy Spark mają tylko prywatne adresy IP, dlatego usługa Dataproc Serverless musi mieć włączony prywatny dostęp do Google w regionie, w którym uruchamiasz zadania Spark. Aby włączyć tę opcję w podsieci default, uruchom to polecenie:

gcloud compute networks subnets update default \
  --region=${REGION} \
  --enable-private-ip-google-access

Możesz sprawdzić, czy prywatny dostęp do Google jest włączony, korzystając z tego polecenia, które zwróci wartość True lub False.

gcloud compute networks subnets describe default \
  --region=${REGION} \
  --format="get(privateIpGoogleAccess)"

Utworzenie zasobnika na dane

Utwórz zasobnik na dane, w którym będą przechowywane zasoby utworzone w tym laboratorium.

Wybierz nazwę zasobnika. Nazwy zasobników muszą być globalnie niepowtarzalne dla wszystkich użytkowników.

export BUCKET=<your-bucket-name>

Utwórz zasobnik w regionie, w którym zamierzasz uruchamiać zadania Spark.

gsutil mb -l ${REGION} gs://${BUCKET}

Zasobnik jest dostępny w konsoli Cloud Storage. Możesz też uruchomić polecenie gsutil ls, aby wyświetlić zasobnik.

Tworzenie stałego serwera historii

Interfejs Spark udostępnia bogaty zestaw narzędzi do debugowania i informacji o zadaniach Spark. Aby wyświetlić interfejs Spark dla ukończonych zadań Dataproc Serverless, musisz utworzyć klaster Dataproc z pojedynczym węzłem, który będzie używany jako stały serwer historii.

Ustaw nazwę stałego serwera historii.

PHS_CLUSTER_NAME=my-phs

Wykonaj te czynności.

gcloud dataproc clusters create ${PHS_CLUSTER_NAME} \
    --region=${REGION} \
    --single-node \
    --enable-component-gateway \
    --properties=spark:spark.history.fs.logDirectory=gs://${BUCKET}/phs/*/spark-job-history

Interfejs Spark i stały serwer historii zostaną omówione bardziej szczegółowo w dalszej części tego samouczka.

3. Uruchamianie zadań Serverless Spark za pomocą wsadów Dataproc

W tym przykładzie będziesz pracować na danych z publicznego zbioru danych New York City (NYC) Citi Bike Trips. NYC Citi Bikes to płatny nowojorski system rowerów publicznych. Wykonasz proste przekształcenia i wydrukujesz 10 najpopularniejszych identyfikatorów stacji rowerów publicznych Citi Bike. W tym przykładzie użyto też oprogramowania sprzęgającego spark-bigquery-connector typu open source, aby bezproblemowo odczytywać i zapisywać dane między Sparkiem a BigQuery.

Skopiuj to repozytorium GitHub i przejdź (cd) do katalogu zawierającego plik citibike.py.

git clone https://github.com/GoogleCloudPlatform/devrel-demos.git
cd devrel-demos/data-analytics/next-2022-workshop/dataproc-serverless

citibike.py

import sys

from pyspark.sql import SparkSession
from pyspark.sql.functions import col
from pyspark.sql.types import BooleanType

if len(sys.argv) == 1:
    print("Please provide a GCS bucket name.")

bucket = sys.argv[1]
table = "bigquery-public-data:new_york_citibike.citibike_trips"

spark = SparkSession.builder \
          .appName("pyspark-example") \
          .config("spark.jars","gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar") \
          .getOrCreate()

df = spark.read.format("bigquery").load(table)

top_ten = df.filter(col("start_station_id") \
            .isNotNull()) \
            .groupBy("start_station_id") \
            .count() \
            .orderBy("count", ascending=False) \
            .limit(10) \
            .cache()

top_ten.show()

top_ten.write.option("header", True).csv(f"gs://{bucket}/citibikes_top_ten_start_station_ids")

Prześlij zadanie do Serverless Spark za pomocą Cloud SDK, który jest domyślnie dostępny w Cloud Shell. Uruchom w powłoce to polecenie, które korzysta z Cloud SDK i Dataproc Batches API do przesyłania zadań Serverless Spark.

gcloud dataproc batches submit pyspark citibike.py \
  --batch=citibike-job \
  --region=${REGION} \
  --deps-bucket=gs://${BUCKET} \
  --jars=gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar \
--history-server-cluster=projects/${GOOGLE_CLOUD_PROJECT}/regions/${REGION}/clusters/${PHS_CLUSTER_NAME} \
  -- ${BUCKET}

Przeanalizujmy poszczególne elementy tej misji:

  • gcloud dataproc batches submit odwołuje się do interfejsu Dataproc Batches API.
  • pyspark oznacza, że przesyłasz zadanie PySpark.
  • --batch to nazwa zadania. Jeśli nie podasz identyfikatora, użyty zostanie losowy identyfikator UUID.
  • --region=${REGION} to region geograficzny, w którym zadanie będzie przetwarzane.
  • --deps-bucket=${BUCKET} to miejsce, do którego przesyłany jest lokalny plik Pythona przed uruchomieniem w środowisku bezserwerowym.
  • --jars=gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar zawiera plik JAR dla spark-bigquery-connector w środowisku wykonawczym Spark.
  • --history-server-cluster=projects/${GOOGLE_CLOUD_PROJECT}/regions/${REGION}/clusters/${PHS_CLUSTER} to pełna i jednoznaczna nazwa stałego serwera historii. W tym miejscu są przechowywane dane zdarzeń Spark (osobne od danych wyjściowych konsoli), które można wyświetlać w interfejsie Spark.
  • Znak -- na końcu oznacza, że wszystko, co znajduje się za nim, będzie argumentami środowiska wykonawczego programu. W tym przypadku przesyłasz nazwę zasobnika, zgodnie z wymaganiami zadania.

Po przesłaniu wsadu pojawią się następujące dane wyjściowe.

Batch [citibike-job] submitted.

Po kilku minutach zobaczysz te dane wyjściowe wraz z metadanymi z zadania.

+----------------+------+
|start_station_id| count|
+----------------+------+
|             519|551078|
|             497|423334|
|             435|403795|
|             426|384116|
|             293|372255|
|             402|367194|
|             285|344546|
|             490|330378|
|             151|318700|
|             477|311403|
+----------------+------+

Batch [citibike-job] finished.

W następnej sekcji dowiesz się, jak znaleźć logi tego zadania.

Dodatkowe funkcje

Spark Serverless oferuje dodatkowe opcje uruchamiania zadań.

  • Możesz utworzyć niestandardowy obraz Dockera, na którym będzie uruchamiane zadanie. To doskonały sposób na dodanie kolejnych zależności, takich jak biblioteki R i Pythona.
  • Aby uzyskać dostęp do metadanych Hive, możesz połączyć zadanie z instancją Dataproc Metastore.
  • Z myślą o dodatkowej kontroli Dataproc Serverless obsługuje konfigurację niewielkiego zbioru usług Spark.

4. Wskaźniki i dostrzegalność Dataproc

W konsoli zadań wsadowych Dataproc znajdziesz listę wszystkich zadań Dataproc Serverless. W konsoli zobaczysz Identyfikator wsadu, Lokalizację, Stan, Czas utworzenia, Czas trwaniaTyp każdego zadania. Kliknij Identyfikator wsadu zadania, aby wyświetlić więcej informacji na jego temat.

Na tej stronie znajdziesz takie informacje jak Monitorowanie, które pokazuje liczbę wykonawców wsadowych Spark używanych przez zadanie na przestrzeni czasu (wskazującą, jak bardzo zadanie zostało przeskalowane automatycznie).

Na karcie Szczegóły znajdziesz więcej metadanych zadania, w tym argumenty i parametry przesłane wraz z zadaniem.

Na tej stronie możesz też uzyskać dostęp do wszystkich logów. Podczas wykonywania zadań Dataproc Serverless generowane są 3 różne zestawy logów:

  • Na poziomie usługi
  • dane wyjściowe konsoli,
  • Rejestrowanie zdarzeń Spark

Na poziomie usługi – obejmują logi wygenerowane przez usługę Dataproc Serverless. Zawierają one m.in. żądania dodatkowych procesorów na potrzeby autoskalowania. Możesz je wyświetlić, klikając Wyświetl logi, co spowoduje otwarcie Cloud Logging.

Dane wyjściowe konsoli można wyświetlić w sekcji Wyniki. To dane wyjściowe wygenerowane przez zadanie, w tym metadane, które Spark drukuje przy rozpoczęciu zadania, lub wszelkie instrukcje print włączone do zadania.

Logowanie zdarzeń Spark jest dostępne w interfejsie Spark. Ponieważ zadanie Spark zostało uruchomione na stałym serwerze historii, możesz uzyskać dostęp do interfejsu Spark, klikając Wyświetl serwer historii usługi Spark. Znajdziesz tam informacje o wcześniej uruchomionych zadaniach Spark. Więcej informacji o interfejsie Spark znajdziesz w oficjalnej dokumentacji Spark.

5. Szablony Dataproc: BQ -> GCS

Szablony Dataproc to narzędzia open source, które pomagają jeszcze bardziej uprościć zadania związane z przetwarzaniem danych w chmurze. Są one otoką usługi Dataproc Serverless i zawierają szablony wielu zadań importu i eksportu danych, w tym:

  • BigQuerytoGCSGCStoBigQuery
  • GCStoBigTable
  • GCStoJDBCJDBCtoGCS
  • HivetoBigQuery
  • MongotoGCSGCStoMongo

Pełna lista jest dostępna w README.

W tej sekcji użyjesz szablonów Dataproc do eksportowania danych z BigQuery do GCS.

Kopiowanie repozytorium

Skopiuj repozytorium i przejdź do folderu python.

git clone https://github.com/GoogleCloudPlatform/dataproc-templates.git
cd dataproc-templates/python

Konfigurowanie środowiska

Teraz ustawisz zmienne środowiskowe. Szablony Dataproc używają zmiennej środowiskowej GCP_PROJECT dla identyfikatora projektu, więc ustaw ją na GOOGLE_CLOUD_PROJECT..

export GCP_PROJECT=${GOOGLE_CLOUD_PROJECT}

Region powinien być ustawiony w środowisku z wcześniejszego etapu. Jeśli nie, ustaw go tutaj.

export REGION=<region>

Szablony Dataproc przetwarzają zadania BigQuery za pomocą narzędzia spark-bigquery-connector i wymagają, aby identyfikator URI był uwzględniony w zmiennej środowiskowej JARS. Ustaw zmienną JARS.

export JARS="gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar"

Konfigurowanie parametrów szablonu

Ustaw nazwę zasobnika tymczasowego, którego usługa ma używać.

export GCS_STAGING_LOCATION=gs://${BUCKET}

Następnie ustawisz zmienne specyficzne dla zadania. W przypadku tabeli wejściowej ponownie odwołasz się do zbioru danych BigQuery NYC Citibike.

BIGQUERY_GCS_INPUT_TABLE=bigquery-public-data.new_york_citibike.citibike_trips

Możesz wybrać csv, parquet, avro lub json. W tym samouczku wybierz CSV. W następnej sekcji dowiesz się, jak używać szablonów Dataproc do konwertowania typów plików.

BIGQUERY_GCS_OUTPUT_FORMAT=csv

Ustaw tryb wyników na overwrite. Możesz wybrać overwrite, append, ignore lub errorifexists.

BIGQUERY_GCS_OUTPUT_MODE=overwrite

Ustaw lokalizację wyjściową GCS jako ścieżkę w zasobniku.

BIGQUERY_GCS_OUTPUT_LOCATION=gs://${BUCKET}/BQtoGCS

Uruchamianie szablonu

Uruchom szablon BIGQUERYTOGCS, określając go poniżej i podając ustawione parametry wejściowe.

./bin/start.sh \
-- --template=BIGQUERYTOGCS \
        --bigquery.gcs.input.table=${BIGQUERY_GCS_INPUT_TABLE} \
        --bigquery.gcs.output.format=${BIGQUERY_GCS_OUTPUT_FORMAT} \
        --bigquery.gcs.output.mode=${BIGQUERY_GCS_OUTPUT_MODE} \
        --bigquery.gcs.output.location=${BIGQUERY_GCS_OUTPUT_LOCATION}

Dane wyjściowe będą dość zaszumione, ale po około minucie zobaczysz te informacje:

Batch [5766411d6c78444cb5e80f305308d8f8] submitted.
...
Batch [5766411d6c78444cb5e80f305308d8f8] finished.

Aby sprawdzić, czy pliki zostały wygenerowane, uruchom to polecenie:

gsutil ls ${BIGQUERY_GCS_OUTPUT_LOCATION}

Spark domyślnie zapisuje dane w wielu plikach, w zależności od ilości danych. W tym przypadku zobaczysz około 30 wygenerowanych plików. Nazwy plików wyjściowych Spark mają przedrostek part, po którym następuje 5-cyfrowy numer (reprezentujący numer części) oraz ciąg znaków skrótu. Gdy danych jest dużo, Spark zazwyczaj zapisuje je w kilku plikach. Przykładowa nazwa pliku to part-00000-cbf69737-867d-41cc-8a33-6521a725f7a0-c000.csv.

6. Szablony Dataproc: CSV do Parquet

Teraz użyjesz szablonów Dataproc, aby przekonwertować dane w GCS z jednego typu pliku na inny za pomocą szablonu GCSTOGCS. Ten szablon korzysta z SparkSQL i umożliwia też przesyłanie zapytania SparkSQL do przetworzenia podczas transformacji w celu dodatkowego przetwarzania.

Potwierdzanie zmiennych środowiskowych

Sprawdź, czy zmienne środowiskowe GCP_PROJECT, REGION i GCS_STAGING_BUCKET są ustawione w poprzedniej sekcji.

echo ${GCP_PROJECT}
echo ${REGION}
echo ${GCS_STAGING_LOCATION}

Ustawianie parametrów szablonu

Teraz ustawisz parametry konfiguracji dla GCStoGCS. Zacznij od lokalizacji plików wejściowych. Pamiętaj, że jest to katalog, a nie konkretny plik, ponieważ wszystkie pliki w katalogu zostaną przetworzone. Ustaw tę wartość na BIGQUERY_GCS_OUTPUT_LOCATION.

GCS_TO_GCS_INPUT_LOCATION=${BIGQUERY_GCS_OUTPUT_LOCATION}

Ustaw format pliku wejściowego.

GCS_TO_GCS_INPUT_FORMAT=csv

Ustaw żądany format wyjściowy. Możesz wybrać format Parquet, JSON, Avro lub CSV.

GCS_TO_GCS_OUTPUT_FORMAT=parquet

Ustaw tryb wyników na overwrite. Możesz wybrać overwrite, append, ignore lub errorifexists.

GCS_TO_GCS_OUTPUT_MODE=overwrite

Ustaw lokalizację wyjściową.

GCS_TO_GCS_OUTPUT_LOCATION=gs://${BUCKET}/GCStoGCS

Uruchamianie szablonu

Uruchom szablon GCStoGCS.

./bin/start.sh \
-- --template=GCSTOGCS \
        --gcs.to.gcs.input.location=${GCS_TO_GCS_INPUT_LOCATION} \
        --gcs.to.gcs.input.format=${GCS_TO_GCS_INPUT_FORMAT} \
        --gcs.to.gcs.output.format=${GCS_TO_GCS_OUTPUT_FORMAT} \
        --gcs.to.gcs.output.mode=${GCS_TO_GCS_OUTPUT_MODE} \
        --gcs.to.gcs.output.location=${GCS_TO_GCS_OUTPUT_LOCATION}

Dane wyjściowe będą dość zaszumione, ale po około minucie powinien wyświetlić się komunikat o powodzeniu podobny do tego poniżej.

Batch [c198787ba8e94abc87e2a0778c05ec8a] submitted.
...
Batch [c198787ba8e94abc87e2a0778c05ec8a] finished.

Aby sprawdzić, czy pliki zostały wygenerowane, uruchom to polecenie:

gsutil ls ${GCS_TO_GCS_OUTPUT_LOCATION}

Ten szablon umożliwia też dostarczanie zapytań SparkSQL przez przekazywanie do szablonu parametrów gcs.to.gcs.temp.view.name i gcs.to.gcs.sql.query, co pozwala uruchamiać zapytania SparkSQL na danych przed zapisaniem ich w GCS.

7. Zwalnianie miejsca

Aby uniknąć niepotrzebnych opłat na koncie Google Cloud Platform po ukończeniu tego laboratorium:

  1. Usuń zasobnik Cloud Storage dla utworzonego środowiska.
gsutil rm -r gs://${BUCKET}
  1. Usuń klaster Dataproc używany na potrzeby stałego serwera historii.
gcloud dataproc clusters delete ${PHS_CLUSTER_NAME} \
  --region=${REGION}
  1. Usuń zadania Dataproc Serverless. Otwórz konsolę Batches, kliknij pole obok każdego zadania, które chcesz usunąć, a następnie kliknij USUŃ.

Jeśli projekt został utworzony specjalnie na potrzeby tego ćwiczenia, możesz go też usunąć:

  1. W konsoli GCP otwórz stronę Projekty.
  2. Na liście projektów wybierz projekt, który chcesz usunąć, i kliknij Usuń.
  3. W polu wpisz identyfikator projektu i kliknij Wyłącz, aby usunąć projekt.

8. Co dalej?

Poniższe materiały zawierają dodatkowe informacje o tym, jak korzystać z Serverless Spark: