O Celeborn gerencia dados intermediários de shuffle data e spill data como um remote shuffle service (RSS). Implante-o em um cluster do Container Service for Kubernetes (ACK) para usar o RSS em jobs do Spark.
Benefícios
Para frameworks de big data como MapReduce, Spark e Flink, o Celeborn como RSS oferece:
Escrita de shuffle baseada em push: nós mapper dispensam armazenamento em disco local, ideal para arquiteturas cloud-native com
storage and compute separation.Leitura de shuffle baseada em merge: a mesclagem de dados ocorre nos nós
workerem vez dos reducers, o que evita E/S aleatória de arquivos pequenos e sobrecarga de rede por transferências reduzidas.high availability: nósmasterdo Celeborn usam o protocolo Raft para garantirhigh availabilitye estabilidade.fault tolerance:dual replicasreduzem significativamente a probabilidade de falhas na busca de dados.
Pré-requisitos
Add-on ack-spark-operator implantado.
Cliente kubectl conectado ao cluster ACK. Para mais informações, consulte Conectar-se a um cluster ACK usando kubectl.
Bucket do OSS criado. Para mais informações, consulte Criar um bucket.
-
ossutil instalado e configurado. Para mais informações, consulte Instalar o ossutil e Configurar o ossutil.
Para mais detalhes sobre o comando ossutil, consulte Guia de início rápido da ferramenta de linha de comando ossutil.
node poolscriados e configurados conforme o ambiente de cluster descrito abaixo. Consulte Criar e gerenciar um node pool.
Ambiente do cluster
Este exemplo usa a seguinte configuração de cluster ACK:
-
Processo
masterimplantado nonode poolceleborn-master com a seguinte configuração:Nome do
node pool: celeborn-masterNúmero de nós: 3
ECS instance type: g8i.2xlargelabel: celeborn.apache.org/role=mastertaint: celeborn.apache.org/role=master:NoScheduleArmazenamento de dados por nó: /mnt/celeborn_ratis (1024 GB)
-
Processo
workerimplantado nonode poolceleborn-worker com a seguinte configuração:Nome do
node pool: celeborn-workerNúmero de nós: 5
ECS instance type: g8i.4xlargelabel: celeborn.apache.org/role=workertaint: celeborn.apache.org/role=worker:NoSchedule-
Armazenamento de dados por nó:
/mnt/disk1 (1024 GB)
/mnt/disk2 (1024 GB)
/mnt/disk3 (1024 GB)
/mnt/disk4 (1024 GB)
Visão geral do procedimento
Siga estas etapas para implantar o Celeborn em um cluster ACK e executar um job de exemplo do Spark.
-
Crie uma
container imagedo CelebornBaixe uma
releasedo Celeborn, crie umacontainer imagee envie-a para o seuimage repositorypara implantar o componente ack-celeborn. -
Implante o componente ack-celeborn
Use o
Helm chartack-celeborn disponível noMarketplacedoACKpara implantar um cluster Celeborn com acontainer imagecriada. -
Crie uma
Sparkcontainer imageCrie uma
Sparkcontainer imagecom as dependências do Celeborn e doOSSe envie-a para o seuimage repository. -
Prepare e envie dados de teste para o
OSSGere um conjunto de dados de teste para o job PageRank e faça upload para o
OSS. -
Execute um job de exemplo do
SparkExecute um job PageRank de exemplo e configure-o para usar o Celeborn como RSS.
-
(Opcional) Limpe os recursos
Após concluir o tutorial, remova o job do
Sparke outros recursos desnecessários para evitar cobranças.
Etapa 1: Criar uma container image do Celeborn
Baixe a release necessária (por exemplo, versão 0.5.2) no site oficial do Celeborn. Substitua <IMAGE-REGISTRY> e <IMAGE-REPOSITORY> pelo seu registro de imagem e nome da imagem. Modifique PLATFORMS para definir a arquitetura de destino. Consulte Implantar o Celeborn no Kubernetes. O comando docker buildx requer Docker 19.03 ou posterior. Consulte Instalar e usar Docker e Docker Compose.
CELEBORN_VERSION=0.5.2 # The Celeborn version.
IMAGE_REGISTRY=<IMAGE-REGISTRY> # The image registry, for example, docker.io.
IMAGE_REPOSITORY=<IMAGE-REPOSITORY> # The image name, for example, apache/celeborn.
IMAGE_TAG=${CELEBORN_VERSION} # The image tag. This example uses the Celeborn version as the tag.
PLATFORMS=linux/amd64 # The image platform architecture. To support multiple platforms, separate them with commas, for example, linux/amd64,linux/arm64.
# Download the release package.
wget https://downloads.apache.org/celeborn/celeborn-${CELEBORN_VERSION}/apache-celeborn-${CELEBORN_VERSION}-bin.tgz
# Extract the package.
tar -zxvf apache-celeborn-${CELEBORN_VERSION}-bin.tgz
# Go to the working directory.
cd apache-celeborn-${CELEBORN_VERSION}-bin
# Use Docker Buildx to build the image and push it to the image repository.
docker buildx build \
--output=type=registry \
--push \
--platform=${PLATFORMS} \
--tag=${IMAGE_REGISTRY}/${IMAGE_REPOSITORY}:${IMAGE_TAG} \
-f docker/Dockerfile \
.
Etapa 2: Implantar o componente ack-celeborn
Acesse o ACK console. No painel de navegação à esquerda, clique em .
Na página Marketplace, clique na aba App Catalog, selecione ack-celeborn e, na página ack-celeborn, clique em Deploy.
No painel Create, selecione um cluster e namespace e clique em Next.
-
Na página Parameters, configure os parâmetros e clique em OK.
image: # Replace this with the address of the Celeborn image that you built in Step 1. registry: docker.io # The image registry. repository: apache/celeborn # The image name. tag: 0.5.2 # The image tag. celeborn: celeborn.client.push.stageEnd.timeout: 120s celeborn.master.ha.enabled: true celeborn.master.ha.ratis.raft.server.storage.dir: /mnt/celeborn_ratis celeborn.master.heartbeat.application.timeout: 300s celeborn.master.heartbeat.worker.timeout: 120s celeborn.master.http.port: 9098 celeborn.metrics.enabled: true celeborn.metrics.prometheus.path: /metrics/prometheus celeborn.rpc.dispatcher.numThreads: 4 celeborn.rpc.io.clientThreads: 64 celeborn.rpc.io.numConnectionsPerPeer: 2 celeborn.rpc.io.serverThreads: 64 celeborn.shuffle.chunk.size: 8m celeborn.worker.fetch.io.threads: 32 celeborn.worker.flusher.buffer.size: 256K celeborn.worker.http.port: 9096 celeborn.worker.monitor.disk.enabled: false celeborn.worker.push.io.threads: 32 celeborn.worker.storage.dirs: /mnt/disk1:disktype=SSD:capacity=1024Gi,/mnt/disk2:disktype=SSD:capacity=1024Gi,/mnt/disk3:disktype=SSD:capacity=1024Gi,/mnt/disk4:disktype=SSD:capacity=1024Gi master: replicas: 3 env: - name: CELEBORN_MASTER_MEMORY value: 28g - name: CELEBORN_MASTER_JAVA_OPTS value: -XX:-PrintGC -XX:+PrintGCDetails -XX:+PrintGCTimeStamps -XX:+PrintGCDateStamps -Xloggc:gc-master.out -Dio.netty.leakDetectionLevel=advanced - name: CELEBORN_NO_DAEMONIZE value: "1" - name: TZ value: Asia/Shanghai volumeMounts: - name: celeborn-ratis mountPath: /mnt/celeborn_ratis resources: requests: cpu: 7 memory: 28Gi limits: cpu: 7 memory: 28Gi volumes: - name: celeborn-ratis hostPath: path: /mnt/celeborn_ratis type: DirectoryOrCreate nodeSelector: celeborn.apache.org/role: master tolerations: - key: celeborn.apache.org/role operator: Equal value: master effect: NoSchedule worker: replicas: 5 env: - name: CELEBORN_WORKER_MEMORY value: 28g - name: CELEBORN_WORKER_OFFHEAP_MEMORY value: 28g - name: CELEBORN_WORKER_JAVA_OPTS value: -XX:-PrintGC -XX:+PrintGCDetails -XX:+PrintGCTimeStamps -XX:+PrintGCDateStamps -Xloggc:gc-worker.out -Dio.netty.leakDetectionLevel=advanced - name: CELEBORN_NO_DAEMONIZE value: "1" - name: TZ value: Asia/Shanghai volumeMounts: - name: disk1 mountPath: /mnt/disk1 - name: disk2 mountPath: /mnt/disk2 - name: disk3 mountPath: /mnt/disk3 - name: disk4 mountPath: /mnt/disk4 resources: requests: cpu: 14 memory: 56Gi limits: cpu: 14 memory: 56Gi volumes: - name: disk1 hostPath: path: /mnt/disk1 type: DirectoryOrCreate - name: disk2 hostPath: path: /mnt/disk2 type: DirectoryOrCreate - name: disk3 hostPath: path: /mnt/disk3 type: DirectoryOrCreate - name: disk4 hostPath: path: /mnt/disk4 type: DirectoryOrCreate nodeSelector: celeborn.apache.org/role: worker tolerations: - key: celeborn.apache.org/role operator: Equal value: worker effect: NoScheduleA tabela a seguir descreve os principais parâmetros. Para a lista completa, consulte a seção ConfigMaps na página do ack-celeborn.
-
Aguarde a conclusão da implantação do Celeborn. Se encontrar problemas com pods, consulte Solução de problemas de Pod.
kubectl get -n celeborn statefulsetSaída esperada:
NAME READY AGE celeborn-master 3/3 68s celeborn-worker 5/5 68s
Etapa 3: Criar uma container image do Spark
Este exemplo usa o Spark 3.5.3. Crie um Dockerfile com o conteúdo abaixo para gerar uma container image e enviá-la ao seu image repository.
ARG SPARK_IMAGE=<SPARK_IMAGE> # Replace <SPARK_IMAGE> with your Spark base image.
FROM ${SPARK_IMAGE}
# Add dependency for Hadoop Aliyun OSS support
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aliyun/3.3.4/hadoop-aliyun-3.3.4.jar ${SPARK_HOME}/jars
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/com/aliyun/oss/aliyun-sdk-oss/3.17.4/aliyun-sdk-oss-3.17.4.jar ${SPARK_HOME}/jars
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/org/jdom/jdom2/2.0.6.1/jdom2-2.0.6.1.jar ${SPARK_HOME}/jars
# Add dependency for Celeborn
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/org/apache/celeborn/celeborn-client-spark-3-shaded_2.12/0.5.1/celeborn-client-spark-3-shaded_2.12-0.5.1.jar ${SPARK_HOME}/jars
Etapa 4: Enviar dados de teste para o OSS
Para preparar os dados de teste e enviá-los ao OSS, consulte Etapa 1: Preparar dados de teste e enviá-los ao OSS.
Etapa 5: Criar um Secret do OSS
Para criar um Secret com as access credentials do OSS, consulte Etapa 3: Criar um Secret para armazenar credenciais de acesso do OSS.
Etapa 6: Enviar um job de exemplo do Spark
Crie um arquivo de manifesto SparkApplication chamado spark-pagerank.yaml com o conteúdo a seguir. Substitua <SPARK_IMAGE> pela sua imagem criada na Etapa 3: Criar uma container image do Spark, e substitua <OSS_BUCKET> e <OSS_ENDPOINT> pelo seu OSS bucket e endpoint. Para configurações do Spark, consulte a documentação do Celeborn.
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
image: <SPARK_IMAGE> # The Spark image. Replace <SPARK_IMAGE> with your Spark image name.
mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.3.jar
mainClass: org.apache.spark.examples.SparkPageRank
arguments:
- oss://<OSS_BUCKET>/data/pagerank_dataset.txt # The input test dataset. Replace <OSS_BUCKET> with your OSS bucket name.
- "10" # The number of iterations.
sparkVersion: 3.5.3
hadoopConf:
fs.AbstractFileSystem.oss.impl: org.apache.hadoop.fs.aliyun.oss.OSS
fs.oss.impl: org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystem
fs.oss.endpoint: <OSS_ENDPOINT> # The OSS endpoint. For example, the internal endpoint for OSS in the China (Beijing) region is oss-cn-beijing-internal.aliyuncs.com.
fs.oss.credentials.provider: com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider
sparkConf:
spark.shuffle.manager: org.apache.spark.shuffle.celeborn.SparkShuffleManager
spark.serializer: org.apache.spark.serializer.KryoSerializer
spark.celeborn.master.endpoints: celeborn-master-0.celeborn-master-svc.celeborn.svc.cluster.local,celeborn-master-1.celeborn-master-svc.celeborn.svc.cluster.local,celeborn-master-2.celeborn-master-svc.celeborn.svc.cluster.local
spark.celeborn.client.spark.shuffle.writer: hash
spark.celeborn.client.push.replicate.enabled: "false"
spark.sql.adaptive.localShuffleReader.enabled: "false"
spark.sql.adaptive.enabled: "true"
spark.sql.adaptive.skewJoin.enabled: "true"
spark.shuffle.sort.io.plugin.class: org.apache.spark.shuffle.celeborn.CelebornShuffleDataIO
spark.dynamicAllocation.shuffleTracking.enabled: "false"
spark.executor.userClassPathFirst: "false"
driver:
cores: 1
coreLimit: 1200m
memory: 512m
serviceAccount: spark-operator-spark
envFrom:
- secretRef:
name: spark-oss-secret
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-oss-secret
restartPolicy:
type: Never
(Opcional) Etapa 7: Limpar recursos
Após concluir este tutorial, exclua estes recursos para evitar cobranças.
Exclua o job do Spark:
kubectl delete sparkapplication spark-pagerank
Exclua o recurso Secret:
kubectl delete secret spark-oss-secret
Referências
Para enviar jobs do
Sparkcom o Spark Operator, consulte Usar o Spark Operator para executar jobs do Spark.Para visualizar o histórico de jobs do
Spark, consulte Usar o Spark History Server para visualizar informações sobre jobs do Spark.Consulte a documentação do Apache Celeborn.