Uruchamiamy Apache Spark na Kubernetes

Drodzy czytelnicy, dzień dobry. Dzisiaj porozmawiamy trochę o Apache Spark i jego możliwościach rozwoju.

Uruchamiamy Apache Spark na Kubernetes

W dzisiejszym świecie Big Data, Apache Spark jest de facto standardem przy tworzeniu zadań przetwarzania wsadowego. Oprócz tego, jest również wykorzystywany do tworzenia aplikacji strumieniowych, działających w koncepcji micro batch, które przetwarzają i dostarczają dane w małych porcjach (Spark Structured Streaming). Tradycyjnie był częścią ogólnego stosu Hadoop, wykorzystując jako menedżera zasobów YARN (lub, w niektórych przypadkach, Apache Mesos). Do 2020 roku jego użycie w tradycyjnej formie dla większości firm staje pod dużym znakiem zapytania z powodu braku odpowiednich dystrybucji Hadoop — rozwój HDP i CDH został wstrzymany, CDH jest niewystarczająco przemyślany i ma wysokie koszty, a inni dostawcy Hadoop albo zaprzestali działalności, albo mają niepewną przyszłość. Dlatego rosnące zainteresowanie wśród społeczności oraz dużych firm zyskuje uruchamianie Apache Spark za pomocą Kubernetes — jako standard w orkiestracji kontenerów i zarządzaniu zasobami w chmurach prywatnych i publicznych, rozwiązuje problem niewygodnego planowania zasobów zadań Spark na YARN i oferuje stabilnie rozwijającą się platformę z wieloma komercyjnymi i otwartymi dystrybucjami dla firm każdej wielkości. Dodatkowo, na fali popularności wiele z nich już zdążyło zbudować kilka instalacji i zdobyć doświadczenie w jego użyciu, co ułatwia migrację.

Od wersji 2.3.0 Apache Spark zyskał oficjalne wsparcie dla uruchamiania zadań w klastrze Kubernetes i dzisiaj porozmawiamy o obecnej dojrzałości tego podejścia, różnych jego zastosowaniach oraz pułapkach, z którymi można się spotkać przy wdrażaniu.

Przede wszystkim przyjrzymy się procesowi tworzenia zadań i aplikacji opartych na Apache Spark oraz wyróżnimy typowe przypadki, w których trzeba uruchomić zadanie w klastrze Kubernetes. W trakcie przygotowywania tego posta wykorzystano dystrybucję OpenShift i przedstawione zostaną polecenia, które są aktualne dla jego narzędzia wiersza poleceń (oc). Dla innych dystrybucji Kubernetes można używać odpowiednich poleceń standardowego narzędzia wiersza poleceń Kubernetes (kubectl) lub ich analogów (na przykład dla oc adm policy).

Pierwsza opcja użycia — spark-submit

W trakcie tworzenia zadań i aplikacji deweloper musi uruchamiać zadania do debugowania transformacji danych. Teoretycznie można do tych celów używać stubów, ale praca z rzeczywistymi (choćby testowymi) instancjami systemów końcowych okazała się w tym typie zadań szybsza i lepsza. W przypadku debugowania na rzeczywistych instancjach systemów końcowych możliwe są dwa scenariusze działania:

  • deweloper uruchamia zadanie Spark lokalnie w trybie standalone;

    Uruchamiamy Apache Spark na Kubernetes

  • deweloper uruchamia zadanie Spark na klastrze Kubernetes w konturze testowym.

    Uruchamiamy Apache Spark na Kubernetes

Pierwsza opcja jest dopuszczalna, ale wiąże się z pewnymi wadami:

  • dla każdego dewelopera konieczne jest zapewnienie dostępu z miejsca pracy do wszystkich niezbędnych mu instancji systemów końcowych;
  • na komputerze roboczym musi być wystarczająca ilość zasobów do uruchomienia opracowywanego zadania.

Druga opcja nie ma tych wad, ponieważ wykorzystanie klastra Kubernetes pozwala przydzielić odpowiednią pulę zasobów do uruchamiania zadań i zapewnić wymagane dostępności do instancji systemów końcowych, elastycznie udostępniając do nich dostęp za pomocą modelu ról Kubernetes dla wszystkich członków zespołu deweloperskiego. Wyróżnimy to jako pierwszą opcję użycia — uruchamianie zadań Spark z lokalnego komputera dewelopera w klasterze Kubernetes w konturze testowym.

Opowiemy więcej o procesie konfiguracji Sparka do lokalnego uruchamiania. Aby zacząć korzystać z Sparka, należy go zainstalować:

mkdir /opt/spark
cd /opt/spark
wget http://mirror.linux-ia64.org/apache/spark/spark-2.4.5/spark-2.4.5.tgz
tar zxvf spark-2.4.5.tgz
rm -f spark-2.4.5.tgz

Zbieramy niezbędne pakiety do pracy z Kubernetes:

cd spark-2.4.5/
./build/mvn -Pkubernetes -DskipTests clean package

Pełne zbudowanie zajmuje dużo czasu, a do tworzenia obrazów Docker i uruchamiania ich na klastrze Kubernetes w rzeczywistości potrzebne są tylko pliki jar z katalogu „assembly/”, dlatego można zbudować tylko ten podprojekt:

./build/mvn -f ./assembly/pom.xml -Pkubernetes -DskipTests clean package

Aby uruchomić zadania Spark w Kubernetes, należy stworzyć obraz Docker, który będzie używany jako bazowy. Możliwe są tutaj 2 podejścia:

  • Stworzony obraz Docker zawiera kod wykonywalny zadania Spark;
  • Stworzony obraz zawiera tylko Spark i niezbędne zależności, kod wykonywalny jest zdalnie umieszczany (na przykład w HDFS).

