Të dashur lexues, përshëndetje. Sot do të flasim pak për Apache Spark dhe perspektivat e zhvillimit të tij.

Në botën moderne të Big Data, Apache Spark është de facto standard për zhvillimin e detyrave të përpunimit të paketave të dhënash. Përveç kësaj, ai përdoret gjithashtu për të krijuar aplikacione streaming që funksionojnë në konceptin micro batch, duke përpunuar dhe dërguar të dhënat në sasi të vogla (Spark Structured Streaming). Tradicionalisht, ai ka qenë një pjesë e stack-ut të përgjithshëm Hadoop, duke përdorur YARN si menaxher resurshesh (ose, në disa raste, Apache Mesos). Deri në vitin 2020, përdorimi i tij në formën tradicionale për shumicën e kompanive është me shumë dyshime për shkak të mungesës së distribucioneve të mira Hadoop - zhvillimi i HDP dhe CDH është ndalur, CDH është i papërpunuar dhe ka një kosto të lartë, ndërsa ofruesit e tjerë të Hadoop ose kanë ndaluar së ekzistuari, ose kanë një të ardhme të paqartë. Prandaj, një interes në rritje për komunitetin dhe kompanitë e mëdha është lëshimi i Apache Spark me ndihmën e Kubernetes - duke u bërë standard në orkestrimin e kontejnerëve dhe menaxhimin e burimeve në cloud privat dhe publik, ai zgjidh problemin e planifikimit të vështirë të resurseve të detyrave Spark në YARN dhe ofron një platformë që evoluon stabilisht me shumë distribucione komerciale dhe të hapura për kompani të të gjitha madhësive dhe llojeve. Për më tepër, me rritjen e popullaritetit, shumica e tyre tashmë ka arritur të krijojë disa instalime dhe të rrisin ekspertizën e tyre në përdorimin e tij, që e bën migrimin më të lehtë.
Deri në versionin 2.3.0, Apache Spark mori mbështetje zyrtare për ekzekutimin e detyrave në klasterin Kubernetes dhe sot, do të flasim për maturinë aktuale të këtij qasjes, variacione të ndryshme të përdorimit të tij dhe pengesat me të cilat do të përballeni gjatë implementimit.
Së pari, le të shqyrtojmë procesin e zhvillimit të detyrave dhe aplikacioneve mbi bazën e Apache Spark dhe të theksojmë rastet tipike kur është e nevojshme të ekzekutohet një detyrë në klasterin Kubernetes. Gjatë përgatitjes së këtij postimi, distribucioni që përdoret është OpenShift dhe do të jepen komandat që janë relevante për utilitarin e tij të komandave (oc). Për distribucione të tjera të Kubernetes mund të përdoren komandat përkatëse të utilitarit standard të komandave Kubernetes (kubectl) ose analogët e tyre (p.sh. për oc adm policy).
Varianti i parë i përdorimit - spark-submit
Gjatë zhvillimit të detyrave dhe aplikacioneve, zhvilluesi ka nevojë të ekzekutojë detyrat për debugging-un e transformimeve të dhënash. Teorikisht, për këto qëllime mund të përdoren stub-e, por zhvillimi me implikimin e mjeteve reale (edhe nëse janë prova) ka treguar se është më i shpejtë dhe më cilësor në këtë klasë detyrash. Në rastin kur po bëjmë debugging mbi mjete reale, janë të mundshme dy skenarë veprimi:
- zhvilluesi ekzekuton detyrën Spark lokalisht në modalitetin standalone;

- zhvilluesi ekzekuton detyrën Spark në klasterin Kubernetes në konturin e provës.

