Instale o Spark Operator no ACK para gerenciar todo o ciclo de vida dos jobs do Spark usando manifests declarativos do Kubernetes.
Pré-requisitos
Um cluster ACK Pro ou cluster ACK Serverless Pro executando Kubernetes 1.24 ou posterior. Consulte Criar um cluster gerenciado pelo ACK, Criar um cluster ACK Serverless e Atualizar manualmente clusters ACK.
Cliente kubectl conectado ao cluster ACK. Para mais informações, consulte Conectar-se a um cluster ACK usando kubectl.
Como funciona
O Spark Operator automatiza o ciclo de vida de jobs do Spark no Kubernetes usando recursos CustomResourceDefinition (CRD), como SparkApplication e ScheduledSparkApplication. Ele utiliza funcionalidades nativas do Kubernetes, como dimensionamento automático, verificações de integridade e gerenciamento de recursos. O ACK fornece o ack-spark-operator, baseado no projeto kubeflow/spark-operator. Consulte Spark Operator | Kubeflow.
Benefícios:
Gerenciamento simplificado: Automatize a implantação e o ciclo de vida de jobs do Spark com configurações declarativas do Kubernetes.
Suporte a multilocação: Use namespaces e cotas de recursos do Kubernetes para isolar recursos. Selecione nós específicos para executar cargas de trabalho do Spark em recursos dedicados.
Provisionamento elástico de recursos: Escale com recursos elásticos, como Elastic Container Instance (ECI) ou pools de nós elásticos durante picos de carga, equilibrando desempenho e custo.
Casos de uso:
Análise de dados: Use o Spark para análises interativas e limpeza de dados.
Computação em lote: Execute jobs em lote agendados para processar grandes conjuntos de dados.
Processamento em tempo real: Processe fluxos de dados em tempo real com o Spark Streaming.
Visão geral do procedimento
O fluxo de trabalho abrange a implantação do Spark Operator, o envio de jobs, o monitoramento da execução e o gerenciamento do ciclo de vida do job.
Implantar o componente ack-spark-operator: Instale o Spark Operator no cluster ACK.
Enviar um job do Spark: Crie e envie um manifesto de job do Spark.
Monitorar o job do Spark: Verifique o status do job, o status dos pods e os logs.
Acessar a interface web do Spark: Visualize os detalhes da execução do job no navegador.
Atualizar o job do Spark: Modifique e reaplique o manifesto do job.
Excluir o job do Spark: Remova o job e libere os recursos.
Etapa 1: Implantar o componente ack-spark-operator
Faça login no ACK console. No painel de navegação à esquerda, clique em .
Na página Marketplace, clique na aba App Catalog, pesquise e selecione ack-spark-operator.
Na página do ack-spark-operator, 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.
Os principais parâmetros estão listados abaixo. Para obter a lista completa, consulte a aba ConfigMaps na página do ack-spark-operator.
Parâmetro
Descrição
Padrão
controller.replicasNúmero de réplicas do controlador.
1
webhook.replicasNúmero de réplicas do webhook.
1
spark.jobNamespacesNamespaces onde os jobs do Spark podem ser executados. Uma string vazia permite todos os namespaces. Separe vários valores com vírgulas (
,).-
["default"](padrão) -
[""](todos os namespaces) -
["ns1","ns2","ns3"](vários namespaces)
spark.serviceAccount.nameO Spark Operator cria uma ServiceAccount chamada
spark-operator-sparke os recursos RBAC necessários em cada namespace especificado porspark.jobNamespaces. Se personalizado, especifique o novo nome ao enviar jobs do Spark.spark-operator-spark -
Etapa 2: Enviar um job do Spark
Crie um manifesto SparkApplication para enviar um job do Spark.
-
Crie o seguinte manifesto e salve-o como
spark-pi.yaml.apiVersion: sparkoperator.k8s.io/v1beta2 kind: SparkApplication metadata: name: spark-pi namespace: default # The namespace must be in the list of namespaces specified by spark.jobNamespaces. spec: type: Scala mode: cluster image: registry-cn-hangzhou.ack.aliyuncs.com/ack-demo/spark:3.5.4 imagePullPolicy: IfNotPresent mainClass: org.apache.spark.examples.SparkPi mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.4.jar arguments: - "1000" sparkVersion: 3.5.4 driver: cores: 1 coreLimit: 1200m memory: 512m serviceAccount: spark-operator-spark # If you customized the ServiceAccount name, change the value accordingly. executor: instances: 1 cores: 1 coreLimit: 1200m memory: 512m restartPolicy: type: Never -
Envie o job do Spark:
kubectl apply -f spark-pi.yamlSaída esperada:
sparkapplication.sparkoperator.k8s.io/spark-pi created
Etapa 3: Monitorar o job do Spark
Verifique o status, os pods e os logs do job do Spark.
-
Verifique o status do job do Spark:
kubectl get sparkapplication spark-piSaída esperada:
NAME STATUS ATTEMPTS START FINISH AGE spark-pi SUBMITTED 1 2024-06-04T03:17:11Z <no value> 15s -
Verifique o status dos pods. Este comando filtra os pods pelo rótulo
sparkoperator.k8s.io/app-name=spark-pi:kubectl get pod -l sparkoperator.k8s.io/app-name=spark-piSaída esperada:
NAME READY STATUS RESTARTS AGE spark-pi-driver 1/1 Running 0 49s spark-pi-7272428fc8f5f392-exec-1 1/1 Running 0 13sApós a conclusão do job, o driver exclui todos os pods executores.
-
Visualize os detalhes do job do Spark:
kubectl describe sparkapplication spark-pi -
Visualize as últimas 20 linhas dos logs do pod driver:
kubectl logs --tail=20 spark-pi-driverSaída esperada:
24/05/30 10:05:30 INFO TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool 24/05/30 10:05:30 INFO DAGScheduler: ResultStage 0 (reduce at SparkPi.scala:38) finished in 7.942 s 24/05/30 10:05:30 INFO DAGScheduler: Job 0 is finished. Cancelling potential speculative or zombie tasks for this job 24/05/30 10:05:30 INFO TaskSchedulerImpl: Killing all running tasks in stage 0: Stage finished 24/05/30 10:05:30 INFO DAGScheduler: Job 0 finished: reduce at SparkPi.scala:38, took 8.043996 s Pi is roughly 3.1419522314195225 24/05/30 10:05:30 INFO SparkContext: SparkContext is stopping with exitCode 0. 24/05/30 10:05:30 INFO SparkUI: Stopped Spark web UI at http://spark-pi-1e18858fc8f56b14-driver-svc.default.svc:4040 24/05/30 10:05:30 INFO KubernetesClusterSchedulerBackend: Shutting down all executors 24/05/30 10:05:30 INFO KubernetesClusterSchedulerBackend$KubernetesDriverEndpoint: Asking each executor to shut down 24/05/30 10:05:30 WARN ExecutorPodsWatchSnapshotSource: Kubernetes client has been closed. 24/05/30 10:05:30 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped! 24/05/30 10:05:30 INFO MemoryStore: MemoryStore cleared 24/05/30 10:05:30 INFO BlockManager: BlockManager stopped 24/05/30 10:05:30 INFO BlockManagerMaster: BlockManagerMaster stopped 24/05/30 10:05:30 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped! 24/05/30 10:05:30 INFO SparkContext: Successfully stopped SparkContext 24/05/30 10:05:30 INFO ShutdownHookManager: Shutdown hook called 24/05/30 10:05:30 INFO ShutdownHookManager: Deleting directory /var/data/spark-14ed60f1-82cd-4a33-b1b3-9e5d975c5b1e/spark-01120c89-5296-4c83-8a20-0799eef4e0ee 24/05/30 10:05:30 INFO ShutdownHookManager: Deleting directory /tmp/spark-5f98ed73-576a-41be-855d-dabdcf7de189
Etapa 4: Acessar a interface web do Spark
A interface web fica disponível apenas enquanto o pod driver está em execução (Running).
Por padrão, controller.uiService.enable é true, o que cria um Service para expor a interface web via redirecionamento de porta. Se definido como false, nenhum Service será criado e você deverá fazer o redirecionamento de porta diretamente do pod driver.
O comando kubectl port-forward é adequado para testes, mas não recomendado para produção devido a riscos de segurança.
-
Encaminhe a porta da interface web para sua máquina local:
-
Redirecionamento via Service
kubectl port-forward services/spark-pi-ui-svc 4040 -
Redirecionamento via pod
kubectl port-forward pods/spark-pi-driver 4040Saída esperada:
Forwarding from 127.0.0.1:4040 -> 4040 Forwarding from [::1]:4040 -> 4040
-
Abra http://127.0.0.1:4040 no navegador.
(Opcional) Etapa 5: Atualizar o job do Spark
Atualize o manifesto do job para modificar os parâmetros do job do Spark.
-
Edite o arquivo
spark-pi.yaml. Por exemplo, definaargumentscomo10000e as instâncias doexecutorcomo2.apiVersion: sparkoperator.k8s.io/v1beta2 kind: SparkApplication metadata: name: spark-pi spec: type: Scala mode: cluster image: registry-cn-hangzhou.ack.aliyuncs.com/ack-demo/spark:3.5.4 imagePullPolicy: IfNotPresent mainClass: org.apache.spark.examples.SparkPi mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.4.jar arguments: - "10000" sparkVersion: 3.5.4 driver: cores: 1 coreLimit: 1200m memory: 512m serviceAccount: spark-operator-spark # If you customized the ServiceAccount name, change the value accordingly. executor: instances: 2 cores: 1 coreLimit: 1200m memory: 512m restartPolicy: type: Never -
Aplique as alterações:
kubectl apply -f spark-pi.yaml -
Verifique o status do job:
kubectl get sparkapplication spark-piO job do Spark será executado novamente. Saída esperada:
NAME STATUS ATTEMPTS START FINISH AGE spark-pi RUNNING 1 2024-06-04T03:37:34Z <no value> 20m
(Opcional) Etapa 6: Excluir o job do Spark
Exclua o job do Spark para liberar seus recursos associados.
kubectl delete -f spark-pi.yaml
Alternativamente:
kubectl delete sparkapplication spark-pi