Na początek zbudujemy obraz Docker, zawierający testowy przykład zadania Spark. Do tworzenia obrazów Docker Spark ma odpowiednie narzędzie o nazwie „docker-image-tool”. Zapoznajmy się z jego instrukcją:

./bin/docker-image-tool.sh --help

Za jego pomocą można tworzyć obrazy Docker i przesyłać je do zdalnych rejestrów, ale domyślnie ma szereg wad:

  • zobowiązany do tworzenia od razu 3 obrazów Docker — dla Spark, PySpark i R;
  • nie pozwala na podanie nazwy obrazu.

Dlatego będziemy korzystać z zmodyfikowanej wersji tego narzędzia, przedstawionej poniżej:

vi bin/docker-image-tool-upd.sh

#!/usr/bin/env bash

function error {
  echo "$@" 1>&2
  exit 1
}

if [ -z "${SPARK_HOME}" ]; then
  SPARK_HOME="$(cd "`dirname "$0"`"/..; pwd)"
fi
. "${SPARK_HOME}/bin/load-spark-env.sh"

function image_ref {
  local image="$1"
  local add_repo="${2:-1}"
  if [ $add_repo = 1 ] && [ -n "$REPO" ]; then
    image="$REPO/$image"
  fi
  if [ -n "$TAG" ]; then
    image="$image:$TAG"
  fi
  echo "$image"
}

function build {
  local BUILD_ARGS
  local IMG_PATH

  if [ ! -f "$SPARK_HOME/RELEASE" ]; then
    IMG_PATH=$BASEDOCKERFILE
    BUILD_ARGS=(
      ${BUILD_PARAMS}
      --build-arg
      img_path=$IMG_PATH
      --build-arg
      datagram_jars=datagram/runtimelibs
      --build-arg
      spark_jars=assembly/target/scala-$SPARK_SCALA_VERSION/jars
    )
  else
    IMG_PATH="kubernetes/dockerfiles"
    BUILD_ARGS=(${BUILD_PARAMS})
  fi

  if [ -z "$IMG_PATH" ]; then
    error "Cannot find docker image. This script must be run from a runnable distribution of Apache Spark."
  fi

  if [ -z "$IMAGE_REF" ]; then
    error "Cannot find docker image reference. Please add -i arg."
  fi

  local BINDING_BUILD_ARGS=(
    ${BUILD_PARAMS}
    --build-arg
    base_img=$(image_ref $IMAGE_REF)
  )
  local BASEDOCKERFILE=${BASEDOCKERFILE:-"$IMG_PATH/spark/docker/Dockerfile"}

  docker build $NOCACHEARG "${BUILD_ARGS[@]}" 
    -t $(image_ref $IMAGE_REF) 
    -f "$BASEDOCKERFILE" .
}

function push {
  docker push "$(image_ref $IMAGE_REF)"
}

function usage {
  cat <<EOF
Usage: $0 [options] [command]
Builds or pushes the built-in Spark Docker image.

Commands:
  build       Build image. Requires a repository address to be provided if the image will be
              pushed to a different registry.
  push        Push a pre-built image to a registry. Requires a repository address to be provided.

Options:
  -f file               Dockerfile to build for JVM based Jobs. By default builds the Dockerfile shipped with Spark.
  -p file               Dockerfile to build for PySpark Jobs. Builds Python dependencies and ships with Spark.
  -R file               Dockerfile to build for SparkR Jobs. Builds R dependencies and ships with Spark.
  -r repo               Repository address.
  -i name               Image name to apply to the built image, or to identify the image to be pushed.  
  -t tag                Tag to apply to the built image, or to identify the image to be pushed.
  -m                    Use minikube's Docker daemon.
  -n                    Build docker image with --no-cache
  -b arg      Build arg to build or push the image. For multiple build args, this option needs to
              be used separately for each build arg.

Using minikube when building images will do so directly into minikube's Docker daemon.
There is no need to push the images into minikube in that case, they'll be automatically
available when running applications inside the minikube cluster.

Check the following documentation for more information on using the minikube Docker daemon:

  https://kubernetes.io/docs/getting-started-guides/minikube/#reusing-the-docker-daemon

Examples:
  - Build image in minikube with tag "testing"
    $0 -m -t testing build

  - Build and push image with tag "v2.3.0" to docker.io/myrepo
    $0 -r docker.io/myrepo -t v2.3.0 build
    $0 -r docker.io/myrepo -t v2.3.0 push
EOF
}

if [[ "$@" = *--help ]] || [[ "$@" = *-h ]]; then
  usage
  exit 0
fi

REPO=
TAG=
BASEDOCKERFILE=
NOCACHEARG=
BUILD_PARAMS=
IMAGE_REF=
while getopts f:mr:t:nb:i: option
do
 case "${option}"
 in
 f) BASEDOCKERFILE=${OPTARG};;
 r) REPO=${OPTARG};;
 t) TAG=${OPTARG};;
 n) NOCACHEARG="--no-cache";;
 i) IMAGE_REF=${OPTARG};;
 b) BUILD_PARAMS=${BUILD_PARAMS}" --build-arg "${OPTARG};;
 esac
done

case "${@: -1}" in
  build)
    build
    ;;
  push)
    if [ -z "$REPO" ]; then
      usage
      exit 1
    fi
    push
    ;;
  *)
    usage
    exit 1
    ;;
esac