Varianti i parë ka të drejta të ekzistojë, por sjell një sërë disavantazhesh:
- çdo zhvillues duhet të sigurojë qasje nga vendi i punës deri te të gjitha mjetet e nevojshme për të;
- në makinerinë e punës nevojiten burime të mjaftueshme për të ekzekutuar detyrën që po zhvillohet.
Varianti i dytë është i liruar nga këto mungesa, pasi përdorimi i klasterit Kubernetes lejon ndarjen e grumbullit të nevojshëm të burimeve për ekzekutimin e detyrave dhe siguron për to qasje të nevojshme ndaj mjetëve përfundimtare, duke ofruar fleksibilitet në qasjen përmes modelit të rolit të Kubernetes për të gjithë anëtarët e ekipit të zhvillimit. Le ta ndërlidhim këtë si varianti i parë i përdorimit - ekzekutimi i detyrave Spark nga makineria lokale e zhvilluesit në klasterin Kubernetes në konturin e provës.
Le të flasim më në detaje për procesin e konfigurimit të Spark për ekzekutimin lokal. Për të filluar përdorimin e Spark, është e nevojshme ta instaloni:
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
Po përgatisim paketat e nevojshme për punën me Kubernetes:
cd spark-2.4.5/
./build/mvn -Pkubernetes -DskipTests clean package
Ndërsa ndërtimi i plotë merr shumë kohë, për të krijuar imazhe Docker dhe për t'i ekzekutuar në klasterin Kubernetes në të vërtetë ne kemi nevojë vetëm për skedarët jar nga katalogu 'assembly/', prandaj mund të ndihmojmë vetëm këtë nënprojekt:
./build/mvn -f ./assembly/pom.xml -Pkubernetes -DskipTests clean package
Për të ekzekutuar detyrat Spark në Kubernetes, nevojitet të krijoni një imazh Docker, i cili do të përdoret si bazë. Këtu janë dy qasje të mundshme:
- Imazhi Docker i krijuar përmban kodin ekzekutiv të detyrës Spark;
- Imazhi i krijuar përmban vetëm Spark dhe varësitë e nevojshme, ndërsa kodi ekzekutiv është i vendosur në distancë (p.sh., në HDFS).
Për të filluar, le të krijojmë një imazh Docker që përmban një shembull testues të detyrës Spark. Për krijimin e imazheve Docker, Spark ka një utilitar përkatës të quajtur «docker-image-tool». Të shqyrtojmë ndihmën e tij:
./bin/docker-image-tool.sh --help
Me të mund të krijoni imazhe Docker dhe t'i ngarkoni ato në regjistra të distancë, por në mënyrë standarde ka disa disavantazhe:
- duhet të krijojë domosdoshmërisht 3 imazhe Docker - për Spark, PySpark dhe R;
- nuk lejon caktimin e emrit të imazhit.
Prandaj, ne do të përdorim një version të modifikuar të këtij utilitari, të shpjeguar më poshtë:
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
Me të mbledhim imazhin bazë Spark, i cili përmban një detyrë testuese për llogaritjen e numrit Pi duke përdoruar Spark (këtu {docker-registry-url} - URL-ja e regjistrit tuaj të imazheve Docker, {repo} - emri i repository-t brenda regjistrit, që përputhet me projektin në OpenShift, {image-name} - emri i imazhit (nëse përdoret një ndarëse tre-niveli, për shembull, si në regjistrin e integruar të imazheve Red Hat OpenShift), {tag} - etiketa e kësaj versioni të imazhit):
./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
Autorizohemi në klasterin OKD përmes utilitarit të konsolës (këtu {OKD-API-URL} - URL API të klasterit OKD):
oc login {OKD-API-URL}
Merrni token-in e përdoruesit aktual për autorizim në Docker Registry:
oc whoami -t
Autorizohemi në regjistrin e brendshëm Docker të klasterit OKD (si fjalëkalim përdorim token-in e marrë nga komandë e mëparshme):
docker login {docker-registry-url}
Ngarkoni imazhin e mbledhur në Docker Registry OKD:
./bin/docker-image-tool-upd.sh -r {docker-registry-url}/{repo} -i {image-name} -t {tag} push
Kontrolloni që imazhi i mbledhur është i disponueshëm në OKD. Për këtë hapni në shfletues URL-në me listën e imazheve të projektit përkatës (këtu {project} - emri i projektit brenda klasterit OpenShift, {OKD-WEBUI-URL} - URL Web konsolës OpenShift) - https://{OKD-WEBUI-URL}/console/project/{project}/browse/images/{image-name}.
Për të ekzekutuar detyra duhet të krijohet një llogari shërbimi me privilegje për ekzekutimin e pods nën root (këtu do ta diskutojmë këtë moment):
oc create sa spark -n {project}
oc adm policy add-scc-to-user anyuid -z spark -n {project}
Ekzekutojmë komandën spark-submit për të publikuar detyrën Spark në klasterin OKD, duke caktuar llogarinë shërbimi të krijuar dhe imazhin 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
Këtu:
âname â emri i detyrĂ«s, i cili do tĂ« marrĂ« pjesĂ« nĂ« formimin e emrave tĂ« pods tĂ« Kubernetes;
âclass â klasa e skedarit ekzekutiv qĂ« thirret gjatĂ« ekzekutimit tĂ« detyrĂ«s;
âconf â parametrat konfigurues tĂ« Spark;
spark.executor.instances â numri i ekzekutorĂ«ve tĂ« Spark qĂ« do tĂ« ekzekutohen;
spark.kubernetes.authenticate.driver.serviceAccountName â emri i llogarisĂ« shĂ«rbimi tĂ« Kubernetes, e pĂ«rdorur gjatĂ« ekzekutimit tĂ« pods (pĂ«r tĂ« pĂ«rcaktuar kontekstin e sigurisĂ« dhe mundĂ«sitĂ« gjatĂ« ndĂ«rveprimit me API Kubernetes);
spark.kubernetes.namespace â hapĂ«sira e emrave tĂ« Kubernetes, ku do tĂ« ekzekutohen pods e drejtuesit dhe ekzekutorĂ«ve;
spark.submit.deployMode â mĂ«nyra e ekzekutimit tĂ« Spark (pĂ«r spark-submit standard pĂ«rdoret «cluster», pĂ«r Spark Operator dhe versione mĂ« tĂ« vonshme tĂ« Spark «client»);
spark.kubernetes.container.image â imazhi Docker, i pĂ«rdorur pĂ«r ekzekutimin e pods;
spark.master â URL API e Kubernetes (specifikohet e jashtme, kur i drejtohemi nga makina lokale);
local:// â rruga deri te skedari ekzekutiv tĂ« Spark brenda imazhit Docker.
Shkoni nĂ« projektin pĂ«rkatĂ«s OKD dhe shqyrtoni pods e krijuara â https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods.
Për të simplifikuar procesin e zhvillimit, mund të përdoret një variant tjetër, ku krijohet një imazh bazë i përbashkët Spark, që përdoret nga të gjitha detyrat për ekzekutim, dhe snapshots e skedareve ekzekutivë publikohen në një ruajtje të jashtme (p.sh., Hadoop) dhe caktohen gjatë thirrjes së spark-submit në formën e një lidhjeje. Në këtë rast, mund të ekzekutoni versione të ndryshme të detyrave Spark pa mbledhur përsëri imazhet Docker, duke përdorur si publikues për imazhet, për shembull, WebHDFS. Dërgo një kërkesë për krijimin e skedarit (këtu {host} - host-i i shërbimit WebHDFS, {port} - porta e shërbimit WebHDFS, {path-to-file-on-hdfs} - rruga e dëshiruar te skedari në HDFS):
curl -i -X PUT "http://{host}:{port}/webhdfs/v1/{path-to-file-on-hdfs}?op=CREATE
Në këtë rast do të merrni një përgjigje si kjo (këtu {location} - kjo është URL që duhet të përdorni për ngarkimin e skedarit):
HTTP/1.1 307 TEMPORARY_REDIRECT
Location: {location}
Content-Length: 0
Ngarko skedarin ekzekutiv të Spark në HDFS (këtu {path-to-local-file} - rruga deri te skedari ekzekutiv të Spark në hostin aktual):
curl -i -X PUT -T {path-to-local-file} "{location}"
Pas kësaj, mund të bëni spark-submit duke përdorur skedarin Spark, të ngarkuar në HDFS (këtu {class-name} - emri i klasës që duhet të ekzekutohet për realizimin e detyrës):
/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}
Duhet të theksohet se për qasje në HDFS dhe për të siguruar funksionimin e detyrës mund të jetë e nevojshme të modifikoni Dockerfile dhe skenarin entrypoint.sh - të shtoni në Dockerfile një direktivë për të kopjuar bibliotekat varësore në drejtori /opt/spark/jars dhe të përfshini skedarin e konfigurimit të HDFS në SPARK_CLASSPATH në entrypoint.sh.
Varianti i dytë i përdorimit - Apache Livy
Pastaj, kur projekti është dizajnuar dhe rezultati duhet të testohet, lind pyetja për ekzekutimin e tij në procesin CI/CD dhe për monitorimin e statusit të tij. Natyrisht, mund ta ekzekutoni atë edhe me thirrjen lokale spark-submit, por kjo e komplikon infrastrukturën CI/CD pasi kërkon instalimin dhe konfigurimin e Spark në agjentët/runners e serverit CI dhe konfigurimin e qasjes në API-në e Kubernetes. Për këtë rast, zgjidhja e synuar është përdorimi i Apache Livy si REST API për ekzekutimin e detyrave Spark, e vendosur brenda klasterit Kubernetes. Me të, mund të ekzekutoni detyra Spark në klasterin Kubernetes duke përdorur kërkesa të zakonshme cURL, që është lehtësisht e realizueshme mbi çdo zgjidhje CI, dhe vendosja e tij brenda klasterit Kubernetes zgjidh çështjen e autentifikimit gjatë ndërveprimit me API-në e Kubernetes.

