Apache Spark auf Kubernetes starten

Liebe Leser, einen guten Tag. Heute sprechen wir ein wenig über Apache Spark und seine Entwicklungsperspektiven.

Apache Spark auf Kubernetes starten

In der heutigen Welt der Big Data ist Apache Spark de facto der Standard für die Entwicklung von Datenverarbeitungsjobs. Darüber hinaus wird es auch zur Erstellung von Streaming-Anwendungen verwendet, die im Mikro-Batch-Konzept arbeiten und Daten in kleinen Portionen verarbeiten und ausliefern (Spark Structured Streaming). Traditionell war es Teil des Hadoop-Ökosystems und nutzte als Ressourcenscheduler YARN (oder in einigen Fällen Apache Mesos). Bis 2020 steht die Verwendung von Hadoop in seiner traditionellen Form für die meisten Unternehmen jedoch in Zweifel, da es an vernünftigen Distributionen mangelt – die Entwicklung von HDP und CDH wurde eingestellt, CDH ist unzureichend ausgearbeitet und hat hohe Kosten, während andere Hadoop-Anbieter entweder ihren Betrieb eingestellt haben oder eine ungewisse Zukunft haben. Daher gewinnt der Start von Apache Spark über Kubernetes – als Standard für die Orchestrierung von Containern und Ressourcenzuteilung in privaten und öffentlichen Clouds – zunehmend an Interesse innerhalb der Community und bei großen Unternehmen, da er das Problem der umständlichen Ressourcenzuteilung von Spark-Jobs unter YARN löst und eine stabil wachsende Plattform mit einer Vielzahl von kommerziellen und Open-Source-Distributionen für Unternehmen jeder Größe und Branche bietet. Darüber hinaus haben die meisten Unternehmen aufgrund der wachsenden Beliebtheit bereits ein paar ihrer Installationen eingerichtet und ihre Expertise in der Nutzung weiter ausgebaut, was den Umstieg erleichtert.

Seit Version 2.3.0 bietet Apache Spark offizielle Unterstützung für die Ausführung von Aufgaben in einem Kubernetes-Cluster. Heute werden wir über die aktuelle Reife dieses Ansatzes, verschiedene Nutzungsmöglichkeiten und potenzielle Fallstricke sprechen, die bei der Implementierung zu erwarten sind.

Zunächst betrachten wir den Entwicklungsprozess von Aufgaben und Anwendungen auf Basis von Apache Spark und heben typische Szenarien hervor, in denen es notwendig ist, eine Aufgabe in einem Kubernetes-Cluster auszuführen. Für diesen Beitrag wird OpenShift als Distribution verwendet, und es werden Befehle relevant für dessen Befehlszeilenwerkzeug (oc) angegeben. Für andere Kubernetes-Distributionen können die entsprechenden Befehle des Standard-Kubernetes-Befehlszeilenwerkzeugs (kubectl) oder deren Äquivalente (z.B. für oc adm policy) verwendet werden.

Die erste Nutzungsoption ist spark-submit.

Bei der Entwicklung von Aufgaben und Anwendungen müssen Entwickler Aufgaben zur Fehlerbehebung von Datentransformationen ausführen. Theoretisch könnten Platzhalter für diese Zwecke verwendet werden, jedoch hat sich die Arbeit mit realen (auch wenn es Testexemplare sind) Endsystemen in dieser Art von Aufgaben als schneller und qualitativ hochwertiger erwiesen. Wenn wir die Fehlerbehebung an realen Exemplaren der Endsysteme durchführen, können zwei Arbeitsabläufe möglich sein:

  • der Entwickler führt die Spark-Aufgabe lokal im Standalone-Modus aus;

    Apache Spark auf Kubernetes starten

  • der Entwickler führt die Spark-Aufgabe in einem Test-Cluster auf Kubernetes aus.

    Apache Spark auf Kubernetes starten

Die erste Variante hat ihre Berechtigung, bringt jedoch einige Nachteile mit sich:

  • Für jeden Entwickler muss der Zugang von seinem Arbeitsplatz zu allen erforderlichen Endsystemexemplaren sichergestellt werden;
  • Auf dem Arbeitsrechner müssen genügend Ressourcen vorhanden sein, um die zu entwickelnde Aufgabe auszuführen.