Za jego pomocą zbieramy podstawowy obraz Spark, zawierający testowe zadanie obliczania liczby Pi przy użyciu Spark (tutaj {docker-registry-url} — URL twojego rejestru obrazów Docker, {repo} — nazwa repozytorium w rejestrze, odpowiadająca projektowi w OpenShift, {image-name} — nazwa obrazu (jeśli używane jest trójwarstwowe oddzielenie obrazów, na przykład jak w zintegrowanym rejestrze obrazów Red Hat OpenShift), {tag} — tag danej wersji obrazu):

./bin/docker-image-tool-upd.sh -f resource-managers/kubernetes/docker/src/main/dockerfiles/spark/Dockerfile -r {docker-registry-url}/{repo} -i {image-name} -t {tag} build

Logujemy się do klastra OKD za pomocą narzędzia konsolowego (tutaj {OKD-API-URL} — URL API klastra OKD):

oc login {OKD-API-URL}

Uzyskamy token aktualnego użytkownika do autoryzacji w Docker Registry:

oc whoami -t

Logujemy się do wewnętrznego Docker Registry klastra OKD (jako hasło używamy tokenu uzyskanego za pomocą poprzedniej komendy):

docker login {docker-registry-url}

Prześlemy zbudowany obraz Docker do Docker Registry OKD:

./bin/docker-image-tool-upd.sh -r {docker-registry-url}/{repo} -i {image-name} -t {tag} push

Sprawdźmy, czy zebrany obraz jest dostępny w OKD. W tym celu otwieramy w przeglądarce URL z listą obrazów odpowiedniego projektu (tutaj {project} — nazwa projektu w klastrze OpenShift, {OKD-WEBUI-URL} — adres URL konsoli Web OpenShift) — https://{OKD-WEBUI-URL}/console/project/{project}/browse/images/{image-name}.

Aby uruchomić zadania, musi zostać utworzone konto usługi z uprawnieniami do uruchamiania podów jako root (dzisiaj omówimy tę kwestię):

oc create sa spark -n {project}
oc adm policy add-scc-to-user anyuid -z spark -n {project}

Wykonamy polecenie spark-submit, aby opublikować zadanie Spark w klastrze OKD, wskazując utworzone konto usługi i obraz Docker:

 /opt/spark/bin/spark-submit --name spark-test --class org.apache.spark.examples.SparkPi --conf spark.executor.instances=3 --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark --conf spark.kubernetes.namespace={project} --conf spark.submit.deployMode=cluster --conf spark.kubernetes.container.image={docker-registry-url}/{repo}/{image-name}:{tag} --conf spark.master=k8s://https://{OKD-API-URL}  local:///opt/spark/examples/target/scala-2.11/jars/spark-examples_2.11-2.4.5.jar

Gdzie:

—name — nazwa zadania, która weźmie udział w tworzeniu nazw podów Kubernetes;

—class — klasa pliku wykonywalnego wywoływana podczas uruchamiania zadania;

—conf — parametry konfiguracyjne Spark;

spark.executor.instances — liczba uruchamianych executorów Spark;

spark.kubernetes.authenticate.driver.serviceAccountName — nazwa konta usługi Kubernetes używanego podczas uruchamiania podów (do określenia kontekstu bezpieczeństwa i uprawnień podczas interakcji z API Kubernetes);

spark.kubernetes.namespace — przestrzeń nazw Kubernetes, w której będą uruchamiane pody sterownika i executorów;

spark.submit.deployMode — sposób uruchamiania Spark (dla standardowego spark-submit używa się „cluster”, dla Spark Operator i nowszych wersji Spark „client”);

spark.kubernetes.container.image — obraz Docker używany do uruchamiania podów;

spark.master — URL API Kubernetes (podawany zewnętrzny, gdyż dostęp ma miejsce z lokalnej maszyny);

local:// — ścieżka do pliku wykonywalnego Spark wewnątrz obrazu Docker.

Przechodzimy do odpowiedniego projektu OKD i studiujemy utworzone pody — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods.

Aby uprościć proces rozwoju, można zastosować jeszcze jeden wariant, w którym tworzy się wspólny obraz bazowy Spark, używany przez wszystkie zadania do uruchamiania, a migawki plików wykonywalnych publikowane są w zewnętrznym magazynie (na przykład Hadoop) i wskazywane przy wywołaniu spark-submit w postaci odnośnika. W takim przypadku można uruchamiać różne wersje zadań Spark bez ponownego budowania obrazów Docker, używając do publikacji obrazów, na przykład WebHDFS. Wysyłamy żądanie utworzenia pliku (tutaj {host} — host usługi WebHDFS, {port} — port usługi WebHDFS, {path-to-file-on-hdfs} — pożądana ścieżka do pliku na HDFS):

curl -i -X PUT "http://{host}:{port}/webhdfs/v1/{path-to-file-on-hdfs}?op=CREATE"

Odpowiedź będzie miała postać (gdzie {location} to URL, który należy użyć do załadowania pliku):

HTTP/1.1 307 TEMPORARY_REDIRECT
Location: {location}
Content-Length: 0

Ładujemy plik wykonywalny Spark do HDFS (gdzie {path-to-local-file} to ścieżka do pliku wykonywalnego Spark na bieżącym hoście):

curl -i -X PUT -T {path-to-local-file} "{location}"

Po tym możemy uruchomić spark-submit z wykorzystaniem pliku Spark załadowanego na HDFS (gdzie {class-name} to nazwa klasy, którą należy uruchomić do wykonania zadania):

