Avviamo Apache Spark su Kubernetes

Cari lettori, buon pomeriggio. Oggi parleremo un po' di Apache Spark e delle sue prospettive di sviluppo.

Avviamo Apache Spark su Kubernetes

Nel mondo moderno, Big Data e Apache Spark sono di fatto standard per lo sviluppo di applicazioni di elaborazione dati batch. Inoltre, viene utilizzato per creare applicazioni di streaming che operano secondo il concetto di micro batch, elaborando e trasferendo i dati in piccole porzioni (Spark Structured Streaming). Tradizionalmente, ha fatto parte del comune stack Hadoop, utilizzando YARN come gestore delle risorse (o, in alcuni casi, Apache Mesos). Entro il 2020, l'utilizzo tradizionale di Hadoop da parte della maggior parte delle aziende è stato messo in discussione a causa della mancanza di adeguate distribuzioni di Hadoop: lo sviluppo di HDP e CDH è stato arrestato, CDH è poco ottimizzato e costoso, e gli altri fornitori di Hadoop hanno o cessato di esistere o hanno un futuro incerto. Pertanto, cresce l'interesse nella comunità e tra le grandi aziende per avviare Apache Spark utilizzando Kubernetes: diventando lo standard per l'orchestrazione dei container e la gestione delle risorse in ambienti cloud privati e pubblici, risolve il problema della complessità nell'allocazione delle risorse per i task Spark su YARN e offre una piattaforma in costante sviluppo con molteplici distribuzioni commerciali e open-source per aziende di ogni dimensione e settore. Inoltre, sulla scia della popolarità, la maggior parte ha già installato un paio di versioni e ha sviluppato competenze nel suo utilizzo, facilitando la transizione.

A partire dalla versione 2.3.0, Apache Spark ha acquisito supporto ufficiale per l'esecuzione di attività nel cluster Kubernetes e oggi discuteremo della maturità attuale di questo approccio, delle varie opzioni di utilizzo e delle insidie che si possono incontrare durante l'implementazione.

Iniziamo esaminando il processo di sviluppo di lavori e applicazioni basati su Apache Spark e sottolineiamo i casi tipici in cui è necessario eseguire un'attività su un cluster Kubernetes. Per la preparazione di questo post, utilizziamo OpenShift come distribuzione e forniremo comandi rilevanti per il suo strumento da riga di comando (oc). Per altre distribuzioni Kubernetes, possono essere utilizzati i comandi corrispondenti dello strumento standard da riga di comando Kubernetes (kubectl) o le loro alternative (ad esempio, per oc adm policy).

La prima opzione di utilizzo — spark-submit

Nello sviluppo di attività e applicazioni, lo sviluppatore ha bisogno di eseguire task per il debug della trasformazione dei dati. In teoria, per questi scopi possono essere utilizzati dei mock, ma lo sviluppo che coinvolge istanze reali (anche se di test) dei sistemi finali si è dimostrato più rapido e di qualità superiore in questa classe di attività. Nel caso in cui eseguiamo il debug su istanze reali dei sistemi finali, sono possibili due scenari operativi:

  • lo sviluppatore esegue un job Spark localmente in modalità standalone;

    Avviamo Apache Spark su Kubernetes

  • lo sviluppatore esegue un job Spark su un cluster Kubernetes in un ambiente di test.

    Avviamo Apache Spark su Kubernetes

La prima opzione è valida, ma comporta diversi svantaggi:

  • ogni sviluppatore deve garantire l'accesso dal proprio posto di lavoro a tutte le istanze finali necessarie;
  • la macchina di lavoro deve disporre di risorse sufficienti per eseguire il job in sviluppo.

La seconda opzione è priva di questi difetti, poiché l'uso di un cluster Kubernetes consente di riservare il pool di risorse necessario per l'esecuzione dei task e di fornire gli accessi richiesti agli istanze dei sistemi finali, offrendo un accesso flessibile tramite il modello di ruoli di Kubernetes a tutti i membri del team di sviluppo. Individuiamolo come prima opzione da utilizzare: esecuzione di task Spark dalla macchina locale dello sviluppatore su un cluster Kubernetes in un ambiente di test.