Die zweite Option ist frei von diesen Nachteilen, da die Verwendung eines Kubernetes-Clusters es ermöglicht, den erforderlichen Pool von Ressourcen für die Ausführung von Aufgaben bereitzustellen und den notwendigen Zugriff auf die Instanzen der Endsysteme sicherzustellen. Gleichzeitig gewährt die rollenbasierte Zugriffskontrolle von Kubernetes allen Mitgliedern des Entwicklungsteams flexiblen Zugang. Wir heben sie als erste Nutzungsmöglichkeit hervor: das Ausführen von Spark-Aufgaben von der lokalen Maschine des Entwicklers auf einem Kubernetes-Cluster in einer Testumgebung.

Lassen Sie uns den Prozess der Konfiguration von Spark für den lokalen Start näher erläutern. Um Spark nutzen zu können, muss es installiert werden:

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

Jetzt sammeln wir die erforderlichen Pakete für die Arbeit mit Kubernetes:

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

Der vollständige Build benötigt viel Zeit, und für die Erstellung von Docker-Images und deren Ausführung auf einem Kubernetes-Cluster sind in Wirklichkeit nur die JAR-Dateien aus dem Verzeichnis „assembly/“ erforderlich. Daher kann nur dieses Teilprojekt gebaut werden:

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

Für den Einsatz von Spark-Jobs in Kubernetes muss ein Docker-Image erstellt werden, das als Basis dient. Es gibt zwei Ansätze dafür:

  • Das erstellte Docker-Image enthält den ausführbaren Code des Spark-Jobs;
  • Das erstellte Image enthält nur Spark und die erforderlichen Abhängigkeiten; der ausführbare Code wird remote platziert (z. B. in HDFS).

Zuerst erstellen wir ein Docker-Image, das ein Testbeispiel für einen Spark-Job enthält. Zum Erstellen von Docker-Images hat Spark ein entsprechendes Tool namens „docker-image-tool“. Lassen Sie uns die Hilfe dazu betrachten:

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

Mit diesem Tool können Docker-Images erstellt und in entfernte Registries hochgeladen werden, jedoch hat es standardmäßig mehrere Nachteile:

  • Es erstellt unbedingt gleichzeitig 3 Docker-Images – für Spark, PySpark und R;
  • Es erlaubt nicht, den Namen des Images anzugeben.

Deshalb werden wir eine modifizierte Version dieses Tools verwenden, die wie folgt aussieht:

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

Damit erstellen wir ein Basismodell von Spark, das eine Testaufgabe zur Berechnung der Zahl Pi mit Spark enthält (hierbei ist {docker-registry-url} die URL Ihres Docker-Image-Registers, {repo} der Name des Repositories im Register, das mit dem Projekt in OpenShift übereinstimmt, {image-name} der Name des Images [falls eine dreistufige Trennung der Images verwendet wird, wie im integrierten Image-Register von Red Hat OpenShift], und {tag} der Tag dieser Version des Images):

./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

Wir melden uns beim OKD-Cluster über das Konsolen-Utility an (hierbei ist {OKD-API-URL} die URL der OKD-Cluster-API):

oc login {OKD-API-URL}

Wir erhalten das Token des aktuellen Benutzers zur Authentifizierung im Docker Registry:

oc whoami -t

Wir melden uns beim internen Docker Registry des OKD-Clusters an (als Passwort verwenden wir das über den vorherigen Befehl erhaltene Token):

docker login {docker-registry-url}

Wir laden das gebaute Docker-Image in das OKD Docker Registry hoch:

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

Wir überprüfen, ob das erstellte Image in OKD verfügbar ist. Dazu öffnen wir im Browser die URL mit der Liste der Images des entsprechenden Projekts (hier ist {project} der Name des Projekts im OpenShift-Cluster, {OKD-WEBUI-URL} ist die URL der OpenShift-Webkonsole) — https://{OKD-WEBUI-URL}/console/project/{project}/browse/images/{image-name}.

Um Aufgaben auszuführen, muss ein Dienstkonto mit den Berechtigungen zum Starten von Pods unter root erstellt werden (dieses Thema werden wir später besprechen):

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

Wir führen den Befehl spark-submit aus, um eine Spark-Aufgabe im OKD-Cluster zu veröffentlichen, indem wir das erstellte Dienstkonto und das Docker-Image angeben:

 /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

Hier:

—name — der Name der Aufgabe, der bei der Namebildung der Kubernetes-Pods verwendet wird;

—class — die Klasse der ausführbaren Datei, die beim Starten der Aufgabe aufgerufen wird;

—conf — die Konfigurationsparameter für Spark;

spark.executor.instances — die Anzahl der gestarteten Spark-Executors;

spark.kubernetes.authenticate.driver.serviceAccountName — der Name des Kubernetes-Dienstkontos, das beim Starten von Pods verwendet wird (zur Bestimmung des Sicherheitskontexts und der Berechtigungen bei der Interaktion mit der Kubernetes-API);

spark.kubernetes.namespace — der Kubernetes-Namespace, in dem die Pods des Treibers und der Executors gestartet werden.

spark.submit.deployMode — die Art und Weise, wie Spark gestartet wird (für das Standard-spark-submit wird «cluster» verwendet, für Spark Operator und spätere Versionen von Spark «client»);

spark.kubernetes.container.image — das Docker-Image, das zum Starten der Pods verwendet wird;

spark.master — die URL API von Kubernetes (außerhalb angegeben, damit die Verbindung von der lokalen Maschine erfolgt);

local:// — der Pfad zur ausführbaren Datei von Spark innerhalb des Docker-Images.

Gehen Sie zu dem entsprechenden OKD-Projekt und überprüfen Sie die erstellten Pods — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods.

Zur Vereinfachung des Entwicklungsprozesses kann eine weitere Variante verwendet werden, bei der ein gemeinsames Basis-Image von Spark erstellt wird, das von allen Jobs zum Starten verwendet wird, und die Snapshots der ausführbaren Dateien in einem externen Speicher (z. B. Hadoop) veröffentlicht werden und beim Aufruf von spark-submit als Link angegeben werden. In diesem Fall können verschiedene Versionen von Spark-Jobs ohne Neuaufbau der Docker-Images gestartet werden, indem beispielsweise WebHDFS zum Veröffentlichen der Images verwendet wird. Senden Sie eine Anfrage zum Erstellen der Datei (hier {host} — der Host des WebHDFS-Dienstes, {port} — der Port des WebHDFS-Dienstes, {path-to-file-on-hdfs} — der gewünschte Pfad zur Datei auf HDFS):

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

Dabei erhalten Sie eine Antwort in der Form (hier ist {location} die URL, die für den Datei-Upload verwendet werden soll):

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

Sie laden die ausführbare Datei Spark in HDFS hoch (hier ist {path-to-local-file} der Pfad zur ausführbaren Spark-Datei auf dem aktuellen Host):

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

Nach diesem Schritt können wir spark-submit mit der auf HDFS hochgeladenen Spark-Datei durchführen (hier ist {class-name} der Name der Klasse, die für die Ausführung der Aufgabe gestartet werden soll):

/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}

Es ist darauf hinzuweisen, dass möglicherweise Änderungen an der Dockerfile und dem Skript entrypoint.sh erforderlich sind, um auf HDFS zuzugreifen und die Aufgabe auszuführen. Fügen Sie eine Direktive zum Kopieren der erforderlichen Bibliotheken in das Verzeichnis /opt/spark/jars in die Dockerfile ein und stellen Sie die HDFS-Konfigurationsdatei in SPARK_CLASSPATH in entrypoint.sh bereit.

Eine zweite Verwendungsmöglichkeit ist Apache Livy.