/opt/spark/bin/spark-submit --name spark-test --class {class-name} --conf spark.executor.instances=3 --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark --conf spark.kubernetes.namespace={project} --conf spark.submit.deployMode=cluster --conf spark.kubernetes.container.image={docker-registry-url}/{repo}/{image-name}:{tag} --conf spark.master=k8s://https://{OKD-API-URL}  hdfs://{host}:{port}/{path-to-file-on-hdfs}

Należy zauważyć, że do uzyskania dostępu do HDFS i zapewnienia działania zadania może być konieczne zmodyfikowanie Dockerfile oraz skryptu entrypoint.sh — dodanie w Dockerfile dyrektywy kopiującej zależne biblioteki do katalogu /opt/spark/jars oraz włączenie pliku konfiguracyjnego HDFS do SPARK_CLASSPATH w entrypoint.sh.

Druga opcja użycia — Apache Livy

Następnie, gdy zadanie zostało opracowane i wymagane jest przetestowanie uzyskanego wyniku, pojawia się pytanie o uruchomienie go w ramach procesu CI/CD i monitorowanie statusów jego realizacji. Oczywiście można uruchomić je także za pomocą lokalnego wywołania spark-submit, ale to komplikuje infrastrukturę CI/CD, ponieważ wymaga instalacji i konfiguracji Sparka na agentach/runnach serwera CI oraz ustawienia dostępu do API Kubernetes. W tym przypadku wybrano realizację z użyciem Apache Livy jako REST API do uruchamiania zadań Spark, umieszczonego wewnątrz klastra Kubernetes. Dzięki temu można uruchamiać zadania Spark na klastrze Kubernetes za pomocą standardowych zapytań cURL, co jest łatwe do zrealizowania w ramach dowolnego rozwiązania CI, a jego umiejscowienie w klastrze Kubernetes rozwiązuje problem autoryzacji przy interakcji z API Kubernetes.

Uruchamiamy Apache Spark na Kubernetes

Wydzielamy to jako drugą opcję użycia — uruchomienie zadań Spark w ramach procesu CI/CD na klastrze Kubernetes w środowisku testowym.

Kilka słów o Apache Livy — działa jako serwer HTTP, oferujący interfejs webowy oraz RESTful API, umożliwiające zdalne uruchamianie spark-submit z wymaganymi parametrami. Tradycyjnie był dołączany do dystrybucji HDP, ale może być również wdrażany w OKD lub dowolnej innej instalacji Kubernetes za pomocą odpowiedniego manifestu i zestawu obrazów Docker, na przykład tego — github.com/ttauveron/k8s-big-data-experiments/tree/master/livy-spark-2.3. Dla naszego przypadku skonstruowano podobny obraz Docker, zawierający Spark wersji 2.4.5 na podstawie poniższego Dockerfile:

FROM java:8-alpine

ENV SPARK_HOME=\/opt\/spark
ENV LIVY_HOME=\/opt\/livy
ENV HADOOP_CONF_DIR=\/etc\/hadoop\/conf
ENV SPARK_USER=spark

WORKDIR \/opt

RUN apk add --update openssl wget bash && 
    wget -P \/opt https:\/\/downloads.apache.org\/spark\/spark-2.4.5\/spark-2.4.5-bin-hadoop2.7.tgz && 
    tar xvzf spark-2.4.5-bin-hadoop2.7.tgz && 
    rm spark-2.4.5-bin-hadoop2.7.tgz && 
    ln -s \/opt\/spark-2.4.5-bin-hadoop2.7 \/opt\/spark

RUN wget http:\/\/mirror.its.dal.ca\/apache\/incubator\/livy\/0.7.0-incubating\/apache-livy-0.7.0-incubating-bin.zip && 
    unzip apache-livy-0.7.0-incubating-bin.zip && 
    rm apache-livy-0.7.0-incubating-bin.zip && 
    ln -s \/opt\/apache-livy-0.7.0-incubating-bin \/opt\/livy && 
    mkdir \/var\/log\/livy && 
    ln -s \/var\/log\/livy \/opt\/livy\/logs && 
    cp \/opt\/livy\/conf\/log4j.properties.template \/opt\/livy\/conf\/log4j.properties

ADD livy.conf \/opt\/livy\/conf
ADD spark-defaults.conf \/opt\/spark\/conf\/spark-defaults.conf
ADD entrypoint.sh \/entrypoint.sh

ENV PATH="\/opt\/livy\/bin:${PATH}"

EXPOSE 8998

ENTRYPOINT ["\/entrypoint.sh"]
CMD ["livy-server"]

Obraz utworzony może być zbudowany i załadowany do posiadanego przez Ciebie repozytorium Docker, na przykład wewnętrznego repozytorium OKD. W celu jego wdrożenia używany jest następujący manifest ({registry-url} — URL rejestru obrazów Docker, {image-name} — nazwa obrazu Docker, {tag} — tag obrazu Docker, {livy-url} — pożądany URL, pod którym będzie dostępny serwer Livy; manifest „Route” stosuje się w przypadku, gdy dystrybucja Kubernetes to Red Hat OpenShift, w przeciwnym razie używa się odpowiedniego manifestu Ingress lub usługi typu NodePort):

---
apiVersion: apps/v1
kind: Deployment
metadata:
  labels:
    component: livy
  name: livy