Ta nxjerrim në pah si variantin e dytë të përdorimit - ekzekutimin e detyrave Spark në procesin CI/CD në klasterin Kubernetes në mjedisin e testimit.
Pak mbi Apache Livy - ai funksionon si një server HTTP, duke ofruar një ndërfaqe Web dhe një RESTful API, që lejon ekzekutimin në distancë të spark-submit, duke kaluar parametrat e nevojshëm. Tradicionalisht, ai ofrohej si pjesë e shpërndarjes HDP, por gjithashtu mund të vihet në funksion në OKD ose në çdo instalim tjetër Kubernetes me ndihmën e manifestit përkatës dhe një seti imazhe Docker, për shembull, ky - . Për rastin tonë, është ndërtuar një imazh Docker të ngjashëm, që përmban Spark versionin 2.4.5 nga Dockerfile i mëposhtëm:
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"]
Imazhi i krijuar mund të ndërtohet dhe ngarkohet në depozanin tuaj Docker ekzistues, për shembull, depozita e brendshme OKD. Për ta vendosur atë, përdoret manifesti i mëposhtëm ({registry-url} - URL e depozitës së imazheve Docker, {image-name} - emri i imazhit Docker, {tag} - etiketa e imazhit Docker, {livy-url} - URL e dëshiruar, me të cilin do të jetë në dispozicion serveri Livy; manifesti "Route" përdoret nëse shpërndarja e Kubernetes është Red Hat OpenShift, në përndryshe përdoret manifesti përkatës Ingress ose Lloji i Shërbimit 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
Pas aplikuar dhe pas nisjes me sukses të pod-it, ndërfaqja grafike e Livy është e aksesueshme në lidhjen: http://{livy-url}/ui. Me Livy, ne mund të publikojmë detyrat tona Spark duke përdorur një kërkesë REST, për shembull, nga Postman. Më poshtë është një shembull koleksioni me kërkesat (në vargun «args» mund të kalohen argumente konfigurimi me variablat e nevojshme për ekzekutimin e detyrës):
{
"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 Dërgo detyrën me 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 Dërgo detyrën pa 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": {}
}
Do tĂ« kryejmĂ« kĂ«rkesĂ«n e parĂ« nga koleksioni, do tĂ« kalojmĂ« nĂ« ndĂ«rfaqen OKD dhe do tĂ« verifikojmĂ« se detyra Ă«shtĂ« nisur me sukses â https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods. NĂ« ndĂ«rfaqen Livy (http://{livy-url}/ui) do tĂ« shfaqet njĂ« sesion, ku pĂ«rmes API-t tĂ« Livy ose ndĂ«rfaqes grafike mund tĂ« ndjekim ecjen e detyrĂ«s dhe tĂ« studiojmĂ« logjet e sesionit.
Tani do tĂ« tregojmĂ« mekanizmin e funksionimit tĂ« Livy. PĂ«r kĂ«tĂ«, do tĂ« studiojmĂ« regjistrat e kontejnerit Livy brenda pod-it me serverin Livy â https://{OKD-WEBUI-URL}/console/project/{project}/browse/pods/{livy-pod-name}?tab=logs. KĂ«to tregojnĂ« se, kur thirret API REST tĂ« Livy nĂ« kontejnerin me emrin "livy", ekzekutohet spark-submit, i ngjashĂ«m me atĂ« qĂ« pĂ«rdorĂ«m mĂ« sipĂ«r (kĂ«tu {livy-pod-name} Ă«shtĂ« emri i pod-it tĂ« krijuar me serverin Livy). NĂ« koleksion pĂ«rfshihet gjithashtu njĂ« kĂ«rkesĂ« e dytĂ«, e cila lejon nisjen e detyrave me vendosje tĂ« largĂ«t tĂ« skedarĂ«ve ekzekutivĂ« Spark pĂ«rmes serverit Livy.
MĂ«nyra e tretĂ« e pĂ«rdorimit â Spark Operator
Tani, pasi detyra Ă«shtĂ« testuar, ngrihet pyetja e nisjes sĂ« saj nĂ« mĂ«nyrĂ« tĂ« rregullt. MĂ«nyra natyrale pĂ«r tĂ« nisur detyra nĂ« klasterin Kubernetes Ă«shtĂ« entiteti CronJob dhe mund ta pĂ«rdorim, por aktualisht ka njĂ« popullaritet mĂ« tĂ« madh pĂ«rdorimi tĂ« operatorĂ«ve pĂ«r menaxhimin e aplikacioneve nĂ« Kubernetes dhe pĂ«r Spark ka njĂ« operator tĂ« mjaftueshĂ«m tĂ« pjekur, qĂ« po ashtu pĂ«rdoret nĂ« zgjidhjet e nivelit Enterprise (p.sh. Lightbend FastData Platform). Ne rekomandojmĂ« ta pĂ«rdorim atĂ« â versioni i tanishĂ«m stabil i Spark (2.4.5) ka mundĂ«si shumĂ« tĂ« kufizuara pĂ«r konfigurimin e nisjeve tĂ« detyrave Spark nĂ« Kubernetes, ndĂ«rsa nĂ« versionin e ardhshĂ«m kryesor (3.0.0) Ă«shtĂ« shpallur mbĂ«shtetje e plotĂ« pĂ«r Kubernetes, por data e daljes sĂ« tij mbetet e panjohur. Spark Operator e kompenson kĂ«tĂ« mungesĂ«, duke shtuar parametra tĂ« rĂ«ndĂ«sishĂ«m konfigurimi (p.sh. montimin e ConfigMap me konfigurimin e aksesit nĂ« Hadoop nĂ« pod-Ă«t Spark) dhe mundĂ«sinĂ« e nisjes sĂ« rregullt tĂ« detyrave sipas njĂ« programi.

Ta ndĂ«rlidhim atĂ« si opsionin e tretĂ« tĂ« pĂ«rdorimit â nisjen e rregullt tĂ« detyrave Spark nĂ« klasterin Kubernetes nĂ« ambientin produktiv.
Spark Operator ka kod burimor tĂ« hapur dhe zhvillohet nĂ« kuadĂ«r tĂ« Google Cloud Platform â . Instalimi i tij mund tĂ« kryhet nĂ« 3 mĂ«nyra:
- Në kuadër të instalimit të Lightbend FastData Platform/Cloudflow;
- Me anë të Helm:
helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator helm install incubator/sparkoperator --namespace spark-operator - Duke përdorur manifestet nga repositori zyrtar (https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/tree/master/manifest). Duhet theksuar se në përmbajtjen e Cloudflow përfshihet operatori me versionin API v1beta1. Nëse përdoret ky tip instalimi, përshkrimet e manifestimeve të aplikacioneve Spark duhet të ndjekin shembujt nga etiketat në Git me versionin përkatës të API, p.sh. "v1beta1-0.9.0-2.4.0". Versioni i operatorit mund të shihet në përshkrimin e CRD, që bën pjesë në operator në fjalorin "versions":
oc get crd sparkapplications.sparkoperator.k8s.io -o yaml
Nëse operatori është instaluar siç duhet, atëherë në projektin përkatës do të shfaqet një pod aktiv me operatorin Spark (p.sh. cloudflow-fdp-sparkoperator në hapësirën e Cloudflow për instalimin e Cloudflow) dhe do të shfaqet lloji përkatës i burimeve Kubernetes me emrin "sparkapplications". Për të studiuar aplikacionet ekzistuese të Spark, mund të përdoret komandën e mëposhtme:
oc get sparkapplications -n {project}
Për të nisur detyra me Spark Operator, nevojiten 3 gjëra:
- të krijohet një imazh Docker që përfshin të gjitha bibliotekat e nevojshme, si dhe skedarët e konfigurimit dhe ekzekutiv. Në skemën përfundimtare, ky është imazhi i krijuar në fazën CI/CD dhe e testuar në klasterin e testit;
- të publikohet imazhi Docker në regjistrin e disponueshëm nga klasteri Kubernetes;
- të krijohet një manifest me tipin "SparkApplication" dhe përshkrimin e detyrës që do të nisë. Shembujt e manifestimeve janë të disponueshme në repositorin zyrtar (p.sh., ). Duhet theksuar disa pika të rëndësishme në lidhje me manifestin:
- në fjalorin "apiVersion" duhet të jetë e caktuar versioni i API, përkatës me versionin e operatorit;
- në fjalorin "metadata.namespace" duhet të jetë e caktuar hapësira emërore, ku do të nisë aplikacioni;
- në fjalorin "spec.image" duhet të jetë e caktuar adresa e imazhit të krijuar në regjistrin e disponueshëm;
- në fjalorin "spec.mainClass" duhet të jetë e caktuar klasa e detyrës Spark, që duhet të niset në fillimin e procesit;
- në fjalorin "spec.mainApplicationFile" duhet të jetë e caktuar rruga për skedarin ekzekutiv jar;
- në fjalorin "spec.sparkVersion" duhet të jetë e caktuar versioni i përdorur i Spark;
- në fjalorin "spec.driver.serviceAccount" duhet të jetë e caktuar llogaria shërbyese brenda përkatës hapësire emërore Kubernetes, që do të përdoret për të nisur aplikacionin;
- në fjalorin "spec.executor" duhet të jetë e caktuar sasia e resurseve të alokuara për aplikacionin;
- në fjalorin «spec.volumeMounts» duhet të tregohen direktorja lokale në të cilën do të krijohen skedarët lokalë të detyrës Spark.
Shembulli i krijimit të manifestit (këtu {spark-service-account} është llogaria e shërbimit brenda klasës Kubernetes për ekzekutimin e detyrave 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"
Në këtë manifest është e specifikuar llogaria e shërbimit për të cilën kërkohet të krijohen lidhjet e nevojshme të rolit përpara se të publikohet manifesti, duke ofruar të drejtat e nevojshme për të bashkëpunuar aplikacionin Spark me API-në e Kubernetes (nëse nevojitet). Në rastin tonë, aplikacionit i nevojiten të drejtat për krijimin e Pod-ëve. Le të krijojmë lidhjen e nevojshme të rolit:
oc adm policy add-role-to-user edit system:serviceaccount:{project}:{spark-service-account} -n {project}
Vlen tĂ« theksohet gjithashtu se nĂ« specifikimin e kĂ«tij manifesti mund tĂ« specifikohet parametri "hadoopConfigMap", i cili lejon tĂ« specifikohet njĂ« ConfigMap me konfigurimin e Hadoop pa pasur nevojĂ« tĂ« vendoset pĂ«rpara skedari pĂ«rkatĂ«s nĂ« imazhin Docker. Ai Ă«shtĂ« gjithashtu i pĂ«rshtatshĂ«m pĂ«r ekzekutimin e rregullt tĂ« detyrave â me anĂ« tĂ« parametrin "schedule" mund tĂ« specifikohet njĂ« orar pĂ«r ekzekutimin e kĂ«saj detyre.
Pas kësaj, ruajmë manifestin tonë në skedarin spark-pi.yaml dhe e aplikojmë atë në klasterin tonë Kubernetes:
oc apply -f spark-pi.yaml
Kështu do të krijohet një objekt i tipit "sparkapplications":
oc get sparkapplications -n {project}
> NAME AGE
> spark-pi 22h
Në këtë mënyrë do të krijohet një pod me aplikacionin, statusi i të cilit do të shfaqet në "sparkapplications" të krijuar. Mund ta shikoni me anë të komandës së mëposhtme:
oc get sparkapplications spark-pi -o yaml -n {project}
Me përfundimin e detyrës, POD-i do të kalojë në statusin "Completed", i cili gjithashtu do të përditësohet në "sparkapplications". Logët e aplikacionit mund të shihen në shfletues ose me anë të komandës së mëposhtme (këtu {sparkapplications-pod-name} është emri i pod-it të detyrës së ekzekutuar):
oc logs {sparkapplications-pod-name} -n {project}
Menaxhimi i detyrave Spark gjithashtu mund të bëhet përmes një utiliti të specializuar sparkctl. Për ta instaluar, klonojmë repositorin me kodin e tij burimor, instalojmë Go dhe e ndërtojmë këtë 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
Le të shqyrtojmë listën e detyrave të ekzekutuara Spark:
sparkctl list -n {project}
Le të krijojmë një përshkrim për detyrën 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"
Le të ekzekutojmë detyrën e përshkruar përmes sparkctl:
sparkctl create spark-app.yaml -n {project}
Le të shqyrtojmë listën e detyrave të ekzekutuara Spark:
sparkctl list -n {project}
Le të shqyrtojmë listën e ngjarjeve të detyrës së ekzekutuar Spark:
sparkctl event spark-pi -n {project} -f
Le të shqyrtojmë statusin e detyrës së ekzekutuar Spark:
sparkctl status spark-pi -n {project}
Në përfundim, do të doja të shqyrtonja disavantazhet e gjetura të përdorimit të versionit aktual të stabilizuar të Spark (2.4.5) në Kubernetes:
- Disavantazhi i parĂ« dhe ndoshta mĂ« i rĂ«ndĂ«sishmi â Ă«shtĂ« mungesa e Data Locality. PavarĂ«sisht tĂ« gjitha disavantazheve, YARN kishte disa pĂ«rparĂ«si nĂ« pĂ«rdorimin e tij, pĂ«r shembull, parimi i dorĂ«zimit tĂ« kodit te tĂ« dhĂ«nat (nuk janĂ« tĂ« dhĂ«nat te kodi). FalĂ« kĂ«tij parimi, detyrat Spark ekzekutoheshin nĂ« nyjat ku ndodheshin tĂ« dhĂ«nat qĂ« merrnin pjesĂ« nĂ« llogaritje, dhe kohĂ«zgjatja pĂ«r transportin e tĂ« dhĂ«nave nĂ«pĂ«r rrjet zvogĂ«lohej ndjeshĂ«m. Duke pĂ«rdorur Kubernetes, pĂ«rballemi me nevojĂ«n pĂ«r tĂ« transportuar tĂ« dhĂ«nat nĂ« rrjet qĂ« janĂ« tĂ« angazhuara nĂ« punĂ«n e detyrĂ«s. NĂ«se kĂ«to janĂ« mjaft tĂ« mĂ«dha, koha e ekzekutimit tĂ« detyrĂ«s mund tĂ« rritet nĂ« mĂ«nyrĂ« tĂ« dukshme, si dhe do tĂ« kĂ«rkohet njĂ« volum i konsiderueshĂ«m hapĂ«sire disku tĂ« alokuar pĂ«r instancat e detyrĂ«s Spark pĂ«r ruajtjen e tyre pĂ«rkohĂ«sisht. Ky disavantazh mund tĂ« reduktohet pĂ«rmes pĂ«rdorimit tĂ« mjeteve software tĂ« specializuara qĂ« sigurojnĂ« lokalitetin e tĂ« dhĂ«nave nĂ« Kubernetes (pĂ«r shembull, Alluxio), por kjo faktikisht do tĂ« thotĂ« se do tĂ« nevojitet ruajtja e njĂ« kopjeje tĂ« plotĂ« tĂ« tĂ« dhĂ«nave nĂ« nyjat e klasterit Kubernetes.
- MangĂ«si e dytĂ« e rĂ«ndĂ«sishme Ă«shtĂ« siguria. NĂ« mĂ«nyrĂ« tĂ« paracaktuar, funksionet lidhur me sigurimin e ekzekutimit tĂ« detyrave Spark janĂ« tĂ« çaktivizuara, dhe pĂ«rdorimi i Kerberos nĂ« dokumentacionin zyrtar nuk Ă«shtĂ« mbuluar (ndonĂ«se parametrat pĂ«rkatĂ«s u shfaqĂ«n nĂ« versionin 3.0.0, çka kĂ«rkon punĂ« shtesĂ«), dhe nĂ« dokumentacionin pĂ«r sigurimin gjatĂ« pĂ«rdorimit tĂ« Spark (https://spark.apache.org/docs/2.4.5/security.html) si depozita çelesh figurojnĂ« vetĂ«m YARN, Mesos dhe Standalone Cluster. NĂ« kĂ«tĂ« rast, pĂ«rdoruesi qĂ« ekzekuton detyrat Spark nuk mund tĂ« pĂ«rcaktohet drejtpĂ«rdrejt â ne thjesht caktojmĂ« njĂ« llogari shĂ«rbimi nĂ«n tĂ« cilĂ«n do tĂ« punojĂ« pod-i, ndĂ«rsa pĂ«rdoruesi zgjidhet nĂ« pĂ«rputhje me politikat e caktuara tĂ« sigurisĂ«. PĂ«r kĂ«tĂ« arsye, pĂ«rdoruesi mund tĂ« jetĂ« ose root, qĂ« nuk Ă«shtĂ« i sigurt nĂ« njĂ« mjedis prodhimi, ose njĂ« pĂ«rdorues me UID rastĂ«sor, qĂ« Ă«shtĂ« i pĂ«rshtatshĂ«m nĂ« shpĂ«rndarjen e tĂ« drejtave tĂ« aksesit nĂ« tĂ« dhĂ«na (i zgjidhshĂ«m duke krijuar Politikat e SigurisĂ« sĂ« Pod-it dhe lidhjen e tyre me llogaritĂ« pĂ«rkatĂ«se tĂ« shĂ«rbimit). Aktualisht, zgjidhja Ă«shtĂ« ose vendosja e tĂ« gjithĂ« skedarĂ«ve tĂ« nevojshĂ«m direkt nĂ« imazhin Docker, ose modifikimi i skriptit tĂ« ekzekutimit tĂ« Spark pĂ«r tĂ« pĂ«rdorur mekanizmin e ruajtjes dhe marrjes sĂ« sekreteve tĂ« miratuar nga organizata juaj.
- Ekzekutimi i detyrave Spark pĂ«rmes Kubernetes zyrtarisht Ă«shtĂ« ende nĂ« njĂ« fazĂ« eksperimentale dhe nĂ« tĂ« ardhmen mund tĂ« ketĂ« ndryshime tĂ« mĂ«dha nĂ« artefaktet e pĂ«rdorura (skedarĂ«t e konfigurimit, imazhet bazohet nĂ« Docker dhe skriptet e ekzekutimit). Dhe nĂ« tĂ« vĂ«rtetĂ« â gjatĂ« pĂ«rgatitjes sĂ« materialit janĂ« testuar versionet 2.3.0 dhe 2.4.5, qĂ« kanĂ« treguar sjellje tĂ« rĂ«ndĂ«sishme tĂ« ndryshme.
Do tĂ« presim pĂ«r azhurnimet â sapo ka dalĂ« njĂ« version i ri i Spark (3.0.0), qĂ« solli ndryshime tĂ« ndjeshme nĂ« funksionimin e Spark nĂ« Kubernetes, por ruajti statusin eksperimental tĂ« mbĂ«shtetjes pĂ«r kĂ«tĂ« menaxher burimesh. Ndoshta, azhurnimet e ardhshme do tĂ« lejojnĂ« tĂ« rekomandohen plotĂ«sisht heqjen dorĂ« nga YARN dhe tĂ« ekzekutohen detyrat Spark nĂ« Kubernetes, pa u shqetĂ«suar pĂ«r sigurinĂ« e sistemit tuaj dhe pa pasur nevojĂ« pĂ«r modifikimin e pĂ«rbĂ«rĂ«sve funksionalĂ«.
Fin.
Burimi: habr.com


