Iniciando Apache Spark en Kubernetes

Estimados lectores, buenos días. Hoy hablaremos un poco sobre Apache Spark y sus perspectivas de desarrollo.

Iniciando Apache Spark en Kubernetes

En el mundo actual de Big Data, Apache Spark es el estándar de facto para desarrollar tareas de procesamiento por lotes. Además, se utiliza para crear aplicaciones de streaming que operan bajo el concepto de micro lotes, procesando y entregando datos en pequeñas porciones (Spark Structured Streaming). Tradicionalmente, también ha sido parte de la pila general de Hadoop, utilizando YARN como gestor de recursos (o, en algunos casos, Apache Mesos). Para el 2020, su uso en su forma tradicional por la mayoría de las empresas se encuentra en gran duda debido a la falta de distribuciones adecuadas de Hadoop: el desarrollo de HDP y CDH ha cesado, CDH está insuficientemente desarrollado y tiene un alto costo, mientras que los demás proveedores de Hadoop han dejado de existir o tienen un futuro incierto. Por lo tanto, la comunidad y las grandes empresas están cada vez más interesadas en ejecutar Apache Spark a través de Kubernetes, que se ha establecido como el estándar en la orquestación de contenedores y la gestión de recursos en nubes privadas y públicas, resolviendo el problema de la planificación incómoda de recursos de las tareas de Spark en YARN y ofreciendo una plataforma en constante evolución con múltiples distribuciones comerciales y de código abierto para empresas de todos los tamaños y sectores. Además, con la ola de popularidad, la mayoría ya ha logrado establecer un par de sus instalaciones y aumentar su experiencia en su uso, lo que facilita la migración.

Desde la versión 2.3.0, Apache Spark cuenta con soporte oficial para ejecutar tareas en clústeres de Kubernetes y hoy hablaremos sobre la madurez actual de este enfoque, las diversas formas de utilizarlo y los desafíos que se presentarán al implementarlo.

Primero, examinaremos el proceso de desarrollo de tareas y aplicaciones basadas en Apache Spark y destacaremos los casos típicos en los que es necesario ejecutar una tarea en un clúster de Kubernetes. Para la preparación de este post, se utiliza OpenShift como distribución y se proporcionarán los comandos relevantes para su herramienta de línea de comandos (oc). Para otras distribuciones de Kubernetes, se pueden usar los comandos correspondientes de la herramienta de línea de comandos estándar de Kubernetes (kubectl) o sus análogos (por ejemplo, para oc adm policy).

La primera opción de uso: spark-submit

Durante el desarrollo de tareas y aplicaciones, el desarrollador necesita ejecutar tareas para depurar la transformación de datos. Teóricamente, se pueden usar simulaciones para estos fines, pero el desarrollo con instancias reales (aunque sea de prueba) de sistemas finales ha demostrado ser más rápido y de mayor calidad en este tipo de tareas. En los casos en los que estamos depurando con instancias reales de sistemas finales, se pueden dar dos escenarios de trabajo:

  • el desarrollador ejecuta la tarea Spark localmente en modo standalone;

    Iniciando Apache Spark en Kubernetes

  • el desarrollador ejecuta la tarea Spark en un clúster de Kubernetes en un entorno de prueba.

    Iniciando Apache Spark en Kubernetes

La primera opción tiene derecho a existir, pero conlleva una serie de desventajas:

  • cada desarrollador necesita asegurarse de tener acceso desde su lugar de trabajo a todas las instancias necesarias de los sistemas finales;
  • la máquina de trabajo requiere una cantidad suficiente de recursos para ejecutar la tarea en desarrollo.

La segunda opción carece de estas desventajas, ya que el uso de un clúster de Kubernetes permite asignar el grupo de recursos necesario para ejecutar tareas y garantizar el acceso necesario a las instancias de los sistemas finales, proporcionando acceso flexible mediante el modelo de roles de Kubernetes a todos los miembros del equipo de desarrollo. Destacaremos esto como la primera opción de uso: ejecutar tareas Spark desde la máquina local del desarrollador en un clúster de Kubernetes en un entorno de prueba.

Hablemos en detalle sobre el proceso de configuración de Spark para ejecución local. Para comenzar a usar Spark, es necesario instalarlo:

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

Reunimos los paquetes necesarios para trabajar con Kubernetes:

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

La compilación completa lleva mucho tiempo, y para crear imágenes Docker y ejecutarlas en un clúster Kubernetes, en realidad solo se necesitan los archivos jar del directorio «assembly/», por lo que se puede compilar solo este subproyecto:

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

Para ejecutar tareas Spark en Kubernetes es necesario crear una imagen Docker que se utilizará como base. Aquí hay 2 enfoques posibles:

  • La imagen Docker creada incluye el código ejecutable de la tarea Spark;
  • La imagen creada incluye solo Spark y las dependencias necesarias, el código ejecutable se ubica de forma remota (por ejemplo, en HDFS).

Primero, crearemos una imagen Docker que contenga un ejemplo de prueba de una tarea Spark. Para crear imágenes Docker, Spark tiene una utilidad correspondiente llamada «docker-image-tool». Revisemos la ayuda de esta herramienta:

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

Con ella, se pueden crear imágenes Docker y cargarlas en registros remotos, pero por defecto tiene varias desventajas:

  • obligatoriamente crea 3 imágenes Docker — para Spark, PySpark y R;
  • no permite especificar el nombre de la imagen.

Por lo tanto, utilizaremos una versión modificada de esta herramienta, que se muestra a continuación:

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 ella, generamos la imagen base de Spark, que contiene una tarea de prueba para calcular el número Pi usando Spark (aquí {docker-registry-url} — es la URL de su registro de imágenes Docker, {repo} — es el nombre del repositorio dentro del registro, que coincide con el proyecto en OpenShift, {image-name} — es el nombre de la imagen (si se utiliza una división de imágenes de tres niveles, por ejemplo, como en el registro de imágenes integrado de Red Hat OpenShift), {tag} — es la etiqueta de esta versión de la imagen):

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

Iniciamos sesión en el clúster OKD utilizando la herramienta de consola (aquí {OKD-API-URL} — es la URL de la API del clúster OKD):

oc login {OKD-API-URL}

Obtenemos el token del usuario actual para la autenticación en Docker Registry:

oc whoami -t

Iniciamos sesión en el Docker Registry interno del clúster OKD (utilizamos el token obtenido mediante el comando anterior como contraseña):

docker login {docker-registry-url}

Cargamos la imagen Docker compilada en Docker Registry OKD:

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

Verifiquemos que la imagen recopilada esté disponible en OKD. Para ello, abramos en el navegador la URL de la lista de imágenes del proyecto correspondiente (aquí {project} es el nombre del proyecto dentro del clúster de OpenShift, {OKD-WEBUI-URL} es la URL de la consola web de OpenShift) — https://{OKD-WEBUI-URL}/console/project/{project}/browse/images/{image-name}.

Para ejecutar tareas, debe crearse una cuenta de servicio con privilegios para iniciar pods como root (discutiremos este tema más adelante):

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

Ejecutaremos el comando spark-submit para publicar una tarea Spark en el clúster de OKD, especificando la cuenta de servicio creada y la imagen de 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

Aquí:

—name — el nombre de la tarea que participará en la formación del nombre de los pods de Kubernetes;

—class — la clase del archivo ejecutable que se invoca al iniciar la tarea;

—conf — parámetros de configuración de Spark;

spark.executor.instances — el número de ejecutores de Spark que se van a iniciar;

spark.kubernetes.authenticate.driver.serviceAccountName — el nombre de la cuenta de servicio de Kubernetes utilizada al iniciar pods (para determinar el contexto de seguridad y capacidades al interactuar con la API de Kubernetes);

spark.kubernetes.namespace — el espacio de nombres de Kubernetes en el que se ejecutarán los pods del controlador y de los ejecutores;

spark.submit.deployMode — el modo de despliegue de Spark (para el spark-submit estándar se utiliza "cluster", para el Spark Operator y versiones posteriores de Spark se usa "client");

spark.kubernetes.container.image — la imagen de Docker utilizada para ejecutar los pods;

spark.master — URL de la API de Kubernetes (se indica la externa, ya que la conexión se realiza desde la máquina local);

local:// — la ruta al archivo ejecutable de Spark dentro de la imagen de Docker.

Vamos al proyecto correspondiente en OKD y revisamos los pods creados — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods.

Para simplificar el proceso de desarrollo, se puede usar otra opción en la que se crea una imagen base común de Spark, utilizada por todas las tareas para ejecutar, y los snapshots de los archivos ejecutables se publican en un almacenamiento externo (por ejemplo, Hadoop) y se indican al llamar a spark-submit como un enlace. En este caso, se pueden ejecutar diferentes versiones de tareas Spark sin reconstruir imágenes de Docker, usando, por ejemplo, WebHDFS para publicar las imágenes. Enviamos una solicitud para crear un archivo (aquí {host} es el host del servicio WebHDFS, {port} es el puerto del servicio WebHDFS, {path-to-file-on-hdfs} es la ruta deseada del archivo en HDFS):

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

Se obtendrá una respuesta del tipo (aquí {location} es la URL que debe usarse para cargar el archivo):

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

Cargando el archivo ejecutable de Spark en HDFS (aquí {path-to-local-file} es la ruta al archivo ejecutable de Spark en el host actual):

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

Después de esto, podemos ejecutar spark-submit utilizando el archivo de Spark que se ha cargado en HDFS (aquí {class-name} es el nombre de la clase que se requiere ejecutar para realizar la tarea):

/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 importante mencionar que para acceder a HDFS y asegurar el funcionamiento de la tarea, puede ser necesario modificar el Dockerfile y el script entrypoint.sh: agregar en el Dockerfile una directiva para copiar las bibliotecas dependientes al directorio /opt/spark/jars e incluir el archivo de configuración de HDFS en SPARK_CLASSPATH en entrypoint.sh.

La segunda opción de uso es Apache Livy

A continuación, cuando la tarea está desarrollada y es necesario probar el resultado obtenido, surge la cuestión de su ejecución dentro del proceso CI/CD y el seguimiento de los estados de su ejecución. Por supuesto, se puede ejecutar localmente usando spark-submit, pero esto complica la infraestructura CI/CD ya que requiere la instalación y configuración de Spark en los agentes/corrunners del servidor CI y la configuración de acceso a la API de Kubernetes. Para este caso, se ha elegido implementar Apache Livy como API REST para ejecutar tareas de Spark, alojado dentro del clúster de Kubernetes. Con esto, se pueden ejecutar tareas de Spark en el clúster de Kubernetes utilizando solicitudes cURL comunes, lo cual es fácilmente realizable en cualquier solución de CI, y su colocación dentro del clúster de Kubernetes resuelve el problema de autenticación al interactuar con la API de Kubernetes.

Iniciando Apache Spark en Kubernetes

Destacamos esto como la segunda opción de uso: ejecutar tareas de Spark dentro del proceso CI/CD en el clúster de Kubernetes en un entorno de prueba.

Un poco sobre Apache Livy: funciona como un servidor HTTP, que proporciona una interfaz web y una API RESTful, permitiendo ejecutar spark-submit de forma remota, pasando los parámetros necesarios. Tradicionalmente, se suministraba como parte de la distribución HDP, pero también puede desplegarse en OKD o cualquier otra instalación de Kubernetes utilizando el manifiesto correspondiente y un conjunto de imágenes Docker, por ejemplo, este — github.com/ttauveron/k8s-big-data-experiments/tree/master/livy-spark-2.3. Para nuestro caso se construyó una imagen Docker similar, que incluye Spark versión 2.4.5 del siguiente Dockerfile:

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

La imagen creada puede ser construida y cargada en su repositorio Docker existente, por ejemplo, el repositorio interno de OKD. Para su despliegue se utiliza el siguiente manifiesto ({registry-url} — URL del registro de imágenes Docker, {image-name} — nombre de la imagen Docker, {tag} — etiqueta de la imagen Docker, {livy-url} — URL deseada donde estará disponible el servidor Livy; el manifiesto "Route" se utiliza en caso de que se use Red Hat OpenShift como distribución de Kubernetes, de lo contrario se utiliza el manifiesto correspondiente de Ingress o Service de tipo 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

Después de su aplicación y el inicio exitoso del pod, la interfaz gráfica de Livy está disponible en el enlace: http://{livy-url}/ui. Con Livy, podemos publicar nuestra tarea Spark utilizando una solicitud REST, por ejemplo, desde Postman. A continuación se presenta un ejemplo de colección con solicitudes (en el array 'args' se pueden enviar argumentos de configuración con variables necesarias para la ejecución de la tarea lanzada):

{
    "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 Enviar trabajo 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 Enviar trabajo sin 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": {}
}

Realizaremos la primera solicitud de la colección, iremos a la interfaz de OKD y verificaremos que la tarea se haya iniciado correctamente — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods. Mientras tanto, en la interfaz de Livy (http://{livy-url}/ui) aparecerá una sesión, dentro de la cual se puede rastrear el progreso de la tarea utilizando la API de Livy o la interfaz gráfica y revisar los registros de la sesión.

Ahora mostraremos el mecanismo de funcionamiento de Livy. Para ello, estudiaremos los registros del contenedor de Livy dentro del pod con el servidor de Livy — https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods/{livy-pod-name}?tab=logs. De ellos se puede ver que al invocar la API REST de Livy en el contenedor llamado "livy", se ejecuta spark-submit, similar al que usamos anteriormente (aquí {livy-pod-name} es el nombre del pod creado con el servidor de Livy). La colección también presenta una segunda solicitud que permite ejecutar tareas con la ubicación remota del archivo ejecutable de Spark a través del servidor de Livy.

La tercera opción de uso es el Spark Operator.

Ahora que la tarea ha sido probada, surge la pregunta de su ejecución regular. La forma nativa para ejecutar tareas de manera regular en un clúster de Kubernetes es a través de la entidad CronJob y se puede utilizar, pero en este momento se ha vuelto más popular el uso de operadores para gestionar aplicaciones en Kubernetes, y para Spark existe un operador bastante maduro que se utiliza en soluciones de nivel empresarial (por ejemplo, Lightbend FastData Platform). Recomendamos utilizarlo: la versión estable actual de Spark (2.4.5) tiene capacidades bastante limitadas para la configuración de la ejecución de tareas de Spark en Kubernetes, mientras que en la próxima versión importante (3.0.0) se promete soporte completo para Kubernetes, aunque la fecha de lanzamiento sigue siendo desconocida. Spark Operator compensa esta limitación al agregar parámetros importantes de configuración (por ejemplo, el montaje de ConfigMap con la configuración de acceso a Hadoop en los pods de Spark) y la posibilidad de ejecutar regularmente tareas según un horario.

Iniciando Apache Spark en Kubernetes
Lo destacaremos como la tercera opción de uso: ejecución regular de tareas de Spark en el clúster de Kubernetes en un entorno de producción.

El Spark Operator tiene código abierto y se desarrolla en el marco de Google Cloud Platform — github.com/GoogleCloudPlatform/spark-on-k8s-operator. Su instalación puede llevarse a cabo de 3 maneras:

  1. Como parte de la instalación de Lightbend FastData Platform/Cloudflow;
  2. Usando Helm:
    helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator
    helm install incubator/sparkoperator --namespace spark-operator
    	

  3. La implementación de manifiestos desde el repositorio oficial (https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/tree/master/manifest). Cabe señalar lo siguiente: Cloudflow incluye un operador con la versión de API v1beta1. Si se utiliza este tipo de instalación, las descripciones de los manifiestos de aplicaciones Spark deben basarse en ejemplos de etiquetas en Git con la versión de API correspondiente, por ejemplo, "v1beta1-0.9.0-2.4.0". La versión del operador se puede consultar en la descripción del CRD incluido en el operador en el diccionario "versions":
    oc get crd sparkapplications.sparkoperator.k8s.io -o yaml
    	

Si el operador está instalado correctamente, aparecerá un pod activo con el operador Spark en el proyecto correspondiente (por ejemplo, cloudflow-fdp-sparkoperator en el espacio Cloudflow para la instalación de Cloudflow) y aparecerá el tipo correspondiente de recursos de Kubernetes con el nombre "sparkapplications". Se puede estudiar las aplicaciones Spark existentes con el siguiente comando:

oc get sparkapplications -n {project}

Para ejecutar tareas utilizando el Spark Operator, se deben realizar 3 cosas:

  • crear una imagen Docker que incluya todas las bibliotecas necesarias, así como archivos de configuración y ejecutables. En el contexto objetivo, esta imagen se crea durante la etapa de CI/CD y se prueba en un clúster de prueba;
  • publicar la imagen Docker en un registro accesible desde el clúster de Kubernetes;
  • formular un manifiesto con el tipo "SparkApplication" y la descripción de la tarea que se va a ejecutar. Se pueden encontrar ejemplos de manifiestos en el repositorio oficial (por ejemplo, github.com/GoogleCloudPlatform/spark-on-k8s-operator/blob/v1beta1-0.9.0-2.4.0/examples/spark-pi.yaml). Es importante destacar los aspectos clave relacionados con el manifiesto:
    1. en el diccionario "apiVersion" debe indicarse la versión de la API correspondiente a la versión del operador;
    2. en el diccionario "metadata.namespace" debe indicarse el espacio de nombres en el que se ejecutará la aplicación;
    3. en el diccionario "spec.image" debe indicarse la dirección de la imagen Docker creada en el registro accesible;
    4. en el diccionario "spec.mainClass" debe indicarse la clase de tarea Spark que se requiere ejecutar al iniciar el proceso;
    5. en el diccionario "spec.mainApplicationFile" debe indicarse la ruta al archivo jar ejecutable;
    6. en el diccionario "spec.sparkVersion" debe indicarse la versión de Spark utilizada;
    7. en el diccionario "spec.driver.serviceAccount" debe indicarse la cuenta de servicio dentro del espacio de nombres Kubernetes correspondiente que se utilizará para ejecutar la aplicación;
    8. en el diccionario "spec.executor" debe indicarse la cantidad de recursos asignados a la aplicación;
    9. En el diccionario «spec.volumeMounts» debe especificarse el directorio local donde se crearán los archivos locales de la tarea Spark.

Ejemplo de creación de un manifiesto (aquí {spark-service-account} es la cuenta de servicio dentro del clúster de Kubernetes para ejecutar tareas 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"

En este manifiesto se especifica la cuenta de servicio para la cual es necesario crear los enlaces de roles requeridos antes de publicar el manifiesto, otorgando los permisos necesarios para que la aplicación Spark interactúe con la API de Kubernetes (si es necesario). En nuestro caso, la aplicación necesita permisos para crear Pods. Crearemos el enlace de rol necesario:

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

También cabe señalar que en la especificación de este manifiesto puede incluirse el parámetro «hadoopConfigMap», que permite especificar un ConfigMap con la configuración de Hadoop sin necesidad de colocar previamente el archivo correspondiente en la imagen de Docker. También es adecuado para la ejecución programada de tareas: mediante el parámetro «schedule» se puede establecer un horario para la ejecución de esta tarea.

Después de eso, guardamos nuestro manifiesto en el archivo spark-pi.yaml y lo aplicamos a nuestro clúster de Kubernetes:

oc apply -f spark-pi.yaml

Esto creará un objeto del tipo «sparkapplications»:

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

Se creará un pod con la aplicación, cuyo estado se mostrará en el «sparkapplications» creado. Se puede consultar con el siguiente comando:

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

Al finalizar la tarea, el POD cambiará a estado «Completed», que también se actualizará en el «sparkapplications». Los registros de la aplicación se pueden consultar en el navegador o usando el siguiente comando (aquí {sparkapplications-pod-name} es el nombre del pod de la tarea en ejecución):

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

La gestión de tareas de Spark también se puede realizar mediante la utilidad especializada sparkctl. Para su instalación, clonamos el repositorio con su código fuente, instalamos Go y compilamos esta utilidad:

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

Examinemos la lista de tareas de Spark en ejecución:

sparkctl list -n {project}

Creamos una descripción para la tarea de 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"

Ejecutemos la tarea descrita utilizando sparkctl:

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

Examinemos la lista de tareas de Spark en ejecución:

sparkctl list -n {project}

Revisemos la lista de eventos de la tarea de Spark en ejecución:

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

Verifiquemos el estado de la tarea de Spark en ejecución:

sparkctl status spark-pi -n {project}

En conclusión, nos gustaría considerar los inconvenientes descubiertos en la explotación de la versión estable actual de Spark (2.4.5) en Kubernetes:

  1. La primera y, posiblemente, la mayor desventaja es la falta de Localidad de Datos. A pesar de todas las desventajas, YARN tenía algunos beneficios en su uso, como el principio de entregar el código a los datos (y no los datos al código). Gracias a ello, las tareas de Spark se ejecutaban en los nodos donde se encontraban los datos implicados en los cálculos, lo que reducía significativamente el tiempo de entrega de datos a través de la red. Al utilizar Kubernetes, nos encontramos con la necesidad de mover a través de la red los datos involucrados en el trabajo de la tarea. Si estos son lo suficientemente grandes, el tiempo de ejecución de la tarea puede aumentar considerablemente, además de requerir un volumen considerable de espacio en disco asignado a las instancias de la tarea Spark para su almacenamiento temporal. Esta desventaja se puede reducir mediante el uso de herramientas de software especializadas que garantizan la localidad de datos en Kubernetes (por ejemplo, Alluxio), pero esto significa efectivamente que es necesario almacenar una copia completa de los datos en los nodos del clúster de Kubernetes.
  2. La segunda desventaja importante es la seguridad. Por defecto, las funciones relacionadas con la seguridad en la ejecución de tareas Spark están desactivadas, la opción de usar Kerberos en la documentación oficial no está cubierta (aunque los parámetros correspondientes aparecieron en la versión 3.0.0, lo que requerirá un trabajo adicional), y en la documentación sobre seguridad al utilizar Spark (https://spark.apache.org/docs/2.4.5/security.html) solo se mencionan YARN, Mesos y Standalone Cluster como almacenes de claves. Además, el usuario bajo el cual se ejecutan las tareas de Spark no puede ser especificado directamente; solo podemos establecer una cuenta de servicio bajo la cual funcionará el pod, y el usuario se elige según las políticas de seguridad configuradas. Debido a esto, o se utiliza el usuario root, lo cual no es seguro en un entorno de producción, o un usuario con un UID aleatorio, lo que resulta incómodo para la distribución de permisos de acceso a los datos (esto se puede resolver creando PodSecurityPolicies y vinculándolas a las cuentas de servicio correspondientes). Actualmente, se aborda ya sea incluyendo todos los archivos necesarios directamente en la imagen de Docker, o modificando el script de inicio de Spark para utilizar el mecanismo de almacenamiento y recuperación de secretos adoptado en su organización.
  3. El lanzamiento de tareas de Spark utilizando Kubernetes sigue estando oficialmente en modo experimental, y en el futuro podrían haber cambios significativos en los artefactos utilizados (archivos de configuración, imágenes base de Docker y scripts de inicio). De hecho, al preparar el material se probaron las versiones 2.3.0 y 2.4.5, y el comportamiento fue bastante diferente.

Estaremos atentos a las actualizaciones: recientemente se lanzó una nueva versión de Spark (3.0.0), que trajo cambios significativos en el funcionamiento de Spark en Kubernetes, aunque manteniendo el estatus experimental de soporte de este gestor de recursos. Es posible que las próximas actualizaciones realmente permitan recomendar plenamente abandonar YARN y ejecutar tareas de Spark en Kubernetes, sin temer por la seguridad de su sistema y sin necesidad de realizar modificaciones manuales en los componentes funcionales.

Fin.

Fuente: habr.com

Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS 🔥 Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS | ProHoster