Wenn die Aufgabe entwickelt ist und das Ergebnis getestet werden soll, stellt sich die Frage nach der Ausführung im Rahmen des CI/CD-Prozesses und der Überwachung des Status der Ausführung. Natürlich kann man sie auch durch einen lokalen Aufruf von spark-submit starten, aber das kompliziert die CI/CD-Infrastruktur, da es die Installation und Konfiguration von Spark auf den CI-Server-Agenten/-Runners und die Einrichtung des Zugangs zur Kubernetes-API erfordert. Für diesen Fall wurde entschieden, Apache Livy als REST-API zum Starten von Spark-Jobs, die im Kubernetes-Cluster gehostet werden, zu verwenden. Damit können Spark-Jobs im Kubernetes-Cluster mit gewöhnlichen cURL-Anfragen gestartet werden, was sich leicht in jede CI-Lösung integrieren lässt. Die Platzierung innerhalb des Kubernetes-Clusters löst das Authentifizierungsproblem bei der Interaktion mit der Kubernetes-API.

Apache Spark auf Kubernetes starten

Lassen Sie uns dies als zweite Verwendung hervorheben – den Start von Spark-Jobs im Rahmen des CI/CD-Prozesses auf dem Kubernetes-Cluster in einem Testumfeld.

Einige Details zu Apache Livy – es fungiert als HTTP-Server und bietet ein Web-Interface sowie ein RESTful API, mit dem Sie spark-submit remote ausführen können, indem Sie die erforderlichen Parameter übergeben. Traditionell wurde es als Teil der HDP-Distribution bereitgestellt, kann jedoch auch in OKD oder jeder anderen Kubernetes-Installation mit dem entsprechenden Manifest und einer Reihe von Docker-Images, wie diesem, bereitgestellt – github.com/ttauveron/k8s-big-data-experiments/tree/master/livy-spark-2.3. Für unseren Fall wurde ein ähnliches Docker-Image erstellt, das Spark in der Version 2.4.5 aus dem folgenden Dockerfile enthält:

VON 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"]

Das erstellte Image kann in Ihr vorhandenes Docker-Repository hochgeladen werden, zum Beispiel in das interne OKD-Repository. Zu seiner Bereitstellung verwenden Sie das folgende Manifest ({registry-url} — URL des Docker-Registrys, {image-name} — Name des Docker-Images, {tag} — Tag des Docker-Images, {livy-url} — die gewünschte URL, unter der der Livy-Server erreichbar sein wird; das Manifest 'Route' wird verwendet, wenn Red Hat OpenShift als Kubernetes-Distrubution eingesetzt wird, andernfalls wird das entsprechende Ingress- oder NodePort-Service-Manifest verwendet):

---
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

Nach der Anwendung und dem erfolgreichen Start des Pods ist die grafische Benutzeroberfläche von Livy unter dem Link: http://{livy-url}/ui erreichbar. Mit Livy können wir unsere Spark-Aufgabe über eine REST-Anfrage veröffentlichen, beispielsweise aus Postman. Ein Beispiel für die Sammlung von Anfragen ist unten dargestellt (im Array „args“ können Konfigurationsargumente mit Variablen übergeben werden, die für die Ausführung der gestarteten Aufgabe erforderlich sind):

{
    "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 Job mit Jar einreichen",
            "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 Job ohne Jar einreichen",
            "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": {}
}