spec:
  progressDeadlineSeconds: 600
  replicas: 1
  revisionHistoryLimit: 10
  selector:
    matchLabels:
      component: livy
  strategy:
    rollingUpdate:
      maxSurge: 25%
      maxUnavailable: 25%
    type: RollingUpdate
  template:
    metadata:
      creationTimestamp: null
      labels:
        component: livy
    spec:
      containers:
        - command:
            - livy-server
          env:
            - name: K8S_API_HOST
              value: localhost
            - name: SPARK_KUBERNETES_IMAGE
              value: 'gnut3ll4/spark:v1.0.14'
          image: '{registry-url}/{image-name}:{tag}'
          imagePullPolicy: Always
          name: livy
          ports:
            - containerPort: 8998
              name: livy-rest
              protocol: TCP
          resources: {}
          terminationMessagePath: /dev/termination-log
          terminationMessagePolicy: File
          volumeMounts:
            - mountPath: /var/log/livy
              name: livy-log
            - mountPath: /opt/.livy-sessions/
              name: livy-sessions
            - mountPath: /opt/livy/conf/livy.conf
              name: livy-config
              subPath: livy.conf
            - mountPath: /opt/spark/conf/spark-defaults.conf
              name: spark-config
              subPath: spark-defaults.conf
        - command:
            - /usr/local/bin/kubectl
            - proxy
            - '--port'
            - '8443'
          image: 'gnut3ll4/kubectl-sidecar:latest'
          imagePullPolicy: Always
          name: kubectl
          ports:
            - containerPort: 8443
              name: k8s-api
              protocol: TCP
          resources: {}
          terminationMessagePath: /dev/termination-log
          terminationMessagePolicy: File
      dnsPolicy: ClusterFirst
      restartPolicy: Always
      schedulerName: default-scheduler
      securityContext: {}
      serviceAccount: spark
      serviceAccountName: spark
      terminationGracePeriodSeconds: 30
      volumes:
        - emptyDir: {}
          name: livy-log
        - emptyDir: {}
          name: livy-sessions
        - configMap:
            defaultMode: 420
            items:
              - key: livy.conf
                path: livy.conf
            name: livy-config
          name: livy-config
        - configMap:
            defaultMode: 420
            items:
              - key: spark-defaults.conf
                path: spark-defaults.conf
            name: livy-config
          name: spark-config
---
apiVersion: v1
kind: ConfigMap
metadata:
  name: livy-config
data:
  livy.conf: |-
    livy.spark.deploy-mode=cluster
    livy.file.local-dir-whitelist=/opt/.livy-sessions/
    livy.spark.master=k8s://http://localhost:8443
    livy.server.session.state-retain.sec = 8h
  spark-defaults.conf: 'spark.kubernetes.container.image        "gnut3ll4/spark:v1.0.14"'
---
apiVersion: v1
kind: Service
metadata:
  labels:
    app: livy
  name: livy
spec:
  ports:
    - name: livy-rest
      port: 8998
      protocol: TCP
      targetPort: 8998
  selector:
    component: livy
  sessionAffinity: None
  type: ClusterIP
---
apiVersion: route.openshift.io/v1
kind: Route
metadata:
  labels:
    app: livy
  name: livy
spec:
  host: {livy-url}
  port:
    targetPort: livy-rest
  to:
    kind: Service
    name: livy
    weight: 100
  wildcardPolicy: None

Po jego zastosowaniu i pomyślnym uruchomieniu, graficzny interfejs Livy jest dostępny pod adresem: http://{livy-url}/ui. Dzięki Livy możemy opublikować nasze zadanie Spark, korzystając z zapytania REST, na przykład z Postmana. Przykład kolekcji z zapytaniami został przedstawiony poniżej (w tablicy «args» mogą być przekazywane argumenty konfiguracyjne z zmiennymi niezbędnymi do działania uruchamianego zadania):

{
    "info": {
        "_postman_id": "be135198-d2ff-47b6-a33e-0d27b9dba4c8",
        "name": "Spark Livy",
        "schema": "https://schema.getpostman.com/json/collection/v2.1.0/collection.json"
    },
    "item": [
        {
            "name": "1 Prześlij zadanie z plikiem jar",
            "request": {
                "method": "POST",
                "header": [
                    {
                        "key": "Content-Type",
                        "value": "application/json"
                    }
                ],
                "body": {
                    "mode": "raw",
                    "raw": "{nt\"file\": \"local:///opt/spark/examples/target/scala-2.11/jars/spark-examples_2.11-2.4.5.jar\", nt\"className\": \"org.apache.spark.examples.SparkPi\",nt\"numExecutors\":1,nt\"name\": \"spark-test-1\",nt\"conf\": {ntt\"spark.jars.ivy\": \"/tmp/.ivy\",ntt\"spark.kubernetes.authenticate.driver.serviceAccountName\": \"spark\",ntt\"spark.kubernetes.namespace\": \"{project}\",ntt\"spark.kubernetes.container.image\": \"{docker-registry-url}/{repo}/{image-name}:{tag}\"nt}n}"
                },
                "url": {
                    "raw": "http://{livy-url}/batches",
                    "protocol": "http",
                    "host": [
                        "{livy-url}"
                    ],
                    "path": [
                        "batches"
                    ]
                }
            },
            "response": []
        },
        {
            "name": "2 Prześlij zadanie bez pliku jar",
            "request": {
                "method": "POST",
                "header": [
                    {
                        "key": "Content-Type",
                        "value": "application/json"
                    }
                ],
                "body": {
                    "mode": "raw",
                    "raw": "{nt\"file\": \"hdfs://{host}:{port}/{path-to-file-on-hdfs}\", nt\"className\": \"{class-name}\",nt\"numExecutors\":1,nt\"name\": \"spark-test-2\",nt\"proxyUser\": \"0\",nt\"conf\": {ntt\"spark.jars.ivy\": \"/tmp/.ivy\",ntt\"spark.kubernetes.authenticate.driver.serviceAccountName\": \"spark\",ntt\"spark.kubernetes.namespace\": \"{project}\",ntt\"spark.kubernetes.container.image\": \"{docker-registry-url}/{repo}/{image-name}:{tag}\"nt},nt\"args\": [ntt\"HADOOP_CONF_DIR=/opt/spark/hadoop-conf\",ntt\"MASTER=k8s://https://kubernetes.default.svc:8443\"nt]n}"
                },
                "url": {
                    "raw": "http://{livy-url}/batches",
                    "protocol": "http",
                    "host": [
                        "{livy-url}"
                    ],
                    "path": [
                        "batches"
                    ]
                }
            },
            "response": []
        }
    ],
    "event": [
        {
            "listen": "prerequest",
            "script": {
                "id": "41bea1d0-278c-40c9-ad42-bf2e6268897d",
                "type": "text/javascript",
                "exec": [
                    ""
                ]
            }
        },
        {
            "listen": "test",
            "script": {
                "id": "3cdd7736-a885-4a2d-9668-bd75798f4560",
                "type": "text/javascript",
                "exec": [
                    ""
                ]
            }
        }
    ],
    "protocolProfileBehavior": {}
}

Wykonaj pierwszy żądanie z kolekcji, przejdź do interfejsu OKD i sprawdź, czy zadanie zostało pomyślnie uruchomione — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods. W interfejsie Livy (http://{livy-url}/ui) pojawi się sesja, w ramach której za pomocą API Livy lub interfejsu graficznego można śledzić postęp wykonania zadania i przeglądać logi sesji.

Teraz zaprezentujemy mechanizm działania Livy. W tym celu przeanalizujemy logi kontenera Livy w podzie z serwerem Livy — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods/{livy-pod-name}?tab=logs. Z nich widać, że przy wywołaniu REST API Livy w kontenerze o nazwie „livy” wykonywane jest spark-submit, analogiczne do tego, które stosowaliśmy wcześniej (tutaj {livy-pod-name} to nazwa utworzonego poda z serwerem Livy). W zbiorze przedstawiony jest również drugi żądanie, umożliwiające uruchamianie zadań z zdalnym umiejscowieniem pliku wykonywalnego Spark z pomocą serwera Livy.

Trzeci sposób użycia — Spark Operator

Teraz, gdy zadanie zostało przetestowane, pojawia się pytanie o jego regularne uruchamianie. Natomiast sposobem natywnym do regularnego uruchamiania zadań w klastrze Kubernetes jest encja CronJob, którą można wykorzystać, ale w chwili obecnej znacznie większą popularnością cieszy się wykorzystywanie operatorów do zarządzania aplikacjami w Kubernetes, a dla Sparka dostępny jest wystarczająco dojrzały operator, który jest, między innymi, stosowany w rozwiązaniach na poziomie Enterprise (na przykład Lightbend FastData Platform). Zalecamy jego użycie — obecna stabilna wersja Sparka (2.4.5) ma dość ograniczone możliwości konfiguracji uruchamiania zadań Spark w Kubernetes, natomiast w następnej wersji m.in. (3.0.0) zapowiedziano pełne wsparcie dla Kubernetes, ale data wydania pozostaje nieznana. Spark Operator rekompensuje tę niedogodność, dodając ważne parametry konfiguracyjne (na przykład montowanie ConfigMap z konfiguracją dostępu do Hadoop w podach Spark) oraz możliwość regularnego uruchamiania zadań zgodnie z harmonogramem.

Uruchamiamy Apache Spark na Kubernetes
Wyróżnijmy go jako trzeci sposób użycia — regularne uruchamianie zadań Spark w klastrze Kubernetes w środowisku produkcyjnym.

Spark Operator ma otwarty kod źródłowy i jest rozwijany w ramach Google Cloud Platform — github.com/GoogleCloudPlatform/spark-on-k8s-operator. Można go zainstalować na 3 sposoby:

  1. W ramach instalacji Lightbend FastData Platform/Cloudflow;
  2. Z użyciem Helm:
    helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator
    helm install incubator/sparkoperator --namespace spark-operator
    	

  3. Wykorzystując manifesty z oficjalnego repozytorium (https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/tree/master/manifest). Należy zauważyć, że w skład Cloudflow wchodzi operator z wersją API v1beta1. Jeśli używany jest ten typ instalacji, opisy manifestów aplikacji Spark powinny być oparte na przykładach z tagów w Git z odpowiednią wersją API, na przykład „v1beta1-0.9.0-2.4.0”. Wersję operatora można sprawdzić w opisie CRD wchodzącego w skład operatora w słowniku „versions”:
    oc get crd sparkapplications.sparkoperator.k8s.io -o yaml
    	

Jeśli operator został poprawnie zainstalowany, w odpowiednim projekcie pojawi się aktywny pod z operatorem Spark (na przykład cloudflow-fdp-sparkoperator w przestrzeni Cloudflow dla instalacji Cloudflow) oraz odpowiedni typ zasobów Kubernetes o nazwie „sparkapplications”. Aby zbadać istniejące aplikacje Spark, można użyć następującej komendy:

oc get sparkapplications -n {project}

Aby uruchomić zadania za pomocą Spark Operator, należy wykonać 3 kroki:

  • stworzyć obraz Docker zawierający wszystkie niezbędne biblioteki oraz pliki konfiguracyjne i wykonywalne. W docelowej wersji jest to obraz utworzony na etapie CI/CD i przetestowany w testowym klastrze;
  • opublikować obraz Docker w rejestrze dostępnym z klastra Kubernetes;
  • sporządzić manifest typu „SparkApplication” z opisem uruchamianego zadania. Przykłady manifestów są dostępne w oficjalnym repozytorium (na przykład, github.com/GoogleCloudPlatform/spark-on-k8s-operator/blob/v1beta1-0.9.0-2.4.0/examples/spark-pi.yaml). Warto zwrócić uwagę na istotne aspekty dotyczące manifestu:
    1. w słowniku „apiVersion” powinna być podana wersja API odpowiadająca wersji operatora;
    2. w słowniku „metadata.namespace” powinno być podane przestrzeń nazw, w której uruchomiona zostanie aplikacja;
    3. w słowniku „spec.image” powinien być podany adres utworzonego obrazu Docker w dostępnym rejestrze;
    4. w słowniku „spec.mainClass” powinien być podany klasa zadania Spark, która ma być uruchomiona podczas uruchamiania procesu;
    5. w słowniku „spec.mainApplicationFile” powinien być podany ścieżka do wykonywalnego pliku jar;
    6. w słowniku „spec.sparkVersion” powinna być podana używana wersja Spark;
    7. w słowniku „spec.driver.serviceAccount” powinna być podana usługa konta w ramach odpowiedniej przestrzeni nazw Kubernetes, która będzie używana do uruchomienia aplikacji;
    8. w słowniku „spec.executor” powinno być podane liczba zasobów przydzielonych aplikacji;
    9. W słowniku „spec.volumeMounts” musi być podany lokalny katalog, w którym będą tworzone lokalne pliki zadania Spark.

Przykład tworzenia manifestu (tutaj {spark-service-account} to konto serwisowe wewnątrz klastra Kubernetes do uruchamiania zadań Spark):

apiVersion: "sparkoperator.k8s.io/v1beta1"
kind: SparkApplication
metadata:
  name: spark-pi
  namespace: {project}
spec:
  type: Scala
  mode: cluster
  image: "gcr.io/spark-operator/spark:v2.4.0"
  imagePullPolicy: Always
  mainClass: org.apache.spark.examples.SparkPi
  mainApplicationFile: "local:////opt/spark/examples/jars/spark-examples_2.11-2.4.0.jar"
  sparkVersion: "2.4.0"
  restartPolicy:
    type: Never
  volumes:
    - name: "test-volume"
      hostPath:
        path: "/tmp"
        type: Directory
  driver:
    cores: 0.1
    coreLimit: "200m"
    memory: "512m"
    labels:
      version: 2.4.0
    serviceAccount: {spark-service-account}
    volumeMounts:
      - name: "test-volume"
        mountPath: "/tmp"
  executor:
    cores: 1
    instances: 1
    memory: "512m"
    labels:
      version: 2.4.0
    volumeMounts:
      - name: "test-volume"
        mountPath: "/tmp"

W tym manifeście wskazano konto serwisowe, dla którego przed opublikowaniem manifestu należy utworzyć odpowiednie powiązania ról, zapewniające dostęp wymagany do interakcji aplikacji Spark z API Kubernetes (jeśli to potrzebne). W naszym przypadku aplikacja potrzebuje uprawnień do tworzenia Podów. Utwórzmy wymagane powiązanie ról:

oc adm policy add-role-to-user edit system:serviceaccount:{project}:{spark-service-account} -n {project}

Warto również zauważyć, że w specyfikacji tego manifestu może być podany parametr „hadoopConfigMap”, który pozwala określić ConfigMap z konfiguracją Hadoop bez wcześniejszego umieszczania odpowiedniego pliku w obrazie Docker. Nadaje się również do regularnego uruchamiania zadań — za pomocą parametru „schedule” można określić harmonogram uruchamiania tego zadania.

Po tym zapisujemy nasz manifest w pliku spark-pi.yaml i stosujemy go do naszego klastra Kubernetes:

oc apply -f spark-pi.yaml

W wyniku tego zostanie utworzony obiekt typu „sparkapplications”:

oc get sparkapplications -n {project}
> NAME       AGE
> spark-pi   22h

W tym momencie zostanie utworzony pod z aplikacją, którego status będzie wyświetlany w utworzonym „sparkapplications”. Można go zobaczyć za pomocą następującego polecenia:

oc get sparkapplications spark-pi -o yaml -n {project}

Po zakończeniu zadania POD przejdzie w status „Completed”, który również zostanie zaktualizowany w „sparkapplications”. Logi aplikacji można zobaczyć w przeglądarce lub za pomocą następującego polecenia (tutaj {sparkapplications-pod-name} to nazwa podu uruchomionego zadania):

oc logs {sparkapplications-pod-name} -n {project}

Zarządzanie zadaniami Spark można także przeprowadzić za pomocą specjalistycznego narzędzia sparkctl. Aby je zainstalować, klonujemy repozytorium z jego kodem źródłowym, instalujemy Go i kompilujemy to narzędzie:

git clone https://github.com/GoogleCloudPlatform/spark-on-k8s-operator.git
cd spark-on-k8s-operator/
wget https://dl.google.com/go/go1.13.3.linux-amd64.tar.gz
tar -xzf go1.13.3.linux-amd64.tar.gz
sudo mv go /usr/local
mkdir $HOME/Projects
export GOROOT=/usr/local/go
export GOPATH=$HOME/Projects
export PATH=$GOPATH/bin:$GOROOT/bin:$PATH
go -version
cd sparkctl
go build -o sparkctl
sudo mv sparkctl /usr/local/bin

Przyjrzyjmy się liście uruchomionych zadań Spark:

sparkctl list -n {project}

Utwórzmy opis dla zadania Spark:

vi spark-app.yaml

apiVersion: "sparkoperator.k8s.io/v1beta1"
kind: SparkApplication
metadata:
  name: spark-pi
  namespace: {project}
spec:
  type: Scala
  mode: cluster
  image: "gcr.io/spark-operator/spark:v2.4.0"
  imagePullPolicy: Always
  mainClass: org.apache.spark.examples.SparkPi
  mainApplicationFile: "local:///opt/spark/examples/jars/spark-examples_2.11-2.4.0.jar"
  sparkVersion: "2.4.0"
  restartPolicy:
    type: Never
  volumes:
    - name: "test-volume"
      hostPath:
        path: "/tmp"
        type: Directory
  driver:
    cores: 1
    coreLimit: "1000m"
    memory: "512m"
    labels:
      version: 2.4.0
    serviceAccount: spark
    volumeMounts:
      - name: "test-volume"
        mountPath: "/tmp"
  executor:
    cores: 1
    instances: 1
    memory: "512m"
    labels:
      version: 2.4.0
    volumeMounts:
      - name: "test-volume"
        mountPath: "/tmp"

Uruchommy opisane zadanie za pomocą sparkctl:

sparkctl create spark-app.yaml -n {project}

Przyjrzyjmy się liście uruchomionych zadań Spark:

sparkctl list -n {project}

Przyjrzyjmy się liście zdarzeń uruchomionego zadania Spark:

sparkctl event spark-pi -n {project} -f

Zbadajmy status uruchomionego zadania Spark:

sparkctl status spark-pi -n {project}

Na zakończenie warto przyjrzeć się wykrytym wadom eksploatacyjnym aktualnej stabilnej wersji Spark (2.4.5) w Kubernetes:

  1. Pierwsza i chyba najważniejsza wada to brak lokalności danych. Pomimo wszystkich niedociągnięć YARN, były też korzyści z jego użycia, na przykład zasada dostarczania kodu do danych (a nie danych do kodu). Dzięki temu zadania Spark były wykonywane na węzłach, na których znajdowały się dane biorące udział w obliczeniach, co znacznie skracało czas potrzebny na przesył danych przez sieć. Przy użyciu Kubernetes mamy do czynienia z koniecznością przesuwania przez sieć danych, które są wykorzystywane w zadaniu. Jeśli są one wystarczająco duże, czas wykonania zadania może znacząco wzrosnąć, a także wymagać sporej ilości przestrzeni dyskowej, przypisanej instancjom zadania Spark na ich tymczasowe przechowywanie. Ten niedostatek można ograniczyć poprzez wykorzystanie specjalistycznych narzędzi, które zapewniają lokalność danych w Kubernetes (na przykład Alluxio), ale to w praktyce oznacza konieczność przechowywania pełnej kopii danych na węzłach klastra Kubernetes.
  2. Drugą istotną wadą jest bezpieczeństwo. Domyślnie funkcje związane z bezpieczeństwem w kontekście uruchamiania zadań Spark są wyłączone, a opcja użycia Kerberosa nie jest omówiona w oficjalnej dokumentacji (choć odpowiednie parametry pojawiły się w wersji 3.0.0, co wymaga dodatkowej obróbki), a w dokumentacji dotyczącej bezpieczeństwa podczas korzystania ze Sparka (https://spark.apache.org/docs/2.4.5/security.html) jako magazyny kluczy wymieniane są tylko YARN, Mesos i Standalone Cluster. Przy tym użytkownik, pod którym uruchamiane są zadania Spark, nie może być podany bezpośrednio — możemy jedynie określić konto serwisowe, pod którym będzie działać pod, a użytkownik jest wybierany na podstawie skonfigurowanych polityk bezpieczeństwa. W związku z tym używa się albo użytkownika root, co nie jest bezpieczne w środowisku produkcyjnym, albo użytkownika z losowym UID, co jest niewygodne przy przydzielaniu praw dostępu do danych (można to rozwiązać przez stworzenie PodSecurityPolicies i ich przypisanie do odpowiednich kont serwisowych). Na obecną chwilę problem ten rozwiązuje się poprzez umieszczenie wszystkich niezbędnych plików bezpośrednio w obrazie Dockera lub modyfikację skryptu uruchamiania Sparka w celu wykorzystania mechanizmu przechowywania i uzyskiwania sekretów, przyjętego w Państwa organizacji.
  3. Uruchamianie zadań Spark za pomocą Kubernetes wciąż pozostaje w fazie eksperymentalnej, a w przyszłości mogą wystąpić znaczące zmiany w używanych artefaktach (plikach konfiguracyjnych, bazowych obrazach Docker i skryptach uruchamiających). Rzeczywiście — podczas przygotowywania materiału przetestowano wersje 2.3.0 i 2.4.5, a ich zachowanie różniło się znacznie.

Czekamy na aktualizacje — niedawno pojawiła się nowa wersja Spark (3.0.0), która przyniosła zauważalne zmiany w działaniu Sparka na Kubernetes, ale zachowała eksperymentalny status wsparcia dla tego menedżera zasobów. Możliwe, że następne aktualizacje rzeczywiście pozwolą całkowicie zrezygnować z YARN i uruchamiać zadania Spark na Kubernetes, nie obawiając się o bezpieczeństwo systemu i bez konieczności samodzielnego dostosowywania funkcjonalnych komponentów.

Fin.

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster