Chers lecteurs, bonjour. Aujourd'hui, nous allons parler un peu d'Apache Spark et de ses perspectives de développement.

Dans le monde moderne du Big Data, Apache Spark est de facto le standard pour le dĂ©veloppement de tĂąches de traitement par lot. De plus, il est Ă©galement utilisĂ© pour crĂ©er des applications de streaming fonctionnant selon le concept de micro-batch, traitant et livrant des donnĂ©es par petites portions (Spark Structured Streaming). Traditionnellement, il faisait partie de la pile Hadoop gĂ©nĂ©rale, utilisant comme gestionnaire de ressources YARN (ou, dans certains cas, Apache Mesos). D'ici 2020, son utilisation en tant que telle pour la plupart des entreprises est en grand question en raison du manque de distributions Hadoop viables â le dĂ©veloppement de HDP et CDH est arrĂȘtĂ©, CDH est insuffisamment dĂ©veloppĂ© et coĂ»teux, et d'autres fournisseurs Hadoop ont soit cessĂ© d'exister, soit ont un avenir incertain. Par consĂ©quent, un intĂ©rĂȘt croissant de la part de la communautĂ© et des grandes entreprises se concentre sur le lancement d'Apache Spark via Kubernetes â devenu le standard pour l'orchestration de conteneurs et la gestion des ressources dans des environnements cloud privĂ©s et publics, il rĂ©sout le problĂšme de la planification difficile des ressources des tĂąches Spark sur YARN et offre une plateforme en dĂ©veloppement stable avec de nombreuses distributions commerciales et open source pour des entreprises de toutes tailles. De plus, sur la vague de sa popularitĂ©, la plupart des acteurs ont dĂ©jĂ acquis quelques installations et dĂ©veloppĂ© de l'expertise dans son utilisation, ce qui facilite la migration.
Depuis la version 2.3.0, Apache Spark dispose d'un support officiel pour l'exécution des tùches dans un cluster Kubernetes. Aujourd'hui, nous allons parler de la maturité actuelle de cette approche, des différentes options d'utilisation et des écueils à éviter lors de son implémentation.
Tout d'abord, examinons le processus de dĂ©veloppement des tĂąches et des applications basĂ©es sur Apache Spark et identifions les cas types dans lesquels il est nĂ©cessaire de lancer une tĂąche sur un cluster Kubernetes. Pour la prĂ©paration de ce post, la distribution utilisĂ©e est OpenShift et les commandes fournies seront pertinentes pour son outil en ligne de commande (oc). Pour d'autres distributions Kubernetes, les commandes correspondantes de l'outil en ligne de commande standard Kubernetes (kubectl) ou leurs Ă©quivalents (par exemple, pour oc adm policy) peuvent ĂȘtre utilisĂ©es.
La premiĂšre option d'utilisation est spark-submit
Dans le processus de dĂ©veloppement de tĂąches et d'applications, le dĂ©veloppeur doit exĂ©cuter des tĂąches pour dĂ©boguer la transformation des donnĂ©es. ThĂ©oriquement, des simulations peuvent ĂȘtre utilisĂ©es Ă cette fin, mais le dĂ©veloppement impliquant des instances rĂ©elles (mĂȘme si elles sont de test) des systĂšmes finaux s'est rĂ©vĂ©lĂ© plus rapide et de meilleure qualitĂ© dans ce type de tĂąches. Dans le cas oĂč nous procĂ©dons au dĂ©bogage sur des instances rĂ©elles des systĂšmes finaux, deux scĂ©narios de fonctionnement sont possibles :
- le développeur exécute la tùche Spark localement en mode standalone ;

- le développeur exécute la tùche Spark sur un cluster Kubernetes dans un environnement de test.

La premiÚre option a le droit d'exister, mais entraßne une série d'inconvénients :
- chaque développeur doit garantir l'accÚs depuis son poste de travail à toutes les instances nécessaires des systÚmes finaux ;
- la machine de travail doit disposer d'un nombre suffisant de ressources pour exécuter la tùche en développement.
La deuxiÚme option est dépourvue de ces inconvénients, car l'utilisation du cluster Kubernetes permet de réserver le pool de ressources nécessaire pour exécuter les tùches et d'assurer les accÚs nécessaires aux instances des systÚmes finaux, en fournissant un accÚs flexible via le modÚle de rÎle Kubernetes pour tous les membres de l'équipe de développement. Nous la mettons en avant en tant que premiÚre option d'utilisation : exécution des tùches Spark depuis la machine locale du développeur sur un cluster Kubernetes dans un environnement de test.
Nous parlerons en dĂ©tail du processus de configuration de Spark pour une exĂ©cution locale. Pour commencer Ă utiliser Spark, il doit ĂȘtre installĂ© :
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
Nous rassemblons les packages nécessaires pour travailler avec Kubernetes :
cd spark-2.4.5/
./build/mvn -Pkubernetes -DskipTests clean package
La construction complÚte prend beaucoup de temps, et pour créer des images Docker et les exécuter sur un cluster Kubernetes, seuls les fichiers jar du répertoire « assembly/ » sont réellement nécessaires, nous pouvons donc construire uniquement ce sous-projet :
./build/mvn -f ./assembly/pom.xml -Pkubernetes -DskipTests clean package
Pour exécuter des tùches Spark dans Kubernetes, il est nécessaire de créer une image Docker qui sera utilisée comme base. Deux approches sont possibles ici :
- L'image Docker créée contient le code exécutable de la tùche Spark ;
- L'image créée ne contient que Spark et ses dépendances nécessaires, le code exécutable étant stocké à distance (par exemple, dans HDFS).
Pour commencer, construisons une image Docker contenant un exemple de tùche Spark. Pour créer des images Docker, Spark dispose d'un outil correspondant appelé « docker-image-tool ». Examinons son aide :
./bin/docker-image-tool.sh --help
Avec cet outil, il est possible de créer des images Docker et de les charger dans des registres distants, mais par défaut, il présente certains inconvénients :
- il crĂ©e systĂ©matiquement 3 images Docker â pour Spark, PySpark et R ;
- il ne permet pas de spécifier le nom de l'image.
Nous allons donc utiliser une version modifiée de cet outil, présentée ci-dessous :
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
Avec cela, nous allons assembler l'image de base de Spark, contenant une tĂąche de test pour calculer le nombre Pi Ă l'aide de Spark (ici {docker-registry-url} â l'URL de votre registre d'images Docker, {repo} â le nom du dĂ©pĂŽt au sein du registre, correspondant au projet dans OpenShift, {image-name} â le nom de l'image (si une sĂ©paration en trois niveaux est utilisĂ©e, par exemple, comme dans le registre d'images intĂ©grĂ© de Red Hat OpenShift), {tag} â le tag de cette version de l'image) :
./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
Nous nous authentifions dans le cluster OKD Ă l'aide de l'outil en ligne de commande (ici {OKD-API-URL} â l'URL de l'API du cluster OKD) :
oc login {OKD-API-URL}
Obtenons le jeton de l'utilisateur courant pour l'authentification dans Docker Registry :
oc whoami -t
Nous nous authentifions dans le Docker Registry interne du cluster OKD (en utilisant le mot de passe généré par la commande précédente) :
docker login {docker-registry-url}
Téléchargons l'image Docker assemblée dans le Docker Registry OKD :
./bin/docker-image-tool-upd.sh -r {docker-registry-url}/{repo} -i {image-name} -t {tag} push
VĂ©rifions que l'image assemblĂ©e est disponible dans OKD. Pour cela, ouvrons dans le navigateur l'URL contenant la liste des images du projet correspondant (ici {project} â le nom du projet dans le cluster OpenShift, {OKD-WEBUI-URL} â l'URL de la console Web OpenShift) â https://{OKD-WEBUI-URL}/console/project/{project}/browse/images/{image-name}.
Pour exĂ©cuter des tĂąches, un compte de service avec des privilĂšges de lancement de pods sous root doit ĂȘtre créé (nous discuterons de ce point ultĂ©rieurement) :
oc create sa spark -n {project}
oc adm policy add-scc-to-user anyuid -z spark -n {project}
Exécutons la commande spark-submit pour soumettre la tùche Spark dans le cluster OKD, en spécifiant le compte de service créé et l'image 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
Ici :
--name â le nom de la tĂąche, qui participera Ă la formation des noms des pods Kubernetes ;
âclass â classe du fichier exĂ©cutable appelĂ© lors du lancement de la tĂąche ;
âconf â paramĂštres de configuration de Spark ;
spark.executor.instances â nombre d'instances de Spark Ă exĂ©cuter ;
spark.kubernetes.authenticate.driver.serviceAccountName â nom du compte de service Kubernetes utilisĂ© lors du lancement des pods (pour dĂ©finir le contexte de sĂ©curitĂ© et les autorisations lors de l'interaction avec l'API Kubernetes) ;
spark.kubernetes.namespace â espace de noms Kubernetes dans lequel les pods du driver et des exĂ©cutants seront lancĂ©s ;
spark.submit.deployMode â mode de dĂ©ploiement de Spark (pour le spark-submit standard, 'cluster' est utilisĂ©, pour Spark Operator et les versions plus rĂ©centes de Spark, c'est 'client') ;
spark.kubernetes.container.image â image Docker utilisĂ©e pour le lancement des pods ;
spark.master â URL de l'API Kubernetes (indiquĂ©e pour que l'accĂšs se fasse depuis la machine locale) ;
local:// â chemin vers le fichier exĂ©cutable Spark Ă l'intĂ©rieur de l'image Docker.
AccĂ©dez au projet OKD correspondant et examinez les pods créés â https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods.
Pour simplifier le processus de dĂ©veloppement, une autre option peut ĂȘtre utilisĂ©e, oĂč une image de base commune de Spark est créée, utilisĂ©e par toutes les tĂąches pour l'exĂ©cution, et les instantanĂ©s des fichiers exĂ©cutables sont publiĂ©s dans un stockage externe (comme Hadoop) et spĂ©cifiĂ©s lors de l'appel de spark-submit sous forme de lien. Dans ce cas, diffĂ©rentes versions des tĂąches Spark peuvent ĂȘtre exĂ©cutĂ©es sans reconstruction d'images Docker, en utilisant par exemple WebHDFS pour la publication des images. Envoyez une demande de crĂ©ation de fichier (ici {host} â hĂŽte du service WebHDFS, {port} â port du service WebHDFS, {path-to-file-on-hdfs} â chemin souhaitĂ© vers le fichier sur HDFS) :
curl -i -X PUT "http://{host}:{port}/webhdfs/v1/{path-to-file-on-hdfs}?op=CREATE
Vous recevrez alors une rĂ©ponse de ce type (ici {location} â l'URL Ă utiliser pour tĂ©lĂ©charger le fichier) :
HTTP/1.1 307 TEMPORARY_REDIRECT
Location: {location}
Content-Length: 0
TĂ©lĂ©chargez le fichier exĂ©cutable Spark sur HDFS (ici {path-to-local-file} â chemin vers le fichier exĂ©cutable Spark sur l'hĂŽte actuel) :
curl -i -X PUT -T {path-to-local-file} "{location}"
AprĂšs cela, nous pouvons faire un spark-submit en utilisant le fichier Spark tĂ©lĂ©chargĂ© sur HDFS (ici {class-name} â nom de la classe Ă exĂ©cuter pour la tĂąche) :
/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}
Ă noter qu'un accĂšs Ă HDFS et un fonctionnement de la tĂąche peuvent nĂ©cessiter des modifications du Dockerfile et du script entrypoint.sh â ajout d'une directive dans le Dockerfile pour copier les bibliothĂšques dĂ©pendantes dans le rĂ©pertoire /opt/spark/jars et inclure le fichier de configuration HDFS dans SPARK_CLASSPATH dans entrypoint.sh.
DeuxiĂšme cas d'utilisation â Apache Livy
Ensuite, lorsque la tĂąche est dĂ©veloppĂ©e et qu'il est nĂ©cessaire de tester le rĂ©sultat, la question de son exĂ©cution dans le cadre du processus CI/CD et du suivi de son statut se pose. Bien sĂ»r, il est possible de lâexĂ©cuter en utilisant un appel local spark-submit, mais cela complique l'infrastructure CI/CD car cela nĂ©cessite l'installation et la configuration de Spark sur les agents/exĂ©cuteurs du serveur CI et la configuration d'accĂšs Ă l'API Kubernetes. Pour ce cas, la rĂ©alisation ciblĂ©e choisie est l'utilisation d'Apache Livy comme API REST pour lancer des tĂąches Spark, hĂ©bergĂ©es Ă l'intĂ©rieur du cluster Kubernetes. Avec lui, il est possible de lancer des tĂąches Spark sur le cluster Kubernetes en utilisant des requĂȘtes cURL classiques, ce qui est facilement rĂ©alisable avec n'importe quelle solution CI, et son hĂ©bergement Ă l'intĂ©rieur du cluster Kubernetes rĂ©sout le problĂšme d'authentification lors de l'interaction avec l'API Kubernetes.

Soulignons-le comme un deuxiĂšme cas d'utilisation â exĂ©cution de tĂąches Spark dans le cadre du processus CI/CD sur le cluster Kubernetes dans un environnement de test.
Un peu sur Apache Livy â il fonctionne comme un serveur HTTP, fournissant une interface Web et une API RESTful, permettant de lancer spark-submit Ă distance en passant les paramĂštres nĂ©cessaires. Traditionnellement, il Ă©tait inclus dans la distribution HDP, mais peut Ă©galement ĂȘtre dĂ©ployĂ© sur OKD ou toute autre installation Kubernetes Ă l'aide de son manifeste et d'un ensemble d'images Docker appropriĂ©es, par exemple â . Pour notre cas, une image Docker similaire a Ă©tĂ© construite, incluant Spark version 2.4.5 Ă partir du Dockerfile suivant :
FROM java:8-alpine
ENV SPARK_HOME=/opt/spark
ENV LIVY_HOME=/opt/livy
ENV HADOOP_CONF_DIR=/etc/hadoop/conf
ENV SPARK_USER=spark
WORKDIR /opt
RUN apk add --update openssl wget bash &&
wget -P /opt https://downloads.apache.org/spark/spark-2.4.5/spark-2.4.5-bin-hadoop2.7.tgz &&
tar xvzf spark-2.4.5-bin-hadoop2.7.tgz &&
rm spark-2.4.5-bin-hadoop2.7.tgz &&
ln -s /opt/spark-2.4.5-bin-hadoop2.7 /opt/spark
RUN wget http://mirror.its.dal.ca/apache/incubator/livy/0.7.0-incubating/apache-livy-0.7.0-incubating-bin.zip &&
unzip apache-livy-0.7.0-incubating-bin.zip &&
rm apache-livy-0.7.0-incubating-bin.zip &&
ln -s /opt/apache-livy-0.7.0-incubating-bin /opt/livy &&
mkdir /var/log/livy &&
ln -s /var/log/livy /opt/livy/logs &&
cp /opt/livy/conf/log4j.properties.template /opt/livy/conf/log4j.properties
ADD livy.conf /opt/livy/conf
ADD spark-defaults.conf /opt/spark/conf/spark-defaults.conf
ADD entrypoint.sh /entrypoint.sh
ENV PATH="/opt/livy/bin:${PATH}"
EXPOSE 8998
ENTRYPOINT ["/entrypoint.sh"]
CMD ["livy-server"]
L'image créée peut ĂȘtre construite et tĂ©lĂ©chargĂ©e dans votre dĂ©pĂŽt Docker existant, par exemple, le dĂ©pĂŽt interne OKD. Le dĂ©ploiement se fait Ă l'aide du manifeste suivant ({registry-url} â URL du registre d'images Docker, {image-name} â nom de l'image Docker, {tag} â tag de l'image Docker, {livy-url} â URL souhaitĂ©e Ă laquelle le serveur Livy sera accessible ; le manifeste "Route" est appliquĂ© si Red Hat OpenShift est utilisĂ© comme distribution Kubernetes, sinon un manifeste Ingress ou un Service de type NodePort correspondant est utilisĂ© :
---
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
AprĂšs son utilisation et le lancement rĂ©ussi du pod, l'interface graphique de Livy est accessible via le lien : http://{livy-url}/ui. Avec Livy, nous pouvons publier notre tĂąche Spark en utilisant une requĂȘte REST, par exemple depuis Postman. Un exemple de collection avec des requĂȘtes est prĂ©sentĂ© ci-dessous (dans le tableau « args », des arguments de configuration avec des variables nĂ©cessaires Ă l'exĂ©cution de la tĂąche peuvent ĂȘtre transmis) :
{
"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 Soumettre un travail avec 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 Soumettre un travail sans 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": {}
}
ExĂ©cutons la premiĂšre requĂȘte de la collection, allons dans l'interface OKD et vĂ©rifions que la tĂąche a Ă©tĂ© lancĂ©e avec succĂšs â https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods. Pendant ce temps, une session apparaĂźtra dans l'interface Livy (http://{livy-url}/ui), oĂč il est possible de suivre l'avancement de la tĂąche et d'examiner les journaux de session Ă l'aide de l'API Livy ou de l'interface graphique.
Nous allons maintenant examiner le fonctionnement de Livy. Pour cela, nous Ă©tudierons les journaux du conteneur Livy Ă l'intĂ©rieur du pod avec le serveur Livy â https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods/{livy-pod-name}?tab=logs. Il en ressort que lors de l'appel de l'API REST Livy, le conteneur nommĂ© « livy » exĂ©cute un spark-submit, similaire Ă celui que nous avons utilisĂ© prĂ©cĂ©demment (ici, {livy-pod-name} est le nom du pod créé avec le serveur Livy). La collection prĂ©sente Ă©galement une deuxiĂšme requĂȘte, permettant de lancer des tĂąches avec un fichier exĂ©cutable Spark hĂ©bergĂ© Ă distance Ă l'aide du serveur Livy.
La troisiĂšme option d'utilisation â Spark Operator
Maintenant que la tĂąche a Ă©tĂ© testĂ©e, la question se pose de son lancement rĂ©gulier. La mĂ©thode native pour le lancement rĂ©gulier de tĂąches dans un cluster Kubernetes est l'entitĂ© CronJob, et il est possible de l'utiliser, mais actuellement, l'utilisation d'opĂ©rateurs pour gĂ©rer les applications dans Kubernetes est beaucoup plus populaire, et pour Spark, il existe un opĂ©rateur suffisamment mature, qui, entre autres, est utilisĂ© dans des solutions de niveau Entreprise (par exemple, Lightbend FastData Platform). Nous vous recommandons de l'utiliser â la version stable actuelle de Spark (2.4.5) a des capacitĂ©s de configuration assez limitĂ©es pour le lancement de tĂąches Spark dans Kubernetes, alors que dans la prochaine version majeure (3.0.0), un support complet de Kubernetes est annoncĂ©, mais la date de sortie reste inconnue. Spark Operator compense ce manque en ajoutant d'importants paramĂštres de configuration (par exemple, le montage de ConfigMap avec la configuration d'accĂšs Ă Hadoop dans les pods Spark) et la possibilitĂ© de lancer des tĂąches de façon rĂ©guliĂšre selon un emploi du temps.

Nous le mettons en avant en tant que troisiĂšme option d'utilisation â le lancement rĂ©gulier de tĂąches Spark sur un cluster Kubernetes dans un environnement de production.
Spark Operator est open-source et dĂ©veloppĂ© dans le cadre de Google Cloud Platform â . Son installation peut ĂȘtre rĂ©alisĂ©e de 3 maniĂšres :
- Dans le cadre de l'installation de Lightbend FastData Platform/Cloudflow ;
- Ă l'aide de Helm :
helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator helm install incubator/sparkoperator --namespace spark-operator - L'utilisation des manifestes depuis le dĂ©pĂŽt officiel (https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/tree/master/manifest). Ă noter que Cloudflow inclut un opĂ©rateur avec la version API v1beta1. Si ce type d'installation est utilisĂ©, la description des manifestes des applications Spark doit ĂȘtre basĂ©e sur des exemples provenant des balises dans Git avec la version API correspondante, par exemple, « v1beta1-0.9.0-2.4.0 ». La version de l'opĂ©rateur peut ĂȘtre consultĂ©e dans la description CRD qui fait partie de l'opĂ©rateur dans le dictionnaire « versions » :
oc get crd sparkapplications.sparkoperator.k8s.io -o yaml
Si l'opérateur est installé correctement, un pod actif avec l'opérateur Spark apparaßtra dans le projet correspondant (par exemple, cloudflow-fdp-sparkoperator dans l'espace Cloudflow pour l'installation de Cloudflow), ainsi qu'un type de ressource Kubernetes nommé « sparkapplications ». Pour examiner les applications Spark existantes, vous pouvez utiliser la commande suivante :
oc get sparkapplications -n {project}
Pour exécuter des tùches à l'aide de Spark Operator, trois choses sont nécessaires :
- créer une image Docker incluant toutes les bibliothÚques nécessaires ainsi que les fichiers de configuration et exécutables. Dans l'image cible, il s'agit d'une image créée lors de l'étape CI/CD et testée sur un cluster de test ;
- publier l'image Docker dans un registre accessible depuis le cluster Kubernetes ;
- émettre un manifeste de type « SparkApplication » décrivant la tùche à exécuter. Des exemples de manifestes sont disponibles dans le dépÎt officiel (par exemple, ). Il est important de noter quelques points concernant le manifeste :
- dans le dictionnaire « apiVersion », la version API correspondant Ă la version de l'opĂ©rateur doit ĂȘtre spĂ©cifiĂ©e ;
- dans le dictionnaire « metadata.namespace », l'espace de noms dans lequel l'application sera exĂ©cutĂ©e doit ĂȘtre spĂ©cifiĂ© ;
- dans le dictionnaire « spec.image », l'adresse de l'image Docker créée dans le registre accessible doit ĂȘtre spĂ©cifiĂ©e ;
- dans le dictionnaire « spec.mainClass », la classe de la tĂąche Spark qui doit ĂȘtre exĂ©cutĂ©e lors du lancement du processus doit ĂȘtre indiquĂ©e ;
- dans le dictionnaire « spec.mainApplicationFile », le chemin vers le fichier jar exĂ©cutable doit ĂȘtre indiquĂ© ;
- dans le dictionnaire « spec.sparkVersion », la version de Spark utilisĂ©e doit ĂȘtre spĂ©cifiĂ©e ;
- dans le dictionnaire « spec.driver.serviceAccount », le compte de service dans l'espace de noms Kubernetes correspondant Ă utiliser pour exĂ©cuter l'application doit ĂȘtre spĂ©cifiĂ© ;
- dans le dictionnaire « spec.executor », le nombre de ressources allouĂ©es Ă l'application doit ĂȘtre prĂ©cisĂ©.
- Dans le dictionnaire «spec.volumeMounts», un rĂ©pertoire local doit ĂȘtre spĂ©cifiĂ©, dans lequel les fichiers locaux de la tĂąche Spark seront créés.
Exemple de création d'un manifeste (ici {spark-service-account} est le compte de service au sein du cluster Kubernetes pour l'exécution des tùches 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"
Ce manifeste indique un compte de service pour lequel il est nécessaire de créer les liaisons de rÎles requises avant de publier le manifeste, afin de fournir les droits d'accÚs nécessaires pour permettre à l'application Spark d'interagir avec l'API Kubernetes (si nécessaire). Dans notre cas, l'application a besoin des droits pour créer des pods. Créons la liaison de rÎle requise :
oc adm policy add-role-to-user edit system:serviceaccount:{project}:{spark-service-account} -n {project}
Il convient Ă©galement de noter que dans la spĂ©cification de ce manifeste, le paramĂštre «hadoopConfigMap» peut ĂȘtre spĂ©cifiĂ©, permettant de dĂ©signer un ConfigMap avec la configuration Hadoop sans avoir Ă prĂ©alablement insĂ©rer le fichier correspondant dans l'image Docker. Cela convient Ă©galement pour le lancement rĂ©gulier des tĂąches â le paramĂštre «schedule» peut ĂȘtre utilisĂ© pour dĂ©finir un calendrier de lancement de cette tĂąche.
AprĂšs cela, nous sauvegardons notre manifeste dans le fichier spark-pi.yaml et l'appliquons Ă notre cluster Kubernetes :
oc apply -f spark-pi.yaml
Un objet de type «sparkapplications» sera créé :
oc get sparkapplications -n {project}
> NAME AGE
> spark-pi 22h
Un pod avec l'application sera créé, dont le statut sera affiché dans le «sparkapplications» créé. On peut le consulter avec la commande suivante :
oc get sparkapplications spark-pi -o yaml -n {project}
Ă la fin de la tĂąche, le POD passera au statut «Completed», qui sera Ă©galement mis Ă jour dans le «sparkapplications». Les journaux de l'application peuvent ĂȘtre consultĂ©s dans le navigateur ou avec la commande suivante (ici {sparkapplications-pod-name} est le nom du pod de la tĂąche exĂ©cutĂ©e) :
oc logs {sparkapplications-pod-name} -n {project}
La gestion des tĂąches Spark peut Ă©galement ĂȘtre rĂ©alisĂ©e Ă l'aide de l'outil spĂ©cialisĂ© sparkctl. Pour l'installer, nous clonons le dĂ©pĂŽt contenant son code source, installons Go et compilons cet utilitaire :
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
Examinons la liste des tĂąches Spark en cours :
sparkctl list -n {project}
Créons une description pour une tùche 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"
Lançons la tùche décrite à l'aide de sparkctl :
sparkctl create spark-app.yaml -n {project}
Examinons la liste des tĂąches Spark en cours :
sparkctl list -n {project}
Examinons la liste des événements de la tùche Spark en cours :
sparkctl event spark-pi -n {project} -f
Examinons le statut de la tĂąche Spark en cours :
sparkctl status spark-pi -n {project}
En conclusion, examinons les inconvénients découverts lors de l'exploitation de la version stable actuelle de Spark (2.4.5) dans Kubernetes :
- Le premier et probablement le principal inconvĂ©nient est l'absence de localitĂ© des donnĂ©es. MalgrĂ© toutes les faiblesses de YARN, il y avait aussi des avantages Ă son utilisation, comme le principe de livraison de code aux donnĂ©es (et non l'inverse). GrĂące Ă cela, les tĂąches Spark Ă©taient exĂ©cutĂ©es sur les nĆuds oĂč se trouvaient les donnĂ©es impliquĂ©es dans les calculs, ce qui rĂ©duisait considĂ©rablement le temps de transfert des donnĂ©es sur le rĂ©seau. Avec Kubernetes, nous sommes confrontĂ©s Ă la nĂ©cessitĂ© de dĂ©placer des donnĂ©es requises pour le traitement des tĂąches Ă travers le rĂ©seau. Si elles sont suffisamment volumineuses, le temps d'exĂ©cution peut augmenter considĂ©rablement et il peut Ă©galement ĂȘtre nĂ©cessaire de disposer d'une quantitĂ© importante d'espace disque allouĂ© aux instances de tĂąches Spark pour un stockage temporaire. Cet inconvĂ©nient peut ĂȘtre attĂ©nuĂ© grĂące Ă l'utilisation d'outils logiciels spĂ©cialisĂ©s qui assurent la localitĂ© des donnĂ©es dans Kubernetes (comme Alluxio), mais cela implique en fait la nĂ©cessitĂ© de conserver une copie complĂšte des donnĂ©es sur les nĆuds du cluster Kubernetes.
- Le deuxiĂšme inconvĂ©nient important est la sĂ©curitĂ©. Par dĂ©faut, les fonctionnalitĂ©s liĂ©es Ă la sĂ©curitĂ© lors de l'exĂ©cution des tĂąches Spark sont dĂ©sactivĂ©es, l'utilisation de Kerberos dans la documentation officielle n'est pas abordĂ©e (bien que des paramĂštres pertinents soient apparus dans la version 3.0.0, ce qui nĂ©cessitera un travail supplĂ©mentaire), et dans la documentation sur la sĂ©curitĂ© lors de l'utilisation de Spark (https://spark.apache.org/docs/2.4.5/security.html), seuls YARN, Mesos et le cluster autonome figurent parmi les magasins de clĂ©s. De plus, l'utilisateur sous lequel les tĂąches Spark sont exĂ©cutĂ©es ne peut pas ĂȘtre spĂ©cifiĂ© directement â nous ne dĂ©finissons qu'un compte de service sous lequel le pod fonctionnera, et l'utilisateur est sĂ©lectionnĂ© en fonction des politiques de sĂ©curitĂ© configurĂ©es. Par consĂ©quent, soit nous utilisons l'utilisateur root, ce qui n'est pas sĂ»r dans un environnement de production, soit un utilisateur avec un UID alĂ©atoire, ce qui complique la rĂ©partition des droits d'accĂšs aux donnĂ©es (rĂ©soluble en crĂ©ant des PodSecurityPolicies et en les liant aux comptes de service appropriĂ©s). Actuellement, la solution passe soit par la mise en place de tous les fichiers nĂ©cessaires directement dans l'image Docker, soit par la modification du script de lancement de Spark pour utiliser le mĂ©canisme de stockage et d'accĂšs aux secrets adoptĂ© dans votre organisation.
- Le lancement des tùches Spark via Kubernetes est toujours en mode expérimental et des changements significatifs concernant les artefacts utilisés (fichiers de configuration, images de base Docker et scripts de lancement) sont possibles à l'avenir. En effet, lors de la préparation du matériel, les versions 2.3.0 et 2.4.5 ont été testées et leur comportement était sensiblement différent.
Nous attendons des mises Ă jour â une nouvelle version de Spark (3.0.0) a rĂ©cemment Ă©tĂ© publiĂ©e, apportant des changements notables au fonctionnement de Spark sur Kubernetes, tout en conservant le statut expĂ©rimental du support de ce gestionnaire de ressources. Il est possible que les prochaines mises Ă jour permettent effectivement de recommander pleinement d'abandonner YARN et de lancer des tĂąches Spark sur Kubernetes, sans craindre pour la sĂ©curitĂ© de votre systĂšme et sans nĂ©cessiter d'ajouts fonctionnels manuels.
Fin.
Source : habr.com