Wir führen die erste Anfrage aus der Sammlung aus, gehen zum OKD-Interface und überprüfen, ob die Aufgabe erfolgreich gestartet wurde — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods. Gleichzeitig wird im Livy-Interface (http://{livy-url}/ui) eine Sitzung angezeigt, in der die Ausführung der Aufgabe über die Livy-API oder die grafische Benutzeroberfläche verfolgt und die Sitzungsprotokolle eingesehen werden können.

Jetzt zeigen wir, wie Livy funktioniert. Dazu betrachten wir die Protokolle des Livy-Containers innerhalb des Pods mit dem Livy-Server — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods/{livy-pod-name}?tab=logs. Daraus ist ersichtlich, dass beim Aufruf der REST-API von Livy im Container mit dem Namen "livy" der spark-submit-Befehl ausgeführt wird, der dem oben verwendeten entspricht (hier ist {livy-pod-name} der Name des erstellten Pods mit dem Livy-Server). In der Sammlung ist auch eine zweite Anfrage enthalten, die es ermöglicht, Aufgaben mit remote gehosteten Spark-Ausführungsdateien über den Livy-Server zu starten.

Eine dritte Verwendungsmöglichkeit — Spark Operator

Jetzt, da die Aufgabe getestet wurde, stellt sich die Frage nach ihrem regelmäßigen Start. Die native Methode für den regelmäßigen Start von Aufgaben im Kubernetes-Cluster ist die Entität CronJob, die verwendet werden kann. Momentan erfreuen sich jedoch Operatoren zur Verwaltung von Anwendungen in Kubernetes großer Beliebtheit. Für Spark gibt es einen ausreichend ausgereiften Operator, der auch in Enterprise-Lösungen (wie der Lightbend FastData Platform) genutzt wird. Wir empfehlen, diesen zu verwenden – die aktuelle stabile Version von Spark (2.4.5) bietet nur eingeschränkte Möglichkeiten zur Konfiguration der Ausführung von Spark-Aufgaben in Kubernetes, während in der nächsten Hauptversion (3.0.0) die vollständige Unterstützung für Kubernetes angekündigt ist, der Veröffentlichungstermin jedoch unbekannt bleibt. Der Spark Operator kompensiert diese Einschränkungen, indem er wichtige Konfigurationsparameter hinzufügt (wie das Einbinden eines ConfigMap mit Zugangskonfiguration zu Hadoop in die Spark-Pods) und die Möglichkeit, Aufgaben nach einem Zeitplan regelmäßig auszuführen.

Apache Spark auf Kubernetes starten
Lassen Sie uns dies als dritte Option hervorheben – die regelmäßige Ausführung von Spark-Aufgaben in einem produktiven Kubernetes-Cluster.

Der Spark Operator ist Open Source und wird im Rahmen der Google Cloud Platform entwickelt — github.com/GoogleCloudPlatform/spark-on-k8s-operator. Die Installation kann auf drei Arten erfolgen:

  1. Im Rahmen der Installation von Lightbend FastData Platform/Cloudflow;
  2. Mit Helm:
    helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator
    helm install incubator/sparkoperator --namespace spark-operator
    	

  3. Durch Verwendung von Manifests aus dem offiziellen Repository (https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/tree/master/manifest). Es ist zu beachten, dass Cloudflow einen Operator mit API-Version v1beta1 enthält. Wenn dieser Installationstyp verwendet wird, sollten die Beschreibungen der Spark-Anwendungsmanifeste auf Basis der Beispiele aus den Git-Tags mit der entsprechenden API-Version, wie „v1beta1-0.9.0-2.4.0“, basieren. Die Version des Operators kann in der Beschreibung der CRD im 'versions'-Dictionary des Operators eingesehen werden:
    oc get crd sparkapplications.sparkoperator.k8s.io -o yaml
    	

Wenn der Operator korrekt installiert ist, erscheint im entsprechenden Projekt ein aktives Pod mit dem Spark Operator (z.B. cloudflow-fdp-sparkoperator im Cloudflow-Namespace für die Cloudflow-Installation) und ein entsprechender Kubernetes-Ressourcentyp mit dem Namen „sparkapplications“. Die vorhandenen Spark-Anwendungen können mit dem folgenden Befehl untersucht werden:

oc get sparkapplications -n {project}

Um Aufgaben mit dem Spark Operator auszuführen, sind drei Schritte erforderlich:

  • Erstellen Sie ein Docker-Image, das alle erforderlichen Bibliotheken sowie Konfigurations- und ausführbare Dateien enthält. Im Zielbild handelt es sich um ein Bild, das im CI/CD-Prozess erstellt und auf einem Testcluster getestet wurde;
  • Veröffentlichen Sie das Docker-Image in einem Registry, das aus dem Kubernetes-Cluster erreichbar ist;
  • Erstellen Sie ein Manifest des Typs „SparkApplication“ mit der Beschreibung der auszuführenden Aufgabe. Beispiele für Manifeste finden Sie im offiziellen Repository (zum Beispiel, github.com/GoogleCloudPlatform/spark-on-k8s-operator/blob/v1beta1-0.9.0-2.4.0/examples/spark-pi.yaml). Es gibt einige wichtige Punkte zu beachten, die das Manifest betreffen:
    1. Im Dictionary „apiVersion“ muss die API-Version angegeben werden, die der Version des Operators entspricht;
    2. Im Dictionary „metadata.namespace“ muss der Namespace angegeben werden, in dem die Anwendung ausgeführt wird;
    3. Im Dictionary „spec.image“ muss die Adresse des erstellten Docker-Images im verfügbaren Registry angegeben werden;
    4. Im Dictionary „spec.mainClass“ muss die Spark-Klasse angegeben werden, die beim Start des Prozesses ausgeführt werden soll;
    5. Im Dictionary „spec.mainApplicationFile“ muss der Pfad zur ausführbaren Jar-Datei angegeben werden;
    6. Im Wörterbuch «spec.sparkVersion» muss die verwendete Spark-Version angegeben werden;
    7. Im Wörterbuch «spec.driver.serviceAccount» muss das Servicekonto innerhalb des entsprechenden Kubernetes-Namensraums angegeben werden, das zum Starten der Anwendung verwendet wird;
    8. Im Wörterbuch «spec.executor» muss die Anzahl der Ressourcen, die der Anwendung zugewiesen werden, angegeben werden;
    9. Im Wörterbuch «spec.volumeMounts» muss das lokale Verzeichnis angegeben werden, in dem die lokalen Dateien der Spark-Aufgabe erstellt werden.

Beispiel für die Erstellung eines Manifests (hier ist {spark-service-account} das Servicekonto innerhalb des Kubernetes-Clusters für das Starten von Spark-Aufgaben):

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"

In diesem Manifest ist das Dienstkonto angegeben, für das vor der Veröffentlichung des Manifests die erforderlichen Rollenzuweisungen erstellt werden müssen, um den notwendigen Zugriff für die Interaktion der Spark-Anwendung mit der Kubernetes-API zu gewähren (falls erforderlich). In unserem Fall benötigt die Anwendung das Recht, Pods zu erstellen. Lassen Sie uns die notwendige Rollenzuweisung erstellen:

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

Es ist auch erwähnenswert, dass in den Spezifikationen dieses Manifests das Parameter "hadoopConfigMap" angegeben sein kann, das es ermöglicht, eine ConfigMap mit der Hadoop-Konfiguration anzugeben, ohne dass die entsprechende Datei vorher in das Docker-Image eingefügt werden muss. Es ist auch für regelmäßige Aufgaben geeignet – mit dem Parameter "schedule" kann ein Zeitplan für die Ausführung dieser Aufgabe angegeben werden.

Danach speichern wir unser Manifest in der Datei spark-pi.yaml und wenden es auf unseren Kubernetes-Cluster an:

oc apply -f spark-pi.yaml

Dabei wird ein Objekt vom Typ „sparkapplications“ erstellt:

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

Es wird ein Pod mit der Anwendung erstellt, dessen Status im erstellten „sparkapplications“ angezeigt wird. Dies kann mit dem folgenden Befehl überprüft werden:

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

Nach Abschluss der Aufgabe wechselt der POD in den Status „Completed“, der auch in „sparkapplications“ aktualisiert wird. Die Anwendungsprotokolle können im Browser oder mit dem folgenden Befehl eingesehen werden (hierbei ist {sparkapplications-pod-name} der Name des Pods der ausgeführten Aufgabe):

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

Das Management von Spark-Aufgaben kann auch über das spezielle Tool sparkctl erfolgen. Um es zu installieren, klonen wir das Repository mit dem Quellcode, installieren Go und bauen dieses Tool:

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

Lassen Sie uns die Liste der ausgeführten Spark-Aufgaben ansehen:

sparkctl list -n {project}

Erstellen wir eine Beschreibung für die Spark-Aufgabe:

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"

Führen wir die beschriebene Aufgabe mit sparkctl aus:

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

Lassen Sie uns die Liste der ausgeführten Spark-Aufgaben ansehen:

sparkctl list -n {project}

Überprüfen wir die Liste der Ereignisse der gestarteten Spark-Aufgabe:

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

Überprüfen wir den Status der gestarteten Spark-Aufgabe:

sparkctl status spark-pi -n {project}

Abschließend möchten wir die erkannten Nachteile des Betriebs der aktuellen stabilen Version von Spark (2.4.5) in Kubernetes betrachten:

  1. Der erste und wohl größte Nachteil ist das Fehlen der Datenlokalität. Trotz aller Schwächen von YARN gab es auch Vorteile bei seiner Nutzung, wie das Prinzip der Codebereitstellung zu den Daten (nicht die Daten zu dem Code). Dadurch wurden Spark-Jobs auf den Knoten ausgeführt, an denen sich die für die Berechnungen verwendeten Daten befanden, was die Datenübertragungszeit im Netzwerk deutlich reduzierte. Bei der Verwendung von Kubernetes sehen wir uns jedoch der Notwendigkeit gegenüber, Daten über das Netzwerk zu bewegen, die für die Ausführung der Aufgabe benötigt werden. Wenn diese Daten ausreichend groß sind, kann die Laufzeit der Aufgabe erheblich zunehmen, und es kann auch ein erheblicher Menge an Speicherplatz erforderlich sein, der den Spark-Task-Instanzen für die temporäre Speicherung zugewiesen wird. Dieser Nachteil kann durch den Einsatz spezialisierter Software verringert werden, die die Datenlokalität in Kubernetes sicherstellt (beispielsweise Alluxio). Dies bedeutet jedoch faktisch, dass eine vollständige Kopie der Daten auf den Knoten des Kubernetes-Clusters gespeichert werden muss.
  2. Ein weiterer wesentlicher Nachteil ist die Sicherheit. Standardmäßig sind die sicherheitsrelevanten Funktionen zur Ausführung von Spark-Aufgaben deaktiviert; die Verwendung von Kerberos wird in der offiziellen Dokumentation nicht behandelt (obwohl entsprechende Optionen mit Version 3.0.0 verfügbar wurden, was zusätzliche Implementierungen erfordert). In der Sicherheitsdokumentation für die Verwendung von Spark (https://spark.apache.org/docs/2.4.5/security.html) werden als Schlüsselspeicher nur YARN, Mesos und Standalone Cluster genannt. Außerdem kann der Benutzer, unter dem die Spark-Jobs ausgeführt werden, nicht direkt angegeben werden; wir legen lediglich ein Dienstkonto fest, unter dem der Pod läuft, während der Benutzer basierend auf den konfigurierten Sicherheitsrichtlinien ausgewählt wird. Daher wird entweder der Benutzer root verwendet, was in Produktionsumgebungen unsicher ist, oder ein Benutzer mit einer zufälligen UID, was die Zuweisung von Zugriffsrechten auf Daten erschwert (dies lässt sich durch die Erstellung von PodSecurityPolicies und deren Bindung an die entsprechenden Dienstkonten lösen). Derzeit wird entweder versucht, alle erforderlichen Dateien direkt in das Docker-Image zu integrieren, oder das Startskript von Spark wird modifiziert, um den in Ihrer Organisation akzeptierten Mechanismus zur Speicherung und Abfrage von Geheimnissen zu nutzen.
  3. Die Ausführung von Spark-Aufgaben mit Kubernetes befindet sich offiziell noch im Experimentiermodus, und in Zukunft könnten erhebliche Änderungen an den verwendeten Artefakten (Konfigurationsdateien, Basis-Docker-Images und Startskripten) vorgenommen werden. Tatsächlich wurde bei der Vorbereitung des Materials mit den Versionen 2.3.0 und 2.4.5 getestet, und das Verhalten wies erhebliche Unterschiede auf.

Wir erwarten Aktualisierungen – kürzlich wurde die neueste Version von Spark (3.0.0) veröffentlicht, die spürbare Änderungen an der Funktionalität von Spark auf Kubernetes mit sich brachte, dabei jedoch den experimentellen Status des Supports für diesen Ressourcenmanager beibehielt. Möglicherweise werden zukünftige Updates tatsächlich dazu führen, dass wir empfehlen können, YARN abzulehnen und Spark-Aufgaben auf Kubernetes auszuführen, ohne sich um die Sicherheit Ihres Systems sorgen zu müssen und ohne die Notwendigkeit, funktionale Komponenten selbst anzupassen.

Fin.

Quelle: habr.com

Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen 🔥 Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen | ProHoster