É possível enviar um job para um cluster do Serviço Gerenciado para Apache Spark usando uma solicitação HTTP ou programática da API jobs.submit, a ferramenta de linha de comando gcloud da Google Cloud CLI em uma janela de terminal local ou no Cloud Shell ou pelo console doGoogle Cloud aberto em um navegador local. Também é possível executar SSH na instância mestre do cluster e executar um job diretamente na instância sem usar o Serviço Gerenciado para Apache Spark.
Simultaneidade do job : é possível configurar o número máximo de jobs simultâneos do Serviço Gerenciado para Apache Spark
com a propriedade
dataproc:dataproc.scheduler.max-concurrent-jobs
ao criar um cluster. Se esse valor de propriedade não estiver definido,
o limite superior em jobs simultâneos será calculado como
max((masterMemoryMb - 3584) / masterMemoryMbPerJob, 5).
masterMemoryMb é determinado pelo tipo de máquina da VM mestre.
masterMemoryMbPerJob é 1024 por padrão, mas pode ser configurado na criação do cluster com a propriedade de cluster dataproc:dataproc.scheduler.driver-size-mb.
Como enviar um job
Isolamento de classpath e JARs personalizados:não copie "fat JARs" personalizados ou JARs de pacote (como o tempo de execução do Apache Iceberg ou pacotes Google Cloud ) diretamente para diretórios do sistema de cluster, como /usr/lib/spark/jars/. Colocar JARs personalizados em diretórios do sistema polui o classpath do agente do Serviço gerenciado para Apache Spark com dependências transitivas (como bibliotecas Guava ou Hadoop) que entram em conflito com as bibliotecas integradas do agente. Esse conflito pode causar erros de resolução de classe (como ClassNotFoundException ou NoClassDefFoundError) ao realizar operações de gerenciamento de jobs, como cancelar um job, fazendo com que o agente falhe e deixando aplicativos YARN órfãos.
Em vez disso, use uma das seguintes abordagens compatíveis:
- Dependências específicas do job:especifique o caminho do Cloud Storage para seus JARs ao enviar o job usando a flag
--jarsna CLI do Google Cloud, o campo Arquivos JAR no console Google Cloud ou o campojarFileUrisna API. O Spark distribui as dependências para o driver e os executores do job sem contaminar o classpath do agente. - Dependências em todo o cluster:especifique as dependências durante a criação do cluster usando a propriedade do cluster
spark:spark.jars: Isso instrui o Serviço Gerenciado para Apache Spark a configurar--properties="spark:spark.jars=gs://YOUR_BUCKET/jar-1.jar,gs://YOUR_BUCKET/jar-2.jar"
/etc/spark/conf/spark-defaults.conf, distribuindo automaticamente as dependências para os caminhos de classe do driver e do executor do Spark em todos os jobs, enquanto deixa o agente do Serviço Gerenciado para Apache Spark local do nó isolado.
Console
Abra a página Enviar um job do Serviço Gerenciado para Apache Spark no console Google Cloud do seu navegador.
Exemplo de job do Spark
Para enviar um job do Spark de exemplo, preencha os campos na página Enviar um job da seguinte maneira:
- Selecione o nome do Cluster na lista de clusters.
- Defina o Tipo de job para
Spark. - Defina Classe principal ou jar como
org.apache.spark.examples.SparkPi. - Defina Argumentos como o argumento único
1000. - Adicione
file:///usr/lib/spark/examples/jars/spark-examples.jarpara Arquivos jar (ou o campojarFileUrisda API):file:///indica um esquema de LocalFileSystem do Hadoop. O Serviço Gerenciado para Apache Spark instalou/usr/lib/spark/examples/jars/spark-examples.jarno nó mestre do cluster quando criou o cluster. Esse caminho é usado apenas para exemplos pré-instalados fornecidos pelo Serviço Gerenciado para Apache Spark. Não copie JARs personalizados para diretórios do sistema (consulte a observação em Como enviar um job).- Como alternativa, você pode especificar um caminho do Cloud Storage
(
gs://your-bucket/your-jarfile.jar) ou um caminho do sistema de arquivos distribuídos do Hadoop (hdfs://path-to-jar.jar) para um dos seus jars. Se você enviar o job usando a API, especifique esse caminho no campojarFileUris.
Clique em Enviar para iniciar o job. Depois de iniciado, o job será adicionado à lista.
Clique no código da tarefa para abrir a página Jobs, em que você pode conferir a saída do driver do job. Como este job produz linhas de saída longas que
excedem a largura da janela do navegador, você pode marcar a caixa de Quebra de linha para exibir todo
o texto de saída e mostrar o resultado calculado para pi.
Visualize a saída do driver de job na linha de comando usando o comando
gcloud dataproc jobs wait
mostrado abaixo. Para mais informações, consulte
Ver saída do job – COMANDO GCLOUD.
Copie e cole o ID do projeto como o valor para a flag --project e o código da tarefa (mostrado na lista de jobs) como o argumento final.
gcloud dataproc jobs wait job-id \ --project=project-id \ --region=region
Confira snippets da saída do driver do job SparkPi
de exemplo:
... 2015-06-25 23:27:23,810 INFO [dag-scheduler-event-loop] scheduler.DAGScheduler (Logging.scala:logInfo(59)) - Stage 0 (reduce at SparkPi.scala:35) finished in 21.169 s 2015-06-25 23:27:23,810 INFO [task-result-getter-3] cluster.YarnScheduler (Logging.scala:logInfo(59)) - Removed TaskSet 0.0, whose tasks have all completed, from pool 2015-06-25 23:27:23,819 INFO [main] scheduler.DAGScheduler (Logging.scala:logInfo(59)) - Job 0 finished: reduce at SparkPi.scala:35, took 21.674931 s Pi is roughly 3.14189648 ... Job [c556b47a-4b46-4a94-9ba2-2dcee31167b2] finished successfully. driverOutputUri: gs://sample-staging-bucket/google-cloud-dataproc-metainfo/cfeaa033-749e-48b9-... ...
gcloud
Para enviar um job a um cluster do serviço gerenciado para Apache Spark, execute o comando da CLI gcloud gcloud dataproc jobs submit localmente em uma janela de terminal ou no Cloud Shell.
gcloud dataproc jobs submit job-command \ --cluster=cluster-name \ --region=region \ other dataproc-flags \ -- job-args
- Liste os
hello-world.pyacessíveis publicamente localizados no Cloud Storage. Listagem de arquivos:gcloud storage cat gs://dataproc-examples/pyspark/hello-world/hello-world.py
#!/usr/bin/python import pyspark sc = pyspark.SparkContext() rdd = sc.parallelize(['Hello,', 'world!']) words = sorted(rdd.collect()) print(words)
- Envie o job do PySpark para o Serviço Gerenciado para Apache Spark.
Saída do terminal:gcloud dataproc jobs submit pyspark \ gs://dataproc-examples/pyspark/hello-world/hello-world.py \ --cluster=cluster-name \ --region=region
Waiting for job output... … ['Hello,', 'world!'] Job finished successfully.
- Execute o exemplo SparkPi pré-instalado no nó mestre do
cluster do Serviço Gerenciado para Apache Spark. O caminho
file:///usr/lib/spark/examples/jars/spark-examples.jaré apenas para exemplos pré-instalados. Para dependências JAR personalizadas, consulte a observação em Como enviar um job. Saída do terminal:gcloud dataproc jobs submit spark \ --cluster=cluster-name \ --region=region \ --class=org.apache.spark.examples.SparkPi \ --jars=file:///usr/lib/spark/examples/jars/spark-examples.jar \ -- 1000
Job [54825071-ae28-4c5b-85a5-58fae6a597d6] submitted. Waiting for job output… … Pi is roughly 3.14177148 … Job finished successfully. …
REST
Nesta seção, mostramos como enviar um job do Spark para calcular o valor aproximado de
pi usando a API
jobs.submit do Serviço Gerenciado para Apache Spark.
Antes de usar os dados da solicitação abaixo, faça as substituições a seguir:
- project-id: Google Cloud ID do projeto
- region: região do cluster
- clusterName: nome do cluster
Método HTTP e URL:
POST https://dataproc.googleapis.com/v1/projects/project-id/regions/region/jobs:submit
Corpo JSON da solicitação:
{
"job": {
"placement": {
"clusterName": "cluster-name"
},
"sparkJob": {
"args": [
"1000"
],
"mainClass": "org.apache.spark.examples.SparkPi",
"jarFileUris": [
"file:///usr/lib/spark/examples/jars/spark-examples.jar"
]
}
}
}
Para enviar a solicitação, expanda uma destas opções:
Você receberá uma resposta JSON semelhante a esta:
{
"reference": {
"projectId": "project-id",
"jobId": "job-id"
},
"placement": {
"clusterName": "cluster-name",
"clusterUuid": "cluster-Uuid"
},
"sparkJob": {
"mainClass": "org.apache.spark.examples.SparkPi",
"args": [
"1000"
],
"jarFileUris": [
"file:///usr/lib/spark/examples/jars/spark-examples.jar"
]
},
"status": {
"state": "PENDING",
"stateStartTime": "2020-10-07T20:16:21.759Z"
},
"jobUuid": "job-Uuid"
}
Java
Python
Go
Node.js
Enviar um trabalho diretamente no cluster
Se você quiser executar um job diretamente no cluster sem usar o Serviço Gerenciado para Apache Spark, use o SSH no nó mestre do cluster e execute o job no nó mestre.
Depois de estabelecer uma conexão SSH com a instância mestre de VM, execute comandos em uma janela de terminal no nó mestre do cluster para:
- abrir um shell do Spark;
- executar um job do Spark para contar o número de linhas em um arquivo "hello-world" do Python (sete linhas) localizado em um arquivo do Cloud Storage acessível publicamente;
sair do shell.
user@cluster-name-m:~$ spark-shell ... scala> sc.textFile("gs://dataproc-examples" + "/pyspark/hello-world/hello-world.py").count ... res0: Long = 7 scala> :quit
Executar jobs do Bash no Serviço Gerenciado para Apache Spark
Execute um script bash como job do Serviço Gerenciado para Apache Spark porque os mecanismos usados não são compatíveis com um tipo de serviço de nível superior ou porque você precisa configurar ou calcular os argumentos antes de iniciar um job usando hadoop ou spark-submit do seu script.
Exemplo de Python
Suponha que você tenha copiado um script hello.sh bash no Cloud Storage:
gcloud storage cp hello.sh gs://${BUCKET}/hello.shComo o comando pig fs usa caminhos do Hadoop, copie o script do
Cloud Storage para um destino especificado como file:/// para garantir
que ele esteja no sistema de arquivos local em vez do HDFS. Os comandos sh subsequentes
fazem referência ao sistema de arquivos local automaticamente e não exigem o prefixo
file:///.
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
-e='fs -cp -f gs://${BUCKET}/hello.sh file:///tmp/hello.sh; sh chmod 750 /tmp/hello.sh; sh /tmp/hello.sh'Como alternativa, como os jobs do Serviço Gerenciado para Apache Spark enviam o argumento --jars e organizam um arquivo em um diretório temporário criado durante o ciclo de vida do job, especifique o script de shell do Cloud Storage como um argumento --jars:
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
--jars=gs://${BUCKET}/hello.sh \
-e='sh chmod 750 ${PWD}/hello.sh; sh ${PWD}/hello.sh'Observe que o argumento --jars também pode fazer referência a um script local:
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
--jars=hello.sh \
-e='sh chmod 750 ${PWD}/hello.sh; sh ${PWD}/hello.sh'