Approfondiamo il processo di configurazione di Spark per l'esecuzione locale. Per iniziare a utilizzare Spark, è necessario installarlo:

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

Raccogliamo i pacchetti necessari per lavorare con Kubernetes:

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

La compilazione completa richiede molto tempo, e per creare le immagini Docker e avviarle nel cluster Kubernetes sono necessari solo i file jar dalla directory "assembly/", quindi possiamo compilare solo questo sotto-progetto:

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

Per avviare compiti Spark in Kubernetes, è necessario creare un'immagine Docker che verrà utilizzata come base. Sono possibili 2 approcci:

  • L'immagine Docker creata include il codice eseguibile del compito Spark;
  • L'immagine creata include solo Spark e le dipendenze necessarie, mentre il codice eseguibile è posizionato remotamente (ad esempio, in HDFS).

Cominciamo a costruire un'immagine Docker contenente un esempio di compito Spark. Per la creazione delle immagini Docker, Spark ha uno strumento appropriato chiamato «docker-image-tool». Esaminiamo la sua guida:

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

Con esso puoi creare immagini Docker e caricarle in registri remoti, ma per impostazione predefinita presenta alcuni svantaggi:

  • crea necessariamente 3 immagini Docker — per Spark, PySpark e R;
  • non consente di specificare il nome dell'immagine.

Pertanto, utilizzeremo una versione modificata di questo strumento, riportata di seguito:

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

Con questo, creiamo un'immagine base di Spark, contenente un compito di test per calcolare il numero Pi utilizzando Spark (qui {docker-registry-url} è l'URL del tuo registro Docker, {repo} è il nome del repository all'interno del registro, che corrisponde al progetto in OpenShift, {image-name} è il nome dell'immagine - se si utilizza una suddivisione a tre livelli delle immagini, ad esempio come nel registro integrato di Red Hat OpenShift, {tag} è il tag della versione di questa immagine):

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

Autenticati nel cluster OKD utilizzando l'utility della console (qui {OKD-API-URL} è l'URL dell'API del cluster OKD):

oc login {OKD-API-URL}

Otteniamo il token dell'utente corrente per l'autenticazione nel Docker Registry:

oc whoami -t

Autenticati nel Docker Registry interno del cluster OKD (utilizza il token ottenuto con il comando precedente come password):

docker login {docker-registry-url}

Carichiamo l'immagine Docker costruita nel Docker Registry OKD:

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

Verifichiamo che l'immagine creata sia disponibile in OKD. Per fare ciò, apriamo nel browser l'URL con l'elenco delle immagini del progetto corrispondente (qui {project} è il nome del progetto all'interno del cluster OpenShift, {OKD-WEBUI-URL} è l'URL della console Web di OpenShift) — https://{OKD-WEBUI-URL}/console/project/{project}/browse/images/{image-name}.

Per avviare i compiti deve essere creato un account di servizio con privilegi per l'esecuzione dei pod come root (ne discuteremo più avanti):

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

Eseguiamo il comando spark-submit per pubblicare il compito Spark nel cluster OKD, specificando l'account di servizio creato e l'immagine 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

Qui:

—name — il nome del compito che verrà utilizzato per formare i nomi dei pod Kubernetes;

—class — la classe del file eseguibile richiamato all'avvio del compito;

—conf — parametri di configurazione di Spark;

spark.executor.instances — il numero di esecutori Spark da avviare;

spark.kubernetes.authenticate.driver.serviceAccountName — il nome dell'account di servizio Kubernetes utilizzato per avviare i pod (per determinare il contesto di sicurezza e le capacità durante l'interazione con l'API Kubernetes);

spark.kubernetes.namespace — lo spazio dei nomi Kubernetes in cui verranno avviati i pod del driver e degli esecutori;

spark.submit.deployMode — modo di avvio di Spark (per il tradizionale spark-submit si utilizza «cluster», per Spark Operator e versioni più recenti di Spark si utilizza «client»);

spark.kubernetes.container.image — immagine Docker utilizzata per avviare i pod;

spark.master — URL dell'API Kubernetes (indicato esternamente per accedere da una macchina locale);

local:// — percorso del file eseguibile Spark all'interno dell'immagine Docker.

Accediamo al progetto OKD corrispondente e analizziamo i pod creati — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods.

Per semplificare il processo di sviluppo, è possibile utilizzare un'altra opzione, in cui si crea un'immagine base condivisa di Spark, utilizzata da tutti i compiti per l'esecuzione, e gli snapshot dei file eseguibili vengono pubblicati in un archivio esterno (ad esempio, Hadoop) e specificati durante la chiamata a spark-submit come link. In questo caso, si possono eseguire diverse versioni dei compiti Spark senza ricompilare le immagini Docker, utilizzando, ad esempio, WebHDFS per la pubblicazione delle immagini. Inviamo una richiesta per creare un file (qui {host} è l'host del servizio WebHDFS, {port} è la porta del servizio WebHDFS, {path-to-file-on-hdfs} è il percorso desiderato del file su HDFS):

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

Il risultato sarà del tipo (dove {location} è l'URL da utilizzare per caricare il file):

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

Caricando il file eseguibile Spark su HDFS (dove {path-to-local-file} è il percorso del file eseguibile Spark sull'host attuale):

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

Dopo ciò, possiamo eseguire spark-submit utilizzando il file Spark caricato su HDFS (dove {class-name} è il nome della classe che deve essere eseguita per completare il compito):

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

Si deve notare che per accedere a HDFS e garantire il funzionamento del compito, potrebbe essere necessario modificare il Dockerfile e lo script entrypoint.sh — aggiungere al Dockerfile una direttiva per copiare le librerie dipendenti nella directory /opt/spark/jars e includere il file di configurazione HDFS in SPARK_CLASSPATH in entrypoint.sh.

La seconda opzione è utilizzare Apache Livy

Successivamente, quando un'attività è stata sviluppata e si rende necessario testare il risultato ottenuto, sorgono domande sul suo avvio all'interno del processo CI/CD e sul monitoraggio degli stati di esecuzione. Certamente, è possibile eseguirla anche tramite una chiamata locale a spark-submit, ma questo complica l'infrastruttura CI/CD poiché richiede l'installazione e la configurazione di Spark sugli agenti/runners del server CI e la configurazione dell'accesso all'API di Kubernetes. Per questo caso, la soluzione scelta è l'utilizzo di Apache Livy come API REST per l'avvio delle attività Spark, collocata all'interno del cluster Kubernetes. Con Livy, è possibile avviare attività Spark sul cluster Kubernetes utilizzando normali richieste cURL, facilmente implementabili su qualsiasi soluzione CI, e la sua collocazione all'interno del cluster Kubernetes risolve il problema dell'autenticazione durante l'interazione con l'API di Kubernetes.

Avviamo Apache Spark su Kubernetes

Identifichiamolo come la seconda opzione di utilizzo: avvio delle attività Spark all'interno del processo CI/CD nel cluster Kubernetes in un ambiente di test.

Un po' su Apache Livy: funge da server HTTP, fornendo un'interfaccia Web e un'API RESTful che consente di eseguire remotamente spark-submit, passando i parametri necessari. Tradizionalmente, veniva fornito con la distribuzione HDP, ma può essere distribuito anche in OKD o in qualsiasi altra installazione di Kubernetes utilizzando il manifesto appropriato e un set di immagini Docker, ad esempio questo — github.com/ttauveron/k8s-big-data-experiments/tree/master/livy-spark-2.3. Per il nostro caso, è stata costruita un'immagine Docker simile, che include Spark versione 2.4.5 dal seguente Dockerfile:

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

L'immagine creata può essere assemblata e caricata nel tuo repository Docker esistente, ad esempio, il repository interno di OKD. Per il suo deployment si utilizza il seguente manifesto ({registry-url} — URL del registro delle immagini Docker, {image-name} — nome dell'immagine Docker, {tag} — tag dell'immagine Docker, {livy-url} — URL desiderato per l'accesso al server Livy; il manifesto "Route" viene applicato se viene utilizzata la distribuzione Kubernetes Red Hat OpenShift, altrimenti viene utilizzato il manifesto Ingress o il servizio di tipo NodePort corrispondente):

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

Dopo l'applicazione e l'avvio riuscito del pod, l'interfaccia grafica di Livy è accessibile al link: http://{livy-url}/ui. Con Livy, possiamo pubblicare il nostro job Spark utilizzando una richiesta REST, ad esempio, da Postman. Un esempio di collezione con le richieste è mostrato qui sotto (nel array «args» possono essere passati argomenti di configurazione con variabili necessarie per l'esecuzione del job avviato):

{
    "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 Invia lavoro con 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 Invia lavoro senza 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": {}
}

Eseguiamo la prima richiesta dalla collezione, accediamo all'interfaccia OKD e verifichiamo che il task sia stato avviato con successo — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods. In questo modo, nell'interfaccia di Livy (http://{livy-url}/ui) verrà creata una sessione, all'interno della quale è possibile monitorare il progresso dell'attività e visualizzare i log attraverso l'API Livy o l'interfaccia grafica.

Ora mostriamo come funziona Livy. A tal fine, esploriamo i log del contenitore di Livy all'interno del pod con il server Livy — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods/{livy-pod-name}?tab=logs. Da essi si evince che, quando si chiama l'API REST di Livy, viene eseguito spark-submit nel contenitore chiamato "livy", simile a quello utilizzato in precedenza (dove {livy-pod-name} è il nome del pod creato con il server Livy). Nella collezione è presente anche una seconda richiesta, che consente di avviare attività con la posizione remota del file eseguibile Spark attraverso il server Livy.

La terza opzione di utilizzo — Spark Operator

Ora che il compito è stato testato, sorge la questione della sua esecuzione regolare. Un modo nativo per eseguire compiti in modo regolare in un cluster Kubernetes è tramite l'entità CronJob, e può essere utilizzato, ma attualmente stanno guadagnando grande popolarità gli operatori per la gestione delle applicazioni in Kubernetes. Esiste un operatore sufficientemente maturo anche per Spark, che viene utilizzato in soluzioni di livello enterprise (ad esempio, Lightbend FastData Platform). Consigliamo di utilizzarlo: l'attuale versione stabile di Spark (2.4.5) ha capacità di configurazione molto limitate per l'esecuzione dei compiti Spark in Kubernetes, mentre nella prossima versione principale (3.0.0) è annunciato il supporto completo per Kubernetes, ma la data di uscita rimane sconosciuta. L'operatore Spark compensa questa mancanza, aggiungendo parametri di configurazione importanti (come il montaggio di ConfigMap con la configurazione di accesso a Hadoop nei pod Spark) e la possibilità di eseguire i compiti regolarmente secondo un piano.

Avviamo Apache Spark su Kubernetes
Sottolineiamo questo come il terzo caso d'uso: esecuzione regolare dei compiti Spark in un cluster Kubernetes in un ambiente produttivo.

Lo Spark Operator è open source e sviluppato all'interno della Google Cloud Platform — github.com/GoogleCloudPlatform/spark-on-k8s-operator. Può essere installato in 3 modi:

  1. Nell'ambito dell'installazione della Lightbend FastData Platform/Cloudflow;
  2. Utilizzando Helm:
    helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator
    helm install incubator/sparkoperator --namespace spark-operator
    	

  3. Applicando i manifesti dal repository ufficiale (https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/tree/master/manifest). È importante notare che Cloudflow include un operatore con la versione API v1beta1. Se viene utilizzato questo tipo di installazione, le descrizioni dei manifesti delle applicazioni Spark devono essere basate sugli esempi dei tag in Git con la versione API corrispondente, ad esempio «v1beta1-0.9.0-2.4.0». La versione dell'operatore può essere visualizzata nella descrizione del CRD incluso nell'operatore nel dizionario «versions»:
    oc get crd sparkapplications.sparkoperator.k8s.io -o yaml
    	

Se l'operatore è stato installato correttamente, nel progetto corrispondente apparirà un pod attivo con l'operatore Spark (ad esempio, cloudflow-fdp-sparkoperator nello spazio Cloudflow per l'installazione di Cloudflow) e apparirà il corrispondente tipo di risorsa Kubernetes chiamato «sparkapplications». Le applicazioni Spark disponibili possono essere esplorate con il seguente comando:

oc get sparkapplications -n {project}

Per avviare i lavori utilizzando Spark Operator è necessario eseguire 3 operazioni:

  • creare un'immagine Docker che includa tutte le librerie necessarie, così come i file di configurazione ed esecuzione. Nella configurazione prevista, si tratta di un'immagine creata nella fase CI/CD e testata su un cluster di prova;
  • pubblicare l'immagine Docker in un registro accessibile dal cluster Kubernetes;
  • formare un manifesto con il tipo "SparkApplication" e la descrizione del lavoro da avviare. Esempi di manifesti sono disponibili nel repository ufficiale (per esempio, github.com/GoogleCloudPlatform/spark-on-k8s-operator/blob/v1beta1-0.9.0-2.4.0/examples/spark-pi.yaml). È importante notare alcuni punti riguardanti il manifesto:
    1. nel dizionario "apiVersion" deve essere specificata la versione API corrispondente alla versione dell'operatore;
    2. nel dizionario "metadata.namespace" deve essere specificato lo spazio dei nomi in cui verrà eseguita l'applicazione;
    3. nel dizionario "spec.image" deve essere fornito l'indirizzo dell'immagine Docker creata nel registro accessibile;
    4. nel dizionario "spec.mainClass" deve essere specificato il nome della classe Spark che deve essere avviata durante l'esecuzione del processo;
    5. nel dizionario "spec.mainApplicationFile" deve essere fornito il percorso al file jar eseguibile;
    6. nel dizionario «spec.sparkVersion» deve essere specificata la versione di Spark utilizzata;
    7. nel dizionario «spec.driver.serviceAccount» deve essere indicato l'account di servizio all'interno dello spazio dei nomi Kubernetes corrispondente, che sarà utilizzato per eseguire l'applicazione;
    8. nel dizionario «spec.executor» deve essere specificato il numero di risorse allocate per l'applicazione;
    9. nel dizionario «spec.volumeMounts» deve essere specificata la directory locale in cui verranno creati i file locali del compito Spark.

Esempio di manifestazione (qui {spark-service-account} è l'account di servizio all'interno del cluster Kubernetes per eseguire i compiti 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"

In questo manifesto viene specificato un account di servizio per il quale è necessario creare i collegamenti ai ruoli richiesti prima della pubblicazione del manifesto, fornendo i necessari diritti per l'interazione dell'applicazione Spark con l'API Kubernetes (se necessario). Nel nostro caso, l'applicazione ha bisogno dei diritti per creare Pod. Creiamo il collegamento al ruolo necessario:

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

Vale anche la pena notare che nella specifica di questo manifesto può essere specificato il parametro "hadoopConfigMap", che consente di indicare un ConfigMap con la configurazione di Hadoop senza la necessità di inserire preventivamente il file corrispondente nell'immagine Docker. È anche adatto per l'esecuzione regolare dei compiti: con il parametro "schedule" può essere specificato un programma di esecuzione di questo compito.

Dopo di che salviamo il nostro manifesto nel file spark-pi.yaml e lo applichiamo al nostro cluster Kubernetes:

oc apply -f spark-pi.yaml

Verrà quindi creato un oggetto di tipo "sparkapplications":

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

Verrà creato un pod con l'applicazione, il cui stato sarà visualizzato nel "sparkapplications" creato. Può essere visualizzato con il seguente comando:

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

Al termine del task, il POD passerà allo stato "Completato", che verrà anche aggiornato in "sparkapplications". I log dell'applicazione possono essere visualizzati nel browser o utilizzando il seguente comando (dove {sparkapplications-pod-name} è il nome del pod del task in esecuzione):

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

La gestione dei task Spark può essere effettuata anche tramite un'utilità specializzata chiamata sparkctl. Per installarla, cloniamo il repository del suo codice sorgente, installiamo Go e compiliamo l'utilità:

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

Esaminiamo l'elenco dei task Spark in esecuzione:

sparkctl list -n {project}

Creiamo una descrizione per il task 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"

Eseguiamo il compito descritto utilizzando sparkctl:

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

Esaminiamo l'elenco dei task Spark in esecuzione:

sparkctl list -n {project}

Esaminiamo l'elenco degli eventi del task Spark avviato:

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

Esaminiamo lo stato del task Spark avviato:

sparkctl status spark-pi -n {project}

In conclusione, vorremmo esaminare gli svantaggi rilevati nell'uso della versione stabile attuale di Spark (2.4.5) su Kubernetes:

  1. Il primo e, probabilmente, il principale svantaggio è l'assenza di Località dei Dati. Nonostante tutti i difetti, YARN aveva anche vantaggi, come il principio della consegna del codice ai dati (anziché dei dati al codice). Grazie a questo, le attività di Spark venivano eseguite sui nodi dove si trovavano i dati coinvolti nei calcoli, riducendo notevolmente il tempo di trasferimento dei dati attraverso la rete. Con Kubernetes, ci confrontiamo con la necessità di trasferire i dati utilizzati nel lavoro delle attività attraverso la rete. Se questi sono di dimensioni considerevoli, il tempo di esecuzione dell'attività può aumentare notevolmente, così come potrebbe essere necessario un volume significativo di spazio su disco assegnato alle istanze dell'attività Spark per il loro storage temporaneo. Questo svantaggio può essere mitigato attraverso l'uso di strumenti software specializzati che assicurano la località dei dati in Kubernetes (ad esempio, Alluxio), ma ciò implica di fatto la necessità di mantenere una copia completa dei dati sui nodi del cluster Kubernetes.
  2. Il secondo importante svantaggio è la sicurezza. Le funzionalità legate alla sicurezza relative all'esecuzione dei task Spark sono disattivate per impostazione predefinita, e l'uso di Kerberos non è trattato nella documentazione ufficiale (anche se le relative opzioni sono state introdotte nella versione 3.0.0, richiedendo ulteriori elaborazioni). Nella documentazione sulla sicurezza nell'uso di Spark (https://spark.apache.org/docs/2.4.5/security.html) come archivi delle chiavi sono elencati solo YARN, Mesos e Standalone Cluster. Inoltre, l'utente con cui vengono eseguiti i task Spark non può essere specificato direttamente: possiamo solo definire un account di servizio sotto il quale opererà, e l'utente è scelto in base alle politiche di sicurezza configurate. Pertanto, si utilizza l'utente root, il che non è sicuro in un ambiente di produzione, oppure un utente con un UID casuale, il che complica la gestione dei diritti di accesso ai dati (risolvibile creando PodSecurityPolicies e collegandole agli account di servizio corrispondenti). Attualmente, la soluzione è quella di inserire tutti i file necessari direttamente nell'immagine Docker, oppure modificare lo script di avvio di Spark per utilizzare il meccanismo di archiviazione e recupero dei segreti adottato nella vostra organizzazione.
  3. L'esecuzione di attività Spark tramite Kubernetes è ancora ufficialmente in fase sperimentale e in futuro ci potrebbero essere cambiamenti significativi negli artefatti utilizzati (file di configurazione, immagini Docker di base e script di avvio). Infatti, durante la preparazione del materiale sono state testate le versioni 2.3.0 e 2.4.5, il cui comportamento differisce notevolmente.

Attendiamo aggiornamenti: è recentemente stata rilasciata una nuova versione di Spark (3.0.0), che ha apportato notevoli modifiche al funzionamento di Spark su Kubernetes, mantenendo però lo stato sperimentale del supporto per questo gestore delle risorse. È possibile che i prossimi aggiornamenti consentano di raccomandare completamente di abbandonare YARN e di eseguire le attività Spark su Kubernetes, senza compromettere la sicurezza del sistema e senza la necessità di sviluppare autonomamente componenti funzionali.

Fine.

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster